《Kafka 消费者开发与注意事项》

《Kafka 消费者开发与注意事项》

本文介绍 Kafka 消费者的核心概念、Java 消费示例与易踩坑点。重点是同一个 group.id 内消费者数量与分区数的关系、offset 重置策略等。

1 消费者核心概念

  • Consumer Group(消费组):同一 group.id 的消费者共同消费一个或多个 topic,实现负载均衡与故障转移。
  • 分区分配:一个分区同一时刻只能被某个消费组内的一个消费者消费。
  • offset 提交:消费者消费完记录后提交位移,用于故障恢复时续读。

2 Java 消费示例

import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class SimpleConsumer { public static void main(String[] args) { Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "192.168.1.10:9092"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "app-log-group"); // offset 超出范围时重置为 earliest(从最早开始)或 latest(只读新增) props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) { consumer.subscribe(Collections.singletonList("test")); while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); for (ConsumerRecord<String, String> record : records) { System.out.printf("partition=%s offset=%d value=%s%n", record.partition(), record.offset(), record.value()); } } } } }

3 注意事项

1 同一个 group.id 内的消费者数量不能超过分区数

参考:Kafka: The Definitive Guide Ch04

分区与消费者是一对一关系,一个分区最多只能被一个消费者消费。当消费者数量超过分区数时,多出的消费者会处于空闲(没有可供消费的分区),白白占用资源。

  • 消费者数量 ≤ 分区数,才能充分并行消费。
  • 需要更高吞吐应增加分区数,而不是盲目增加消费者。

2 offset 越界(OffsetOutOfRangeException)

消息因保留期删除,或从旧 offset 消费时,可能触发 offset 越界。建议设置正确的重置策略:

props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"); // 或 earliest

3 同一 group.id 被多个应用复用的冲突

多个 application 共用同一个 group.id 会导致消费者组分配错乱(详见《Kafka 常见问题排查》)。

4 参考

阅读 — · 全站 —
🎸 我的歌单 0 首