Kafka 3.6 Streams应用Broker宕机恢复后数据丢失问题咨询
环境与应用配置
我们运行基于Kafka 3.6的KStream应用,处理流程如下:
- 从3个Topic消费数据
- 经过重分区(repartitioning)、
selectKey及reduce操作生成KTable - 将3个KTable执行
LeftJoin后输出至目标Kafka Topic
集群与应用核心配置:
- 默认启用Exactly-Once Semantics (EOS/AOS)
- 集群配置:副本因子(RF)=3、同步副本(ISR)=3、最小同步副本(MISR)=2
- 状态存储使用持久化存储
- 应用具体配置:
settings.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); settings.setProperty(AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, System.getenv(SCHEMA_REGISTRY_URL) != null ? System.getenv(SCHEMA_REGISTRY_URL) : properties.getProperty(SCHEMA_REGISTRY_PROP_NAME)); settings.setProperty(StreamsConfig.STATE_DIR_CONFIG, properties.getProperty(STATESTORE_DIR)); settings.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3); settings.put(StreamsConfig.STATESTORE_CACHE_MAX_BYTES_CONFIG, 30 * 1024 * 1024L); settings.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); settings.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, GenericAvroSerde.class); settings.put(ProducerConfig.MAX_REQUEST_SIZE_CONFIG, 20000000); settings.put(ProducerConfig.ACKS_CONFIG, "all"); settings.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, 20000000); settings.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 18000000); settings.put(ProducerConfig.MAX_BLOCK_MS, 18000000); settings.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "zstd");
问题场景
应用部署在3台Broker的集群上,某次所有Broker宕机约20分钟,恢复后预期KStream会处理宕机期间堆积的消息并正常运行,但实际出现数据丢失。
在添加delivery.timeout.ms及max.block.ms超时配置前,曾触发以下异常:
[ERROR] [xxxx-yyyyyy-xxxx-9d5de3f8-b3d6-4dd6-bc4c-cf8320df3bb2-StreamThread-1] [org.apache.kafka.streams.processor.internals.RecordCollectorImpl.recordSendError(RecordCollectorImpl.java:322)]: stream-thread [xxxx-yyyyyy-xxxx-9d5de3f8-b3d6-4dd6-bc4c-cf8320df3bb2-StreamThread-1] stream-task [1_3] Error encountered sending record to topic xxxx-yyyyyy-xxxx-KSTREAM-REDUCE-STATE-STORE-0000000011-changelog for task 1_3 due to: org.apache.kafka.clients.producer.BufferExhaustedException: Failed to allocate 16384 bytes within the configured max blocking time 60000 ms. Total memory: 33554432 bytes. Available memory: 0 bytes. Poolable size: 16384 bytes The broker is either slow or in bad state (like not having enough replicas) in responding the request, or the connection to broker was interrupted sending the request or receiving the response. Consider overwriting `max.block.ms` and /or `delivery.timeout.ms` to a larger value to wait longer for such scenarios and avoid timeout errors org.apache.kafka.clients.producer.BufferExhaustedException: Failed to allocate 16384 bytes within the configured max blocking time 60000 ms. Total memory: 33554432 bytes. Available memory: 0 bytes. Poolable size: 16384 bytes [ERROR] [xxxx-yyyyyy-xxxx-9d5de3f8-b3d6-4dd6-bc4c-cf8320df3bb2-StreamThread-1] [org.apache.kafka.streams.processor.internals.ProcessorStateManager.flushCache(ProcessorStateManager.java:526)]: stream-thread [xxxx-yyyyyy-xxxx-9d5de3f8-b3d6-4dd6-bc4c-cf8320df3bb2-StreamThread-1] stream-task [1_3] Failed to flush cache of store KSTREAM-REDUCE-STATE-STORE-0000000011: org.apache.kafka.streams.errors.TaskCorruptedException: Tasks [1_3] are corrupted and hence need to be re-initialized
目前应用已无报错,但数据丢失问题仍存在。
问题分析与解答
一、当前场景下数据丢失的可能原因
任务损坏后的状态恢复不完整
之前触发的TaskCorruptedException会强制任务重新初始化,此时Kafka Streams会尝试从changelog主题恢复状态。如果changelog主题在Broker宕机前未完成ISR同步,恢复后部分副本数据缺失,会导致状态存储与消费偏移量不匹配,后续处理时跳过部分消息造成丢失。EOS事务提交的边界异常
Broker集体宕机时,应用可能处于事务提交的中间状态:已消费消息并更新状态,但事务未完成提交。恢复后EOS虽保证事务原子性,但如果本地状态快照与changelog偏移量不一致,会导致已处理消息被回滚,且后续消费不会重新处理这些消息,最终造成丢失。LeftJoin操作的状态不一致
三个KTable的状态存储若存在部分恢复不完整的情况,LeftJoin会因某一侧状态缺失生成错误输出——比如本该关联的数据因状态丢失未被关联,导致输出消息缺失。缓冲区溢出后的隐性数据丢失
调整超时配置前出现的BufferExhaustedException可能已导致部分消息被丢弃。即使后续配置解决了报错,之前丢失的消息也无法通过重启恢复。
二、启用AOS(Exactly-Once Semantics)时可能导致数据丢失的场景
状态存储与changelog双损
本地状态存储损坏,且对应的changelog主题数据丢失(如主题被删除、副本全部损坏),KStream无法恢复状态,只能按auto.offset.reset配置从头或最新位置重新消费,丢失之前的状态数据,导致聚合、Join等操作输出异常。事务与重试配置不当
transaction.timeout.ms过短,Broker处理事务提交缓慢时会触发事务中止,已处理消息被回滚但消费偏移量已提交,导致消息不会被重新处理。- Producer的
retries设为0或retry.backoff.ms不足,遇到临时Broker故障时直接丢弃消息,违反EOS语义。
集群配置不符合EOS要求
- 主题
min.insync.replicas设置低于2,且acks=all时,若ISR数量低于MISR,Producer无法发送消息,若应用未正确处理该异常会导致消息丢失。 - 事务协调器故障且未正确选举新节点,会导致事务无法提交,部分消息永久丢失。
- 主题
应用重启的状态恢复错误
- 本地状态快照与changelog偏移量不匹配时,KStream会重新初始化状态,若
auto.offset.reset设为latest,会从最新偏移量开始消费,丢失中间消息。 - 分布式部署下部分任务恢复状态失败,导致整体处理结果不一致,出现数据丢失。
- 本地状态快照与changelog偏移量不匹配时,KStream会重新初始化状态,若
自定义处理器的异常处理漏洞
自定义Processor、Transformer中若直接捕获异常并跳过消息,未触发任务重启或状态回滚,即使启用EOS也会造成消息丢失。
内容的提问来源于stack exchange,提问作者Justin

