Kafka Stream同谓词多分支消息分发问题及解决方案
Kafka Streams:让两个分支接收同一流的全部消息
我之前也踩过这个坑!你想要从单个KStream读取相同的<key,value>消息,然后分发到两个不同分支,而且两个分支都要接收全部消息,但用branch()方法时却只有第一个分支能收到消息,第二个分支完全没数据,对吧?
你尝试的问题代码
KStream<String, byte[]>[] branches = builder.<String, byte[]>stream("source-topic") .branch((key, val) -> true, // 发往topic1 (key, val) -> true); // 发往topic2 // branch[0]经logic1OnData处理后发送至topic1 branches[0].map(logic1OnData).filter( (key, value) -> { if (key == null || value == null) return false; return value.data() != null; }).to("topic1", Produced.with(Serdes.String(), Serdes.String())); // branch[1]经logic2OnData处理后发送至topic2 branches[1].map(logic2OnData).filter( (key, value) -> { if (key == null || value == null) return false; return value.data() != null; }).to("topic2", Produced.with(Serdes.String(), Serdes.String()));
问题原因
branch()方法的底层逻辑是if-else式的匹配规则:消息会按顺序匹配分支的谓词条件,只要被第一个匹配的分支接收,后续的分支就不会再处理这条消息。哪怕你两个分支的谓词都是(key, val) -> true,第一条匹配后就直接进入第一个分支,第二个分支根本拿不到消息。
正确解决方案
不要用branch()拆分,直接基于原始的KStream分别构建处理链路就可以了。Kafka Streams的流是可复用的,同一个原始流可以被多次处理,每条消息都会流经每个独立的处理链路:
KStream<String, byte[]> stream = builder.<String, byte[]>stream("source-topic"); // 第一条处理链路:发往topic1 stream.map(logic1OnData).filter( (key, value) -> { if (key == null || value == null) return false; return value.data() != null; }).to("topic1", Produced.with(Serdes.String(), Serdes.String())); // 第二条处理链路:发往topic2 stream.map(logic2OnData).filter( (key, value) -> { if (key == null || value == null) return false; return value.data() != null; }).to("topic2", Produced.with(Serdes.String(), Serdes.String()));
修改后,每条从source-topic读取的消息都会被两个map操作分别处理,经过过滤后写入对应的目标topic,完全符合你预期的两个分支都获取相同消息的需求。
内容的提问来源于stack exchange,提问作者Sudarshan sridhar
相关产品推荐
相关产品推荐

