【问题标题】:KTable aggregate forwards the same messagesKTable聚合转发相同的消息
【发布时间】:2020-02-14 18:37:45
【问题描述】:

我正在使用 kafka-streams 将消息聚合到 KTable 中。在我的聚合逻辑中,我总是返回相同的累加器,如下所示:

  streamOfInts
    .groupByKey()
    .aggregate(Accumulator.empty()) {k,v,acc -> acc}
    .toStream()
    .to(...)

我的期望是 - 由于 KTable 的值没有改变 - 不会向下游发送任何值。然而,这种情况并非如此。聚合函数总是转发更新。

确保产生相同(或相等)值的更新不会导致下游转发的最佳方法是什么?

【问题讨论】:

    标签: apache-kafka-streams


    【解决方案1】:

    DSL 运营商在设计 atm 时发出“更新时”而不是“更改时”。有一张 JIRA 票证建议添加“更改时发出”语义 (https://issues.apache.org/jira/browse/KAFKA-8770)。

    作为一种解决方法,您可以使用状态存储实现自定义 transform() -- 对于每个输入记录,您检查存储是新的(-> 发出并放入存储)还是更改(-> 发出和更新商店)。如果它存在并且没有改变,就不要发射任何东西。

    【讨论】:

    • 好的,谢谢,这就是我最终得到的结果:还有一个注意事项:我必须将 ValueTransformerWithKey 与随后的 filter 结合起来,以过滤掉 @987654325 返回的空值@方法。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-03-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多