跳到正文
Joeplover
学习笔记·2025-11-06·约 2 分钟阅读

Hello Kafka:生产者与消费者入门

Kafka Java 客户端入门,演示同步/异步发送、回调函数和消费者轮询等核心 API。

数据流管道

Hello Kafka:生产者与消费者入门

Kafka 是什么

Kafka 是分布式消息队列,和 RabbitMQ 的核心区别在于设计哲学:

特性KafkaRabbitMQ
模型分布式日志(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());
    }
}

生产注意事项

  1. acks 参数:0 最高性能但可能丢消息;1 兼顾性能和安全(默认);all 最安全但吞吐最低
  2. 消息顺序:只在一个 Partition 内保证。需要全局顺序则只能用一个 Partition(牺牲吞吐)
  3. 消费幂等:Kafka 是 at-least-once 语义,消费端需要幂等处理
  4. 版本兼容:客户端版本 ≤ 服务端版本,否则可能连接失败
  5. 生产环境:推荐使用 spring-boot-starter-kafka 替代原生 API