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

