【问题标题】:Kafka Streams TimestampExtractor卡夫卡流时间戳提取器
【发布时间】:2018-10-08 02:26:55
【问题描述】:

大家好,我有一个关于 TimestampExtractor 和 Kafka Streams 的问题......

在我们的应用程序中可能会接收到乱序事件,因此我喜欢根据负载中的业务日期而不是它们放置在主题中的时间点来对事件进行排序。

为此,我编写了一个自定义 TimestampExtractor,以便能够从有效负载中提取时间戳。直到我在这里告诉的一切都运行良好,但是当我为这个主题构建 KTable 时,我发现我收到的事件发生了故障(从业务的角度来看,它不是最后一个事件,而是最后收到的)显示为对象的最后状态,而 ConsumerRecord 具有来自有效负载的时间戳。

我不知道认为 Kafka Stream 会使用 TimestampExtractor 解决这个乱序问题可能是我的错误。

然后在调试过程中,我看到如果 TimestampExtractor 返回 -1 作为结果,Kafka Streams 忽略了消息,并且 TimestampExtractor 还提供了最后接受的事件的时间戳,所以我构建了一个实现以下检查的逻辑(payloadTimestamp

我是否可以处理这样的逻辑或存在其他方法来处理 Kafka 流中的乱序事件....

谢谢解答..

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    目前(Kafka 2.0),KTables 在更新时间戳时不考虑时间戳,因为假设输入主题中没有乱序数据。这种假设的原因是“单写原则”——假设对于压缩的 KTable 输入主题,每个键只有一个生产者,因此,不会有任何乱序数据关于单键。

    这是一个已知问题:https://issues.apache.org/jira/browse/KAFKA-6521

    为了您的修复:执行此“破解”并非 100% 正确或安全:

    • 首先,假设您有两条不同的消息,带有两个不同的密钥<key1, value1, 5>, <key2, value2, 3>。与时间戳为 5 的第一条记录相比,时间戳为 3 的第二条记录晚。但是,两者都有不同的键,因此,您实际上希望将第二条记录放入 KTable 中。只有当您有两条具有相同键的记录时,您才希望删除迟到的数据 IHMO。
    • 其次,如果您有两条记录具有相同的键,而第二条记录乱序,并且在处理第二条记录之前发生崩溃,TimestampExtractor 会丢失第一条记录的时间戳。因此,在重新启动时,它不会丢弃乱序记录。

    要做到这一点,您需要在应用程序逻辑中“手动”过滤,而不是无状态且与密钥无关的 TimestampExtractor。无需通过builder#table() 读取数据,您可以将其作为流读取,然后应用.groupByKey().reduce() 来构建KTable。在你的Reducer逻辑中,你比较新旧记录的时间戳,并返回时间戳较大的记录。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-02-08
      • 2018-03-06
      • 2016-08-03
      • 2018-09-15
      • 2018-03-07
      相关资源
      最近更新 更多