Hello Kafka:生产者与消费者入门
Kafka Java 客户端入门,演示同步/异步发送、回调函数和消费者轮询等核心 API。
Hello Kafka:生产者与消费者入门
Kafka 是什么
Kafka 是分布式消息队列,和 RabbitMQ 的核心区别在于设计哲学:
| 特性 | Kafka | RabbitMQ |
|---|---|---|
| 模型 | 分布式日志(Pull) | 队列(Push/AMQP) |
| 持久化 | 设计为写磁盘 | 可选持久化 |
| 吞吐 | 单机百万条/秒 | 单机万条/秒 |
| 消息顺序 | 单分区内有序 | 单队列有序 |
| 数据保留 | 按时间/大小保留 | 消费后删除 |
Kafka 的高吞吐来自顺序写磁盘——消息不断追加到文件末尾,而不是随机写。
核心概念
Producer → Topic(可以分区) → Consumer Group
Topic:消息的逻辑分类 Partition:Topic 的物理分片,一个 Topic 可以有多个 Partition,实现并行 Consumer Group:组内各消费者消费不同 Partition
Java Producer 示例
Properties props = new Properties();
props.put("bootstrap.servers", "192.168.8.133:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// acks=0:性能最高,可能丢消息
// acks=1:Leader 写入成功即返回(默认)
// acks=all:所有副本写入成功才返回,最安全
props.put("acks", "1");
props.put("retries", 3);
props.put("batch.size", 16384); // 批量发送,提升吞吐
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("my-topic", "key", "value"));
producer.close();
Java Consumer 示例
Properties props = new Properties();
props.put("bootstrap.servers", "192.168.8.133:9092");
props.put("group.id", "my-group");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("enable.auto.commit", "true");
props.put("auto.commit.interval.ms", "1000");
props.put("auto.offset.reset", "earliest"); // 从头开始消费
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("my-topic"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("offset=%d, key=%s, value=%s%n",
record.offset(), record.key(), record.value());
}
}
生产注意事项
- acks 参数:
0最高性能但可能丢消息;1兼顾性能和安全(默认);all最安全但吞吐最低 - 消息顺序:只在一个 Partition 内保证。需要全局顺序则只能用一个 Partition(牺牲吞吐)
- 消费幂等:Kafka 是 at-least-once 语义,消费端需要幂等处理
- 版本兼容:客户端版本 ≤ 服务端版本,否则可能连接失败
- 生产环境:推荐使用
spring-boot-starter-kafka替代原生 API