【问题标题】:Loosing old aggregated records after adding repartitioning in Kafka在 Kafka 中添加分区后丢失旧的聚合记录
【发布时间】:2020-05-24 06:16:44
【问题描述】:

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

【问题讨论】:

标签: java apache-kafka apache-kafka-streams


【解决方案1】:

如果您更改处理拓扑,通常需要重置应用程序并重新处理输入主题中的所有数据以重新计算状态。

在您的情况下,我假设 aggregation 运算符在更改后具有不同的名称,因此不再“找到”其本地状态和更改日志主题。

您可以通过Topology#describe() 比较两种拓扑的名称。

为了顺利升级,您需要通过Materialized.as(...)aggregate() 提供一个固定名称。如果您提供一个固定名称(即在新旧拓扑中相同),问题就会消失。但是,由于您的原始拓扑没有提供固定名称,因此很难摆脱这种情况。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2023-03-26
    • 1970-01-01
    • 2019-01-08
    • 1970-01-01
    • 2020-07-03
    • 1970-01-01
    • 1970-01-01
    • 2018-10-05
    相关资源
    最近更新 更多