【问题标题】:Kafka Streams: reprocessing old data when windowingKafka Streams:窗口化时重新处理旧数据
【发布时间】:2022-01-12 06:53:34
【问题描述】:

拥有一个 Kafka Streams 应用程序,该应用程序通过流连接执行窗口化(使用原始事件时间,而不是挂钟时间),例如1 天。

如果启动此拓扑,并从头开始重新处理数据(如在 lambda 样式架构中),此窗口是否会将旧数据保留在那里?大 例如:如果今天是2022-01-09,我正在接收2021-03-01的数据,这个旧数据会进入表格,还是从一开始就被拒绝?

在这种情况下 - 可以采取哪些策略来重新处理这些数据?

更新使用 Kafka Streams 2.5.0

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    使用流时间应该会看到重新处理与原始运行一样工作。但是,在使用多个主题时有一些警告:流时间将始终是在任何主题(子拓扑和任务的)中找到的最新时间戳。这可能会导致不同的结果。在从时间戳差异很大的分区中消费时,确定有效流时间的不同 Kafka Streams 版本正在不断改进。这使得在不知道 Kafka Streams 版本的情况下很难确定地评论您的问题。

    在您的示例中,只要流时间没有进展太远,旧数据就会被接受。重新处理整个数据集应该可以工作,因为它将线性地通过您的主题。如果旧数据在超过窗口大小 + 宽限期的时间窗口中聚合,Kafka Streams 将拒绝该记录。在这种情况下,Kafka Streams 也会发出错误消息并相应地调整其指标。所以这种行为应该很容易上手。

    如果可行,我建议尝试这种重新处理并查看日志和指标。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-08-21
      • 2017-09-30
      • 1970-01-01
      • 2018-11-10
      • 1970-01-01
      • 2017-01-07
      • 1970-01-01
      相关资源
      最近更新 更多