【问题标题】:How to apply an offset to Tumbling Window, in order to delay the starting of Windows<TimeWindow> in Kafka Streams如何对 Tumbling Window 应用偏移量,以延迟 Kafka Streams 中 Windows<TimeWindow> 的启动
【发布时间】:2023-01-19 07:47:23
【问题描述】:

我正在使用不同的窗口大小计算具有 2022 年 5 月值的数据集的简单均值。使用 1 小时窗口没有问题,而使用 1 周和 1 个月窗口时,记录评估不正确。

正如here所讨论的那样,问题是由于这样的事实自 Unix 纪元 (01-01-1970) 以来,时间被划分为具有指定持续时间的大小相等的块(窗口),然后传入事件被分配到这些块(窗口)中.

所以这意味着使用 31 天的窗口,在 Kafka Streams 中时间是这样划分的:

 01-01-1970 : 31-01-1970
 01-02-1970 : 03-02-1970
 ...
[14-04-2022 : 15-05-2022] <-- Our Window
 16-05-2022 : 15-06-2022
 ...

因此没有所需的 01-05-2022 : 31-05-2022 窗口。

在那个discussion(关于 Flink)中,解决方案是应用 17 天的抵消到 Tumbling Window,以便将窗口从 14-04 转移到 01-05

var monthResult = keyed
        .window(TumblingEventTimeWindows.of(Time.days(31),Time.days(17)))
        .aggregate(new AvgQ1(Config.MONTH))
        .name("Monthly Window Mean AggregateFunction");

但是使用 Kafka Stream,我没有找到偏移函数,也没有找到让我达到相同结果的东西。

这就是我实际定义窗口的方式:

var grouped = keyed
    .groupByKey(Grouped.with(Serdes.Long(), EventSerde.Event()))
    .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(10)))
    .reduce((o, v1) -> o);

【问题讨论】:

  • 你找到解决办法了吗?有同样的问题。

标签: java docker apache-kafka apache-kafka-streams stream-processing


【解决方案1】:

您可以在 Kafka Streams 中使用自定义 TimeWindow 实现。

比照https://github.com/confluentinc/kafka-streams-examples/blob/7.1.1-post/src/test/java/io/confluent/examples/streams/window/DailyTimeWindows.java

有一张将其本地添加到 Kafka Streams 的票,但需求不是很高,因此从未完成。如果你想捡起来,那就太好了!

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-04-02
    • 1970-01-01
    • 2019-03-16
    • 1970-01-01
    • 2018-06-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多