m17*_*1vw 5 java apache-kafka apache-kafka-streams
我有一个 Kafka 流,它从一个主题中获取数据,并且需要将该信息过滤到两个不同的主题。
KStream<String, Model> stream = builder.stream(Serdes.String(), specificAvroSerde, "not-filtered-topic");
stream.filter((key, value) -> new Processor().test(key, value)).to(Serdes.String(), specificAvroSerde, "good-topic");
stream.filterNot((key, value) -> new Processor().test(key, value)).to(Serdes.String(), specificAvroSerde, "bad-topic");
Run Code Online (Sandbox Code Playgroud)
但是,当我这样做时,它会从主题中读取数据两次——不确定随着数据变大,这是否对性能有任何影响。有没有办法只过滤一次并将其推送到两个主题?
您的方法是正确的,并且不会从主题中读取两次数据,也没有进行内部数据复制。您的方法的唯一缺点是,对每条记录都评估了两个过滤器谓词——但是,这非常便宜,不应该是性能问题。
但是,您仍然可以通过使用KStream#branch()它来提高性能,它确实采用多个谓词并依次评估所有谓词并为每个谓词返回一个输入流。如果记录与谓词匹配,则将其放入相应的输出流中并停止评估(即,不会为该单个记录评估进一步的谓词——这确保将每条记录添加到最多一个输出流;或者如果没有谓词匹配)。
因此,您可以只为 提供两个谓词branch():第一个与原始filter()谓词相同,第二个谓词始终返回true。
KStream<String, Model> stream = builder.stream(
Serdes.String(),
specificAvroSerde,
"not-filtered-topic"
);
KStream[] splitStreams = stream.branch(
(key, value) -> new Processor().test(key,value),
(key, value) -> true
);
splitStreams[0].to(Serdes.String(), specificAvroSerde, "good-topic");
splitStreams[1].to(Serdes.String(), specificAvroSerde, "bad-topic");
Run Code Online (Sandbox Code Playgroud)
不过,不确定此代码是否比原始版本更具可读性。我想这是一个品味问题,我个人更喜欢你的原始代码,因为它确实更好地表达了语义。
我添加的版本应该稍微提高 CPU 效率,因为对于满足谓词的所有记录,它只计算一次。并且对于所有不满足结果的记录,一个简单的true将被返回(即,没有第二个谓词评估)。
如果您知道大多数记录将在 中结束splitStream[1],您还可以反转谓词(并splitStream[0]用作“坏流”)以减少对第二个true返回谓词的调用次数。但这些只是微观优化,应该无关紧要。
| 归档时间: |
|
| 查看次数: |
2578 次 |
| 最近记录: |