【发布时间】:2018-08-11 18:07:06
【问题描述】:
如何使用固定大小的基于计数的窗口实现滑动窗口聚合(或转换)?
例如:如果我有如下流数据
input stream = 1,2,3,4,5,6,7,8...
假设时间在这里不相关。假设我的聚合函数是 AVERAGE 并且窗口大小固定为 3 条记录(不是 3 毫秒、3 秒、3 小时等),我希望我的输出流是
output stream = avg(1,2,3), avg(2,3,4), avg(3,4,5), avg(4,5,6), avg(5,6,7)... = 2,3,4,5,6...
Kafka 流工作中记录的 Windows 是“基于时间的”。甚至基类 Window 的构造函数也有以下签名:
Window(long startMs, long endMs)
所以我不确定它是否是进行非基于时间的窗口聚合的正确工具。
Apache Flink 支持count-based sliding and tumbling windows。这正是我所需要的,但我正在 Kafka Streams 中寻找类似的功能。
【问题讨论】: