Kafka Streams中如何获取split()与branch()后的独立下游流?
Kafka Streams 获取 split+branch 后的独立下游流
在Kafka Streams中,原有的branch()方法已被弃用,官方推荐使用split()配合branch()来实现流的分支逻辑。你给出的示例代码已经完成了分支规则定义,要拿到每个分支的独立下游流,只需保存split()返回的StreamSplit对象,再通过分支名称获取即可。
具体实现步骤
- 执行
split()和branch()操作时,保存返回的StreamSplit对象 - 调用
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
相关产品推荐
相关产品推荐

