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

Kafka Streams中如何获取split()与branch()后的独立下游流?

Kafka Streams 获取 split+branch 后的独立下游流

在Kafka Streams中,原有的branch()方法已被弃用,官方推荐使用split()配合branch()来实现流的分支逻辑。你给出的示例代码已经完成了分支规则定义,要拿到每个分支的独立下游流,只需保存split()返回的StreamSplit对象,再通过分支名称获取即可。

具体实现步骤

  1. 执行split()和branch()操作时,保存返回的StreamSplit对象
  2. 调用StreamSplit.stream("分支名称")方法,获取对应分支的独立KStream

完整代码示例

// 执行分支操作并保存StreamSplit对象
StreamSplit<String, ArtikelEvent> splitResult = artikelEvents.split()
        .branch((key, val) -> val.getKategorie().equals("Computer & Büro"), Branched.as("computer"))
        .branch((key, val) -> val.getKategorie().equals("Smartphones"), Branched.as("smartphones"))
        .branch((key, val) -> val.getKategorie().equals("Fotografie"), Branched.as("fotografie"));

// 获取各个分支的独立流
KStream<String, ArtikelEvent> computerStream = splitResult.stream("computer");
KStream<String, ArtikelEvent> smartphonesStream = splitResult.stream("smartphones");
KStream<String, ArtikelEvent> fotografieStream = splitResult.stream("fotografie");

// 后续可对每个分支单独处理,比如输出到Kafka主题
computerStream.to("computer-topic");
smartphonesStream.to("smartphones-topic");
fotografieStream.to("fotografie-topic");

补充说明

  • 若未给分支命名(不使用Branched.as()),也可通过索引(从0开始)调用splitResult.stream(int index)获取分支流,但命名方式更直观、便于维护,推荐优先使用。
  • 分支的过滤条件支持互斥或重叠逻辑,Kafka Streams会按照branch()调用的顺序依次匹配记录,匹配成功的记录会进入对应分支。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 22:41:16