Kafka Java Consumer在Jenkins挂起且本地无法读取消息求助
Kafka Java Consumer 无法读取消息问题排查
你的环境与代码情况
先整理下你的配置与代码细节:
依赖配置
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>1.1.0</version> </dependency>
消费者核心代码
private static String TopicName = "Automation_kafka_test"; LOGGER.info("Initializing the consumer"); KafkaConsumer<String, String> myKafkaCascadeConsumer = new KafkaConsumer<String, String>(KafkaCascadeConsumer.kafkaCascadeConfiguration()); for (Map.Entry<String, Object> entry : KafkaCascadeConsumer.kafkaCascadeConfiguration().entrySet()) { LOGGER.info("Key = "+entry.getKey() + ", Value =" + entry.getValue()); } KafkaConsumerHelper.readKafkaMessages(myKafkaCascadeConsumer, TopicName); myKafkaCascadeConsumer.close(); // 读取消息方法 public static void readKafkaMessages(KafkaConsumer<String, String> myKafkaConsumer, String topicName) { LOGGER.info("Subscribing to Topic =" + topicName); myKafkaConsumer.subscribe(Arrays.asList(topicName)); while (true) { ConsumerRecords<String, String> records = myKafkaConsumer.poll(100); for (ConsumerRecord<String, String> record : records) System.out.printf("offset = %d, key = %s, value = %s", record.offset(), record.key(), record.value()); } }
执行日志
2018-05-30 14:23:21,247 INFO [TestNG-test=Test-1] (US000000_KafkaTest.java:81) - Initializing the consumer 2018-05-30 14:23:21,869 INFO [TestNG-test=Test-1] (US000000_KafkaTest.java:87) - Key = key.deserializer, Value =org.apache.kafka.common.serialization.StringDeserializer 2018-05-30 14:23:21,869 INFO [TestNG-test=Test-1] (US000000_KafkaTest.java:87) - Key = value.deserializer, Value =org.apache.kafka.common.serialization.StringDeserializer 2018-05-30 14:23:21,869 INFO [TestNG-test=Test-1] (US000000_KafkaTest.java:87) - Key = enable.auto.commit, Value =false 2018-05-30 14:23:21,869 INFO [TestNG-test=Test-1] (US000000_KafkaTest.java:87) - Key = group.id, Value =AutomationRamtest1 2018-05-30 14:23:21,870 INFO [TestNG-test=Test-1] (US000000_KafkaTest.java:87) - Key = consumer.timeout.ms, Value =50000 2018-05-30 14:23:21,871 INFO [TestNG-test=Test-1] (US000000_KafkaTest.java:87) - Key = bootstrap.servers, Value =ABsrd00xxx:9092,ABsrd00yyy:9092 ***** masked for privacy*** 2018-05-30 14:23:21,871 INFO [TestNG-test=Test-1] (US000000_KafkaTest.java:87) - Key = auto.commit.interval.ms, Value =1000 2018-05-30 14:23:21,871 INFO [TestNG-test=Test-1] (US000000_KafkaTest.java:87) - Key = auto.offset.reset, Value =earliest 2018-05-30 14:23:21,887 INFO [TestNG-test=Test-1] (KafkaConsumerHelper.java:53) - Subscribing to Topic =Automation_kafka_test
问题分析与解决方案
结合你的情况(本地/Jenkins卡住、CLI能读、生产者正常),我帮你梳理几个核心问题点和修复方案:
无限循环导致的“卡住”假象
你的readKafkaMessages方法用了while(true)无限循环,哪怕没有消息也会一直执行poll(100)。Jenkins看起来卡在订阅步骤,其实是程序进入了无限循环,只是没有消息输出而已;本地运行也是同理,没有消息时不会有任何打印,容易让你误以为程序卡住了。消费者组偏移量的问题
你设置了auto.offset.reset=earliest,但如果消费者组AutomationRamtest1之前已经提交过偏移量,且偏移量已经到了分区末尾,那么即使主题里有旧消息,消费者也不会再读取。而CLI默认使用临时消费者组,所以能读到消息。缺失手动提交偏移量的逻辑
你设置了enable.auto.commit=false,但代码里完全没有手动提交偏移量的逻辑。这会导致:消费者重启后重复读取消息;如果是已有偏移量的消费者组,会一直停留在之前提交的偏移量位置,无法读取新消息。废弃参数的问题
你配置了consumer.timeout.ms=50000,但这个参数在Kafka 0.10.1.0之后已经被废弃,取而代之的是max.poll.interval.ms。虽然这不是当前问题的直接原因,但建议替换成新参数避免潜在问题。
修改后的示例代码
public static void readKafkaMessages(KafkaConsumer<String, String> myKafkaConsumer, String topicName) { LOGGER.info("Subscribing to Topic =" + topicName); myKafkaConsumer.subscribe(Arrays.asList(topicName)); int noMessageCount = 0; final int MAX_NO_MESSAGE_COUNT = 5; // 连续5次没消息就退出循环 while (true) { // Kafka 1.1.0开始推荐使用Duration参数的poll方法 ConsumerRecords<String, String> records = myKafkaConsumer.poll(Duration.ofMillis(100)); if (records.isEmpty()) { noMessageCount++; if (noMessageCount > MAX_NO_MESSAGE_COUNT) { LOGGER.info("连续{}次未获取到消息,退出循环", MAX_NO_MESSAGE_COUNT); break; } continue; } noMessageCount = 0; for (ConsumerRecord<String, String> record : records) { // 加上换行符,让输出更易读 System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value()); } // 手动提交偏移量,确保偏移量被持久化 myKafkaConsumer.commitSync(); LOGGER.info("已提交当前批次消息的偏移量"); } }
验证步骤
- 先修改
group.id为新值,运行代码测试是否能读取到消息 - 验证修改后的代码是否能正常退出,不再出现“卡住”的情况
- 确认偏移量提交逻辑生效,重启消费者后不会重复读取相同消息
内容的提问来源于stack exchange,提问作者rphardu
相关产品推荐
相关产品推荐

