【问题标题】:KStreamWindowAggregate 2.0.1 vs 2.5.0: skipping records instead of processingKStreamWindowAggregate 2.0.1 vs 2.5.0:跳过记录而不是处理
【发布时间】:2020-09-15 10:47:34
【问题描述】:

我最近将我的 kafka 流从 2.0.1 升级到了 2.5.0。结果,我看到了很多类似以下的警告:

org.apache.kafka.streams.kstream.internals.KStreamWindowAggregate$KStreamWindowAggregateProcessor Skipping record for expired window. key=[325233] topic=[MY_TOPIC] partition=[20] offset=[661798621] timestamp=[1600041596350] window=[1600041570000,1600041600000) expiration=[1600059629913] streamTime=[1600145999913]

KStreamWindowAggregate 类中似乎有新的逻辑来检查窗口是否已关闭。如果它已关闭,则跳过消息。与 2.0.1 相比,这些消息仍在处理中。

问题

有没有办法获得与以前相同的行为?在这次升级中,我发现我的数据中有很多空白,但不知道如何解决这个问题,因为以前没有看到这些空白。

我正在使用的聚合函数已经处理了窗口化,因此处理了过期的窗口。这个新逻辑与这个即将到期的窗口有什么关系?

更新

在进一步探索时,我确实发现它与 ms 中的宽限期有关。似乎在我的自定义时间戳提取器(具有使用有效负载中的时间戳而不是普通时间戳的逻辑)中,我能够看到过期窗口警告的传入时间戳确实大于 24 小时相比来自负载的事件时间。

我认为这是由于消费者滞后超过 24 小时造成的。

时间戳提取器提取方法有一个分区时间,根据文档:

partitionTime 提取的当前记录分区的最高有效时间戳˙(如果未知,可能为-1)

那么这是该主题记录的创建时间吗?有没有办法以不再跳过我的记录的方式影响这一点?

【问题讨论】:

标签: apache-kafka apache-kafka-streams


【解决方案1】:

与 2.0.1 相比,这些消息仍在处理中。

这有点令人惊讶(即使我需要仔细检查代码),至少对于默认配置而言。默认情况下,存储 保留时间 设置为 24 小时,因此在 2.0.1 中,早于 24 小时的消息也不应该被处理,因为相应的状态已经被清除。如果您确实将存储保留时间(通过Materialized#withRetention)更改为更大的值,您还需要通过TimeWindows#grace() 方法相应地增加窗口宽限期

我正在使用的聚合函数已经处理了窗口化,因此也处理了过期的窗口。这个新逻辑与这个即将到期的窗口有什么关系?

不确定您的意思是什么或您实际上是如何做到的?新旧逻辑在窗口存储时间(retention time 配置)方面是相似的。新部分是宽限期,如果您愿意,您可以将其增加到与保留时间相同的值。

关于“分区时间”:它是根据TimestampExtractor 返回的任何内容计算的。对于您的情况,这是您从消息有效负载中提取的最大值。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-05-22
    • 2020-06-19
    • 2020-07-20
    • 1970-01-01
    • 2014-10-07
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多