【发布时间】:2020-05-24 06:16:44
【问题描述】:
我有一个 Kafka 流,其中包含以下操作: stream.mapValues.groupByKey.aggregate。聚合基本上是将记录添加到列表中。 现在我将实现更改为: stream.flatMap.groupByKey.aggregate。 FlatMap 正在复制记录:第一条记录与旧实现中的完全相同,第二条记录已更改。因此,在实现重新分区发生变化之后,而在之前,它不是(这很好)。我的问题是,在发布更改后,旧密钥的旧聚合记录消失了。从改变的那一刻起,一切都按原样工作,但我不明白这种行为。据我了解,由于我没有更改密钥,因此它应该与以前位于同一分区上,并且聚合应该继续将消息添加到旧列表中,而不是从头开始。谁能帮我理解为什么会这样?
【问题讨论】:
-
欢迎来到stackoverflow。 stackoverflow.com/help/minimal-reproducible-example
标签: java apache-kafka apache-kafka-streams