You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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("已提交当前批次消息的偏移量");
    }
}

验证步骤

  1. 先修改group.id为新值,运行代码测试是否能读取到消息
  2. 验证修改后的代码是否能正常退出,不再出现“卡住”的情况
  3. 确认偏移量提交逻辑生效,重启消费者后不会重复读取相同消息

内容的提问来源于stack exchange,提问作者rphardu

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 08:25:24