【问题标题】:kafka streams session window retention durationkafka 流会话窗口保留持续时间
【发布时间】:2017-06-07 19:02:01
【问题描述】:

我们正在使用 Kafka 流的 SessionWindows 来聚合相关事件的到达。除了聚合,我们还使用until() API 指定窗口的保留时间。 直播信息
会话窗口(不活动时间)为 1 分钟,传递给until() 的保留时间为 2 分钟。 我们正在使用自定义的TimestampExtractor 来映射事件的时间。

示例:
事件:e1;活动时间:上午 10:00:00;到达时间:下午 2 点(同一天)
事件:e2;活动时间:上午 10:00:30;到达时间下午 2:10(同一天)
第二个事件的到达时间是 e1 到达后的 10 分钟,超过了保留时间 + 不活动时间。但较旧的事件 e1 仍然是聚合的一部分,尽管保留时间为 2 分钟。

问题:
1) kafka 流如何使用until() API 清理状态存储?由于指定为参数的保留值是“窗口将保持多长时间的下限”。究竟何时清除窗口?

2) 是否有后台线程定期清理状态存储?如果是,那么有没有办法确定清除窗口的实际时间。

3) 任何会在保留时间后清除窗口数据的流配置。

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    在我回答您的具体问题之前:请注意,保留时间不是基于系统时间,而是基于“流时间”。 “流时间”是基于 TimestampExtractor 返回的任何内部跟踪的时间进度。无需过多详细说明:对于您的示例,有 2 条记录,当第二条记录到达时,“流时间”将提前 30 秒,因此保留时间尚未过去。

    另请注意,如果没有新数据到达(至少一个分区),“流时间”不会提前。 这适用于 Kafka 0.11.0 及更早版本,但可能会在未来版本中发生变化。

    更新:在 Kafka 2.1 中,stream-time 的计算发生了变化,即使一个分区没有传递数据,stream-time 也可能会提前。详情见KIP-353: Improve Kafka Streams Timestamp Synchronization

    对于您的问题:

    (1) Kafka Streams 将所有存储更新写入更改日志主题和本地 RocksDB 存储。两者都分为具有一定大小的所谓段。如果新数据到达(即“流时间”进展),则创建新段。如果发生这种情况,如果旧段中的所有条记录早于保留时间(即,记录时间戳小于“流时间”减去保留时间),旧段将被删除。

    (2) 因此,没有后台线程,但清理是常规处理的一部分,

    和 (3) 没有强制清除旧记录/窗口的配置。

    如果 all 记录过期,则会删除整个段,因此段中较旧的记录(很可能具有较小/较旧的时间戳)的维护时间比保留时间长。这种设计背后的动机是性能:以每条记录为基础过期太贵了。

    【讨论】:

    • 感谢 Matthias 的详细解释。
    • 我有一个关于**“记录时间戳小于”流时间“减去保留时间”的后续问题**。如果“流时间”是过去,则清除段的条件可能总是错误的。对 POC 的观察是,在事件到达时会创建新的片段,但旧的事件不会被清除,而是被重新洗牌到不同的片段中。
    • 也许为此打开一个新问题?不确定您的意思:流时间已过去。流时间的想法是定义“现在”,关于处理进度。也不确定您所说的“旧事件没有被清除,而是被重新排列到不同的部分”是什么意思。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-12-08
    • 2017-06-20
    • 2012-03-02
    • 1970-01-01
    相关资源
    最近更新 更多