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

KafkaStreams批量消费超5万条消息时停止拉取Topic消息问题排查

Kafka Streams批量消费停滞问题排查

问题背景

过去两周持续排查该问题暂未解决,相关现状如下:

  • 应用为单Pod部署,分配1GB内存,共运行6个Kafka Streams实例,其中2个Streams实例用于处理批量负载(单条消息大小约10KB,总消息量5-10万条)
  • 业务侧NiFi工作流会在3-4分钟内向对应Topic推送5万-11万条消息,通过自定义BulkStreams组件消费时,每次仅能读取2万-3万条记录,之后Streams就会停止从Topic拉取消息
  • Kafka运维团队确认:所有消息均被当前使用的消费组标记为已消费,该消费组内无其他消费者,消费期间也无其他消费组消费该Topic

现有依赖、代码与配置

初始依赖版本

springBootVersion=2.5.12
springCloudVersion=3.0.3
// Kafka
implementation group: "org.springframework.kafka", name: "spring-kafka", version: "2.8.2"
implementation group: "org.apache.kafka", name: "kafka-streams", version: "2.7.1"
implementation group: "org.apache.kafka", name: "kafka-clients", version: "2.7.1"
implementation group: "org.apache.kafka", name: "kafka-streams-test-utils", version: "2.7.1"

BulkStreams核心代码

@Component
@Slf4j
public class BulkStreams extends ConsumerStreams {
    @Value("${kafka.topic}")
    private String kafkaTopic;
    @Value("${application.env}")
    private String environment;
    @Value("${applicationId}")
    private String applicationId;
    @Value("${brokers}")
    private String kafkaServers;
    @Value("${javax.net.ssl.trustStore}")
    private String sslTrustStore;
    @Value("${javax.net.ssl.trustStorePassword}")
    private String sslTrustStorePassword;
    @Value("${connect.ssl.keystore.location}")
    private String sslKeystoreLocation;
    @Value("${connect.ssl.keystore.password}")
    private String sslKeystorePassword;
    @Value("${connect.ssl.key.password}")
    private String sslKeyPassword;
    @Value("${ssl.protocol}")
    private String sslProtocol;
    @Value("${security.protocol}")
    private String securityProtocol;

    @PostConstruct
    public void start() {
        KafkaStreams streams = new KafkaStreams(
                getStreamsBuilder().build(),
                getKafkaConsumerBuilderProperties(
                        applicationId,
                        kafkaServers,
                        "30000",
                        "0",
                        securityProtocol,
                        sslProtocol,
                        sslTrustStore,
                        sslTrustStorePassword,
                        sslKeystoreLocation,
                        sslKeystorePassword,
                        sslKeyPassword,
                        environment));
        try {
            streams.start();
            log.info("Started streams for topic {}", kafkaTopic);
        } catch (Throwable t) {
            log.error("Failed to start streams instance on topic {}. Details: {}", kafkaTopic, t);
        }
    
        @Override
        public void close() {
          }
    }

    public StreamsBuilder getStreamsBuilder() {

        StreamsBuilder builder = new StreamsBuilder();
        final KStream<String, String> source = builder.stream(kafkaTopic);

        source.transform(() -> new Transformer<String, String, KeyValue<Object, Object>>() {
            private ProcessorContext context;

            @Override
            public void init(ProcessorContext context) {
                this.context = context;
            }

            @Override
            public KeyValue<Object, Object> transform(String key, String value) {

                String traceId = String.format("%s | %s", kafkaTopic, LocalDateTime.now());
                try {
                    log.info("message from kafkaTopic= {} for feeds, key= {}, value= {}", kafkaTopic, key, value);

                    //logic to persist data in DB
    
                    return new KeyValue<>(key, value);
                } catch (Exception exception) {
                    log.error("Exception in reading feed: {}, key={}, value: {}", exception, key, value);
                }

                return null;
            }

            @Override
            public void close() {

            }
        });

        return builder;
    }
}

消费者基础配置类代码

@Component
@Slf4j
public class ConsumerStreams {

    @Autowired
    MetricRegistration metricRegistration;

    Properties getKafkaConsumerBuilderProperties(String applicationId,
                                                 String kafkaServers,
                                                 String commitIntervalMsConfig,
                                                 String cacheMaxBytesBuffering,
                                                 String securityProtocol,
                                                 String sslProtocol,
                                                 String sslTrustStore,
                                                 String sslTrustStorePassword,
                                                 String sslKeystoreLocation,
                                                 String sslKeystorePassword,
                                                 String sslKeyPassword,
                                                 String environment) {
        final Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, applicationId);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, applicationId);
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaServers);
        props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, commitIntervalMsConfig); //30000
        props.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, cacheMaxBytesBuffering); //0
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
        //new
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 100);
        props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 500000);
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);

        if (!(environment.equals("local") || environment.equals("dev") || environment.equals("sft"))) {
            props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, securityProtocol);
            props.put(SslConfigs.SSL_PROTOCOL_CONFIG, sslProtocol);
            props.put(SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG, sslTrustStore);
            props.put(SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG, sslTrustStorePassword);
            props.put(SslConfigs.SSL_KEYSTORE_LOCATION_CONFIG, sslKeystoreLocation);
            props.put(SslConfigs.SSL_KEYSTORE_PASSWORD_CONFIG, sslKeystorePassword);
            props.put(SslConfigs.SSL_KEY_PASSWORD_CONFIG, sslKeyPassword);
        }

        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        return props;
    }

    protected void logStream(String key, String value, String topic, MetricNames metric){
        log.debug("message from kafkaTopic={}, key={}, message={} ",
                  topic,
                  key,
                  value);
        metricRegistration.incrementCounter(metric.getName() + Constants.MESSAGE_CONSUMED_COUNTER);
    };

}

已尝试的排查方案

  • 升级Spring与Kafka相关依赖版本,版本配置如下:
springBootVersion=2.6.6
springCloudVersion=3.1.0
// Kafka
implementation group: "org.springframework.kafka", name: "spring-kafka", version: "2.8.7"
implementation group: "org.apache.kafka", name: "kafka-streams", version: "3.0.0"
implementation group: "org.apache.kafka", name: "kafka-clients", version: "3.0.0"
implementation group: "org.apache.kafka", name: "kafka-streams-test-utils", version: "3.0.0"
  • 将COMMIT_INTERVAL_MS_CONFIG降低至1000ms(默认值为3000ms)
  • 将COMMIT_INTERVAL_MS_CONFIG设为1000ms,同时将CACHE_MAX_BYTES_BUFFERING_CONFIG从0调整为10MB至100MB区间
  • 回退COMMIT_INTERVAL_MS_CONFIG为3000ms、CACHE_MAX_BYTES_BUFFERING_CONFIG为0后,新增/调整如下消费参数:
//new
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 300);
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 500000);
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);

参考Kafka重平衡问题相关排查文档,将MAX_POLL_RECORDS_CONFIG分别设置为100、200、300测试,仅将消费停滞的阈值从2万条提升至3万条,未解决根本问题。

测试观测现象

  • 日志无任何报错,Streams实例启动正常,当Topic总消息量低于5万条时可正常消费全量消息
  • 消费停滞前的消费时长约3-4分钟,与NiFi生产消息的时长基本一致
  • 消费过程中Pod CPU使用率为40%-50%,停滞后CPU使用率降至20%
  • 内存无明显尖刺,整体内存使用率稳定在700-800MB区间

根因定位

代码与配置存在5个核心问题,直接导致消费停滞:

  1. 参数硬编码导致调整不生效:BulkStreams.start()方法调用getKafkaConsumerBuilderProperties时,硬编码传入commitIntervalMsConfig="30000"、cacheMaxBytesBuffering="0",之前对提交间隔、缓存大小的所有调整对该批量消费实例完全无效。
  2. 实例无强引用被GC回收:KafkaStreams实例被定义为start()方法的局部变量,方法执行完成后实例无GC Root关联,运行3-4分钟触发GC时实例被回收,KafkaStreams的finalize逻辑会自动关闭所有消费线程,和观测到的停滞时间完全吻合。
  3. 拓扑结构不完整:调用source.transform()后未接收返回值,也未添加任何终端算子(如foreach、to)。Kafka Streams采用反向追溯的方式构建拓扑,无下游sink的处理器会被优化,且transform算子本身设计为向下游转发消息,无输出路径时会导致消费流程异常。同时自动提交机制会提前提交已poll但未处理的消息offset,对应运维观测到的“所有消息被标记为已消费”现象。
  4. 方法定义位置错误:close()方法被错误嵌套在start()方法内部,Java不支持方法嵌套,该代码无法通过编译,即使为粘贴格式问题,现有逻辑也未实现Spring容器销毁时的Streams优雅关闭逻辑。
  5. 日志IO阻塞消费线程:transform方法中用info级别打印单条10KB的全量消息内容,高吞吐下日志同步写入会阻塞消费线程,同时占用大量堆内存触发频繁GC。

修复方案

  • 修复参数硬编码问题:将提交间隔、缓存大小等参数改为配置项注入,调用getKafkaConsumerBuilderProperties时传入配置值,确保参数调整生效。
  • 修复实例引用问题:将KafkaStreams streams定义为BulkStreams类的成员变量,保证应用运行期间实例被强引用持有,不会被GC回收;将close()方法移到start()方法外部作为类成员方法,实现DisposableBean接口或添加@PreDestroy注解,在Bean销毁时调用streams.close()优雅关闭实例。
  • 修复拓扑结构:如果不需要将处理后的消息发往下游,将transform算子替换为无返回值的终端算子process();如果保留transform,需要接收返回的KStream实例,添加to()(发往下游Topic)或foreach()(终端处理)算子保证拓扑连通;处理完消息后可按需调用context.commit()手动提交offset,避免offset超前提交。
  • 优化日志配置:将打印全量消息value的日志级别从info调整为debug,生产环境关闭debug日志,减少日志IO开销和内存占用。
  • 补充性能配置:添加StreamsConfig.NUM_STREAM_THREADS_CONFIG配置,根据负载设置为2-3提升消费吞吐;适当调大ConsumerConfig.FETCH_MAX_BYTES_CONFIG、ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG,适配10KB单条消息的拉取场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 01:06:29