【发布时间】: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)
那么这是该主题记录的创建时间吗?有没有办法以不再跳过我的记录的方式影响这一点?
【问题讨论】:
-
您的宽限期设置为多少?
-
没改,所以我猜默认是24小时
标签: apache-kafka apache-kafka-streams