Kafka Streams leftJoin中自定义Redis状态存储的适配难题
问题描述
我需要在Kafka Streams应用中使用基于Redis的自定义状态存储,但在执行leftJoin操作时遇到了阻碍。
我的代码如下:
val thisStoreSupplier = redisStores.joinWindowStoreSupplier("join-this", window) val otherStoreSupplier = redisStores.joinWindowStoreSupplier("join-other", window) val joined = stream1 .leftJoin( stream2, { foo: Foo, bar: Bar? -> Joined(foo, bar) }, window, StreamJoined .with<String?, Foo?, Bar?>(thisStoreSupplier, otherStoreSupplier) .withKeySerde(Serdes.String()) .withName("my-join") )
理论上这段代码应该正常运行,但查看Kafka Streams内部逻辑后发现,底层会自动创建第三个状态存储,而构造这个存储的逻辑不支持自定义供应商,只能创建内存或RocksDB存储。当前应用处理的事件量极大,内存和RocksDB都无法满足需求,请问如何将这个自动创建的存储也替换为基于Redis的自定义存储?
补充的实际拓扑结构(真实域名已替换为foo/bar):
Topologies: Sub-topology: 0 Source: bar-source (topics: [bar.topic]) --> bar-deserialize Processor: bar-deserialize (stores: []) --> bar-validation-filter <-- bar-source Processor: bar-validation-filter (stores: []) --> KSTREAM-PEEK-0000000003 <-- bar-deserialize Processor: KSTREAM-PEEK-0000000003 (stores: []) --> bar-convert <-- bar-validation-filter Processor: bar-convert (stores: []) --> bar-filter <-- KSTREAM-PEEK-0000000003 Processor: bar-filter (stores: []) --> bar-selectKey <-- bar-convert Processor: bar-selectKey (stores: []) --> KSTREAM-PEEK-0000000007 <-- bar-filter Processor: KSTREAM-PEEK-0000000007 (stores: []) --> bar-repartition-filter <-- bar-selectKey Processor: bar-repartition-filter (stores: []) --> bar-repartition-sink <-- KSTREAM-PEEK-0000000007 Sink: bar-repartition-sink (topic: bar-repartition) <-- bar-repartition-filter Sub-topology: 1 Source: foo-repartition-source (topics: [foo-repartition]) --> KSTREAM-PEEK-0000000022 Source: bar-repartition-source (topics: [bar-repartition]) --> KSTREAM-PEEK-0000000011 Processor: KSTREAM-PEEK-0000000011 (stores: []) --> foo-bar-join-other-windowed <-- bar-repartition-source Processor: KSTREAM-PEEK-0000000022 (stores: []) --> foo-bar-join-this-windowed <-- foo-repartition-source Processor: foo-bar-join-other-windowed (stores: [KSTREAM-OUTEROTHER-0000000026-store]) --> foo-bar-join-outer-other-join <-- KSTREAM-PEEK-0000000011 Processor: foo-bar-join-this-windowed (stores: [KSTREAM-JOINTHIS-0000000025-store]) --> foo-bar-join-this-join <-- KSTREAM-PEEK-0000000022 Processor: foo-bar-join-outer-other-join (stores: [KSTREAM-JOINTHIS-0000000025-store, KSTREAM-OUTERSHARED-0000000025-store]) --> foo-bar-join-merge <-- foo-bar-join-other-windowed Processor: foo-bar-join-this-join (stores: [KSTREAM-OUTEROTHER-0000000026-store, KSTREAM-OUTERSHARED-0000000025-store]) --> foo-bar-join-merge <-- foo-bar-join-this-windowed Processor: foo-bar-join-merge (stores: []) --> foo-bar-transformer <-- foo-bar-join-this-join, foo-bar-join-outer-other-join Processor: foo-bar-transformer (stores: [foo-bar-store]) --> combined-non-null-filter <-- foo-bar-join-merge Processor: combined-non-null-filter (stores: []) --> combined-metrics <-- foo-bar-transformer Processor: combined-metrics (stores: []) --> combined-sink <-- combined-non-null-filter Sink: combined-sink (topic: foobar.combined.temp) <-- combined-metrics Sub-topology: 2 Source: foo-source (topics: [foo.topic]) --> foo-deserialize Processor: foo-deserialize (stores: []) --> foo-validation-filter <-- foo-source Processor: foo-validation-filter (stores: []) --> KSTREAM-PEEK-0000000015 <-- foo-deserialize Processor: KSTREAM-PEEK-0000000015 (stores: []) --> foo-convert <-- foo-validation-filter Processor: foo-convert (stores: []) --> foo-selectKey <-- KSTREAM-PEEK-0000000015 Processor: foo-selectKey (stores: []) --> KSTREAM-PEEK-0000000018 <-- foo-convert Processor: KSTREAM-PEEK-0000000018 (stores: []) --> foo-repartition-filter <-- foo-selectKey Processor: foo-repartition-filter (stores: []) --> foo-repartition-sink <-- KSTREAM-PEEK-0000000018 Sink: foo-repartition-sink (topic: foo-repartition) <-- foo-repartition-filter
解决方案
方案1:手动实现Left Join逻辑,完全控制状态存储
既然Kafka Streams内置的leftJoin会自动创建不受控的状态存储,最直接的办法是绕过该API,自行实现Join逻辑:
- 将两个流分别存入你自定义的Redis窗口状态存储;
- 对其中一个流的每个事件,去另一个Redis存储中查询窗口内匹配的记录;
- 自行处理左Join的核心逻辑(即主流事件即使没有匹配记录也要保留);
- 务必实现窗口过期清理逻辑,避免Redis中数据无限累积。
示例代码思路:
// 1. 将stream1写入Redis窗口存储 stream1.process({ object : Processor<String, Foo> { lateinit var store: WindowStore<String, Foo> lateinit var context: ProcessorContext override fun init(context: ProcessorContext) { this.context = context store = context.getStateStore("join-this") as WindowStore<String, Foo> } override fun process(key: String, value: Foo) { store.put(key, value, context.timestamp()) } override fun close() {} } }, "join-this") // 2. 处理stream2,查询stream1的存储并输出匹配结果 val matchedStream = stream2.transform({ object : Transformer<String, Bar, KeyValue<String, Joined>> { lateinit var thisStore: WindowStore<String, Foo> lateinit var otherStore: WindowStore<String, Bar> lateinit var context: ProcessorContext override fun init(context: ProcessorContext) { this.context = context thisStore = context.getStateStore("join-this") as WindowStore<String, Foo> otherStore = context.getStateStore("join-other") as WindowStore<String, Bar> } override fun transform(key: String, value: Bar): KeyValue<String, Joined>? { otherStore.put(key, value, context.timestamp()) val windowStart = context.timestamp() - window.size() val fooIter = thisStore.fetch(key, windowStart, context.timestamp()) return fooIter.takeIf { it.hasNext() }?.let { KeyValue(key, Joined(it.next().value, value)) } } override fun close() {} } }, "join-this", "join-other") // 3. 处理stream1中无匹配的事件,单独输出 val unmatchedStream = stream1.transform({ object : Transformer<String, Foo, KeyValue<String, Joined>> { lateinit var otherStore: WindowStore<String, Bar> lateinit var context: ProcessorContext override fun init(context: ProcessorContext) { this.context = context otherStore = context.getStateStore("join-other") as WindowStore<String, Bar> } override fun transform(key: String, value: Foo): KeyValue<String, Joined>? { val windowStart = context.timestamp() - window.size() val barIter = otherStore.fetch(key, windowStart, context.timestamp()) return if (!barIter.hasNext()) { KeyValue(key, Joined(value, null)) } else { null // 已有匹配,不重复输出 } } override fun close() {} } }, "join-other") // 4. 合并两个流,得到完整的左Join结果 val finalJoined = matchedStream.merge(unmatchedStream)
方案2:扩展Kafka Streams内部Join实现
如果不想完全重写Join逻辑,可以通过修改或扩展Kafka Streams的内部类实现:
- 定位到创建第三个共享存储的核心逻辑;
- 扩展
KStreamImplJoin类,替换存储创建逻辑,改用你的Redis存储供应商; - 注意这种方法依赖Kafka Streams内部API,版本升级时需要适配。
方案3:配合状态存储扩展机制强化自定义存储
确认你的Redis状态存储完整实现了WindowStoreSupplier接口,并在Streams配置中正确注册自定义存储工厂。不过该方法无法直接解决内置Join自动创建存储的问题,需配合方案1使用,确保所有状态操作都指向Redis存储。
内容的提问来源于stack exchange,提问作者marhewa
相关产品推荐
相关产品推荐

