【问题标题】:How to do aggregation on fixed-size count-based sliding window?如何在固定大小的基于计数的滑动窗口上进行聚合?
【发布时间】: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 中寻找类似的功能。

【问题讨论】:

    标签: apache-kafka-streams


    【解决方案1】:

    如果您不关心时间顺序,您可以实现带有附加状态的自定义Transformer

    StreamsBuilder builder = new StreamsBuilder();
    builder.addStoreStore(...); // add KeyValueStore here
    KStream result = builder.stream("topic").transform(...); // pass in name of your KeyValueStore, too
    

    对于您自定义Transformer,您可以为每个键维护一个List,列表是您的窗口-只要列表小于您的窗口大小,您就可以将新记录附加到列表中-如果它是大小,你触发计算 -- 如果它超过大小,你修剪它并在之后触发计算。

    有关详细信息,请参阅文档:https://kafka.apache.org/10/documentation/streams/developer-guide/processor-api.html(请注意,ProcessorTransformer 基本相同。)

    【讨论】:

      【解决方案2】:

      如果你想使用同样是流引擎的 Apache Storrm,可以将 kafka 作为数据源连接到它。 Storm 新版本提供了一个名为 Tumbling Window 的概念,它为您的拓扑提供了确切数量的元组。这可以很容易地用来解决您的问题。

      更多请关注https://docs.hortonworks.com/HDPDocuments/HDP2/HDP-2.6.0/bk_storm-component-guide/content/storm-windowing-concepts.html

      【讨论】:

      • 看起来翻滚窗口可以基于计数,但滑动窗口总是基于时间。我想要一个基于计数的滑动窗口。
      • 它是双向的,docs.hortonworks.com/HDPDocuments/HDP2/HDP-2.6.0/…,问题是你想要基于计数而不是基于时间
      • 是的。看起来storm支持基于计数的滑动窗口。 Apache Flink 也支持它。想知道为什么 Kafka 流不支持它。
      猜你喜欢
      • 2021-12-11
      • 2021-09-10
      • 1970-01-01
      • 2014-04-20
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-06-23
      相关资源
      最近更新 更多