【发布时间】:2020-02-14 18:37:45
【问题描述】:
我正在使用 kafka-streams 将消息聚合到 KTable 中。在我的聚合逻辑中,我总是返回相同的累加器,如下所示:
streamOfInts
.groupByKey()
.aggregate(Accumulator.empty()) {k,v,acc -> acc}
.toStream()
.to(...)
我的期望是 - 由于 KTable 的值没有改变 - 不会向下游发送任何值。然而,这种情况并非如此。聚合函数总是转发更新。
确保产生相同(或相等)值的更新不会导致下游转发的最佳方法是什么?
【问题讨论】: