Kafka Streams批量Join出现Skipping record for expired segment告警及数据丢失
问题原因
你的猜测完全符合Kafka Streams的底层逻辑:Kafka Streams的流时间默认由任务消费的所有分区的最大事件时间推进。当应用重启追积压、重处理历史数据时,如果某一个分区的消息先被消费到最新位点,流时间会被直接拉到该分区的最新事件时间,后续其他分区消费到的更早时间的事件,会被判定为超出窗口有效范围,直接跳过处理,就会触发你看到的Skipping record for expired segment告警,同时对应join逻辑无法命中,出现结果丢失。
解决方案
1. 配置max.task.idle.ms(最优方案,无额外存储开销)
Kafka Streams 2.4及以上版本提供了max.task.idle.ms配置项,作用是控制任务推进流时间前的等待时长:任务会等待所有分配到的分区都有数据拉取到之后,再更新流时间,避免单个分区的消息把流时间提前带飞。
你可以根据自己场景下积压数据的最大时间跨度设置该值,比如复现场景中两个topic间隔1小时,就可以设置为3600000(1小时),稳态运行时该配置不会产生额外性能开销,也不需要调整窗口的grace时间,完全匹配你低开销的要求。
配置示例:
props.put(StreamsConfig.MAX_TASK_IDLE_MS_CONFIG, 3600000);
2. 低版本适配方案
如果你使用的Kafka Streams版本低于2.4,没有上述配置,可以采用以下两种方案:
- 调整分区分配策略为
StickyAssignor,尽量保证同一个任务分配到的分区的时间跨度差异不会过大,降低流时间跳变的概率 - 重处理历史数据阶段临时调大窗口grace时间,等积压消费完成、应用进入稳态运行后,再把grace值改回原值滚动重启即可,不会长期占用额外存储资源
补充说明
Kafka Streams 3.0及以上版本中,窗口默认的grace时间为0,只要事件时间超出「窗口结束时间+grace时长」就会被直接丢弃,这也是你场景下问题触发的直接原因。使用max.task.idle.ms从流时间推进逻辑层面解决问题,比调整grace时间的资源开销要低得多。
内容的提问来源于stack exchange,提问作者Matt Sheppard

