Kafka Streams表-表内连接产生重复消息问题咨询
Kafka Streams KTable内连接重复输出问题
背景与问题
我需要用Kafka Streams实现两个Kafka主题的连接:
- 最初尝试用KStream(左表)和KTable(右表)做
leftJoin,但频繁出现右值为空的异常情况;窗口连接也不适用我的业务场景。 - 改为将两个主题都建模为KTable做内连接后,发现当向两个输入主题各发送1条同键消息时,结果流会出现2条相同的连接结果,触发两次连接逻辑。
测试与运行差异
- 用
TestTopologyDriver编写单元测试时结果正常,每对同键消息对应1条输出; - 但用
TestCompanion搭配测试代理,或者直接运行应用时,就会出现重复输出。
环境与配置
- 运行环境:Quarkus 2.16.x,Kafka Streams 3.3.2
- 核心配置:
kafka-streams.cache.max.bytes.buffering=10240 kafka-streams.commit.interval.ms=500 kafka-streams.metadata.max.age.ms=500 kafka-streams.auto.offset.reset=earliest kafka-streams.consumer.heartbeat.interval.ms=200
拓扑结构
子拓扑1用于测试时将单条输入消息转发到两个输入主题;子拓扑2实现两个KTable的内连接(外连接过滤空值后也存在同样重复问题):
Sub-topology: 1 Source: KSTREAM-SOURCE-0000000041 (topics: [trigger]) --> KSTREAM-FILTER-0000000044, KSTREAM-FILTER-0000000042 Processor: KSTREAM-FILTER-0000000044 (stores: []) --> KSTREAM-MAPVALUES-0000000045 <-- KSTREAM-SOURCE-0000000041 Processor: KSTREAM-FILTER-0000000042 (stores: []) --> KSTREAM-SINK-0000000043 <-- KSTREAM-SOURCE-0000000041 Processor: KSTREAM-MAPVALUES-0000000045 (stores: []) --> KSTREAM-SINK-0000000046 <-- KSTREAM-FILTER-0000000044 Sink: KSTREAM-SINK-0000000043 (topic: value) <-- KSTREAM-FILTER-0000000042 Sink: KSTREAM-SINK-0000000046 (topic: state) <-- KSTREAM-MAPVALUES-0000000045 Sub-topology: 2 Source: KSTREAM-SOURCE-0000000047 (topics: [value]) --> KSTREAM-TOTABLE-0000000048 Source: KSTREAM-SOURCE-0000000051 (topics: [state]) --> KTABLE-SOURCE-0000000052 Processor: KSTREAM-TOTABLE-0000000048 (stores: [KSTREAM-TOTABLE-STATE-STORE-0000000049]) --> KTABLE-JOINOTHER-0000000055 <-- KSTREAM-SOURCE-0000000047 Processor: KTABLE-SOURCE-0000000052 (stores: [status-STATE-STORE-0000000050]) --> KTABLE-JOINTHIS-0000000054 <-- KSTREAM-SOURCE-0000000051 Processor: KTABLE-JOINOTHER-0000000055 (stores: [status-STATE-STORE-0000000050]) --> KTABLE-MERGE-0000000053 <-- KSTREAM-TOTABLE-0000000048 Processor: KTABLE-JOINTHIS-0000000054 (stores: [KSTREAM-TOTABLE-STATE-STORE-0000000049]) --> KTABLE-MERGE-0000000053 <-- KTABLE-SOURCE-0000000052 Processor: KTABLE-MERGE-0000000053 (stores: []) --> KTABLE-TOSTREAM-0000000056 <-- KTABLE-JOINTHIS-0000000054, KTABLE-JOINOTHER-0000000055 Processor: KTABLE-TOSTREAM-0000000056 (stores: []) --> KSTREAM-MAPVALUES-0000000057 <-- KTABLE-MERGE-0000000053 Processor: KSTREAM-MAPVALUES-0000000057 (stores: []) --> KSTREAM-SINK-0000000058 <-- KTABLE-TOSTREAM-0000000056 Sink: KSTREAM-SINK-0000000058 (topic: joined) <-- KSTREAM-MAPVALUES-0000000057
补充说明
- 两个输入主题的同键消息到达时间间隔极短(仅数毫秒);
- 看到有建议在结果流上加去重处理器,但需要维护第三个状态存储,感觉不合理;
- 未在官方文档中找到该重复触发问题的原因。
问题分析与解决方案
原因定位
这种重复输出是KTable连接的正常行为特性:
- KTable依赖状态存储工作,当第一个KTable的同键消息到达时,若第二个消息还未写入对应状态存储,不会输出连接结果;
- 当第二个KTable的同键消息到达时,会检测到第一个KTable状态存储中的匹配记录,触发连接并输出一条结果;
- 由于两条消息到达间隔极短,结合当前缓存与提交配置,第一个KTable的状态更新可能在第二个消息处理完成后,再次触发连接处理器重新计算,从而输出第二条重复结果。
解决方案
调整缓存与提交配置:
- 调大
kafka-streams.cache.max.bytes.buffering(比如设为10485760即10MB),让更多状态更新在内存缓存中合并,减少下游处理器的触发次数; - 适当增大
kafka-streams.commit.interval.ms(比如设为1000),让状态存储的提交更平缓,避免短时间内多次触发连接计算。
- 调大
使用
suppress()操作去重:- 在KTable转流之后,用
suppress()操作对同键结果进行窗口去重,无需额外维护状态存储(底层复用现有状态):joinedKTable.toStream() .suppress(Suppressed.untilTimeLimit(Duration.ofMillis(500), Suppressed.BufferConfig.unbounded().withKeySerde(keySerde))) .mapValues(...) .to("joined"); - 该操作会在指定时间窗口内只保留同键的最后一条结果,过滤重复输出,时间窗口可根据消息到达间隔调整。
- 在KTable转流之后,用
优化消息发送顺序:
- 如果业务允许,确保两个主题的同键消息按顺序发送,或让其中一个主题的消息延迟发送,避免短时间内两条同键消息同时触发状态更新。
内容的提问来源于stack exchange,提问作者user3904687
相关产品推荐
相关产品推荐

