Java Kafka Consumer拉取消息返回0条及Group Coordinator异常问题
问题:Kafka Consumer拉取消息始终返回0条,Group心跳超时导致成员被移除
运行Java Kafka Consumer代码从Kafka Server拉取消息时,始终返回0条消息,查看Kafka Server日志发现Group相关错误,日志信息如下:
[2022-09-08 05:03:03,785] INFO [GroupCoordinator 1]: Dynamic Member with unknown member id joins group KafkaStudy in Empty state. Created a new member id consumer-1-e99d5a30-346b-4e26-8ff5-7f0e5fac3e9b for this member and add to the group. (kafka.coordinator.group.GroupCoordinator) [2022-09-08 05:03:03,785] INFO [GroupCoordinator 1]: Preparing to rebalance group KafkaStudy in state PreparingRebalance with old generation 27 (__consumer_offsets-15) (reason: Adding new member consumer-1-e99d5a30-346b-4e26-8ff5-7f0e5fac3e9b with group instance id None; client reason: not provided) (kafka.coordinator.group.GroupCoordinator) [2022-09-08 05:03:03,788] INFO [GroupCoordinator 1]: Stabilized group KafkaStudy generation 28 (__consumer_offsets-15) with 1 members (kafka.coordinator.group.GroupCoordinator) [2022-09-08 05:03:03,809] INFO [GroupCoordinator 1]: Assignment received from leader consumer-1-e99d5a30-346b-4e26-8ff5-7f0e5fac3e9b for group KafkaStudy for generation 28. The group has 1 members, 0 of which are static. (kafka.coordinator.group.GroupCoordinator) [2022-09-08 05:03:13,819] INFO [GroupCoordinator 1]: Member consumer-1-e99d5a30-346b-4e26-8ff5-7f0e5fac3e9b in group KafkaStudy has failed, removing it from the group (kafka.coordinator.group.GroupCoordinator) [2022-09-08 05:03:13,820] INFO [GroupCoordinator 1]: Preparing to rebalance group KafkaStudy in state PreparingRebalance with old generation 28 (__consumer_offsets-15) (reason: removing member consumer-1-e99d5a30-346b-4e26-8ff5-7f0e5fac3e9b on heartbeat expiration) (kafka.coordinator.group.GroupCoordinator) [2022-09-08 05:03:13,821] INFO [GroupCoordinator 1]: Group KafkaStudy with generation 29 is now empty (__consumer_offsets-15) (kafka.coordinator.group.GroupCoordinator)
对应的Java代码如下:
Properties properties = new Properties(); properties.put("bootstrap.servers","192.168.226.133:9092"); properties.put("key.deserializer","org.apache.kafka.common.serialization.StringDeserializer"); properties.put("value.deserializer","org.apache.kafka.common.serialization.StringDeserializer"); properties.put("enable.auto.commit", true); properties.put("group.id","KafkaStudy"); properties.put("group.instance.id","1"); consumer = new KafkaConsumer<String, String>(properties); consumer.subscribe(Collections.singleton("quickstart-events")); try{ ConsumerRecords<String,String> records = consumer.poll(100); for (ConsumerRecord<String, String> record : records) {System.out.println(record.value());} } finally{ consumer.close(); }
问题分析
从日志可以看到,Consumer成功加入Group并完成分区分配后,很快因为心跳超时被Coordinator标记为失效并移除。核心原因在于:
- 代码仅调用了一次
poll(100),超时时间仅100毫秒,Consumer还没来得及完成完整的组初始化流程,就进入finally块被关闭,后续无法发送心跳维持成员身份 - 配置的
group.instance.id在日志中显示为None,若使用的Kafka版本低于2.3.0,静态成员特性未被支持,可能引发额外的协调问题,但这不是当前无消息返回的主因
解决方案
- 循环调用poll方法:Consumer需要持续运行并定期调用poll,才能维持与Coordinator的心跳连接,同时拉取消息
- 调整poll超时时间:将超时时间设置为1000毫秒左右,给组协调和消息拉取留足够时间
- 移除不必要的静态成员配置:如果不需要静态成员特性,删除
group.instance.id配置,避免版本兼容性问题
修改后的代码示例:
Properties properties = new Properties(); properties.put("bootstrap.servers","192.168.226.133:9092"); properties.put("key.deserializer","org.apache.kafka.common.serialization.StringDeserializer"); properties.put("value.deserializer","org.apache.kafka.common.serialization.StringDeserializer"); properties.put("enable.auto.commit", true); properties.put("group.id","KafkaStudy"); // 移除group.instance.id配置(可选,若不需要静态成员) KafkaConsumer<String, String> consumer = new KafkaConsumer<>(properties); consumer.subscribe(Collections.singleton("quickstart-events")); try { // 循环拉取消息,保持Consumer活跃 while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> record : records) { System.out.println(record.value()); } } } finally { consumer.close(); }
内容的提问来源于stack exchange,提问作者Kingfight98
相关产品推荐
相关产品推荐

