【发布时间】: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