You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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逻辑:

  1. 将两个流分别存入你自定义的Redis窗口状态存储;
  2. 对其中一个流的每个事件,去另一个Redis存储中查询窗口内匹配的记录;
  3. 自行处理左Join的核心逻辑(即主流事件即使没有匹配记录也要保留);
  4. 务必实现窗口过期清理逻辑,避免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的内部类实现:

  1. 定位到创建第三个共享存储的核心逻辑;
  2. 扩展KStreamImplJoin类,替换存储创建逻辑,改用你的Redis存储供应商;
  3. 注意这种方法依赖Kafka Streams内部API,版本升级时需要适配。

方案3:配合状态存储扩展机制强化自定义存储

确认你的Redis状态存储完整实现了WindowStoreSupplier接口,并在Streams配置中正确注册自定义存储工厂。不过该方法无法直接解决内置Join自动创建存储的问题,需配合方案1使用,确保所有状态操作都指向Redis存储。


内容的提问来源于stack exchange,提问作者marhewa

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 05:25:06