【问题标题】:Storm: Min/max aggregation across several sliding windows with varying sizesStorm:跨多个大小不同的滑动窗口的最小/最大聚合
【发布时间】:2017-05-04 16:32:33
【问题描述】:

我想知道在 Apache Storm 中解决以下问题的最佳做法是什么。

我有一个单独的 spout,它生成一个带有明确时间戳的整数值流。目标是在此流上使用三个滑动窗口执行最小/最大聚合:

  • 最后一小时
  • 最后一天,即过去 24 小时

最后一小时很容易:

topology.setBolt("1h", ...)
    .shuffleGrouping("spout")
    .withWindow(Duration.hours(1), Duration.seconds(10))
    .withTimestampField("timestamp"));

但是,在较长时期内,我担心窗口的队列大小。当我像上一小时聚合一样直接从 spout 使用元组时,每个元组都会最终进入队列。

一种可能性是从预先聚合的“1h”bolt 中消耗元组。但是,由于我使用的是显式时间戳,因此来自“1h”螺栓的迟到的元组将被忽略。 1 小时的延迟不是一个选项,因为这会延迟窗口的评估。有没有办法“允许”迟到的元组而不影响结果的及时性?

当然,我也可以每小时存储一个聚合,然后计算过去 24 小时的最小值,包括“1h”流中的最新值。但我很好奇是否有办法使用 Storm 方法正确地做到这一点。

更新 1

感谢 arunmahadevan 的回答,我更改了 1h min 螺栓以在相应的 1h 窗口中发出所有元组的最大时间戳的最小元组。这样,消费螺栓不会因为迟到而丢弃元组。我还引入了一个新字段original-timestamp 来保留最小元组的原始时间戳。

更新 2

我终于找到了一种更好的方法,即只在 1 小时分钟的螺栓中发出状态变化。只要没有收到新的元组,Storm 就不会在消耗螺栓中提前时间,因此可以防止迟到问题。此外,我可以保留原始时间戳,而无需将其复制到单独的字段中。

【问题讨论】:

    标签: apache-storm


    【解决方案1】:

    我认为定期从“1h”到“24h”螺栓发出最小值应该可以工作并保持“24h”队列大小受到检查。

    如果您配置了延迟,则仅在该延迟之后(即当事件时间超过滑动间隔 + 延迟时)调用螺栓的执行。

    假设如果“1h”bolt 配置了 1 分钟的延迟,则只有在事件时间超过 02:01 之后,才会为 01:00 - 02:00 之间的元组调用执行。 (即,螺栓已经看到时间戳 >= 02:01 的事件)。然而,执行将仅在 01:00 和 02:00 之间接收元组。

    现在,如果您计算最后一小时的最小值并将结果发送到“24 小时”螺栓,该螺栓的滑动间隔为 1 小时且延迟 = 0,一旦传入事件的时间戳跨越下一小时,它将触发。如果您发出时间戳为 02:00 的 01:00-02:00 分钟,“24h”窗口将在收到 min 事件后立即触发(对于前一天 02:00 到 02:00 之间的事件)因为事件时间超过了下一个小时,并且配置的延迟为 0。

    【讨论】:

    • 好的,你回答中真正帮助我的关键部分是发出最小的和正在考虑的窗口的最大时间戳。
    • 您知道如何实现平均聚合吗?这是不同的,因为我不能在 1h 螺栓上使用滑动窗口,因为这会导致相同的元组平均被合并多次。有没有办法将基于翻滚窗口的预聚合与源 spout 中的最新元素结合起来?一旦它们被来自翻滚窗口的聚合覆盖,我会以某种方式需要从队列中逐出源 spout 元素。但是 Storm 中似乎没有自定义驱逐策略的可能性......
    猜你喜欢
    • 2012-05-30
    • 2017-04-03
    • 1970-01-01
    • 2021-12-11
    • 2016-05-27
    • 2019-03-13
    • 2017-09-19
    • 2013-10-27
    • 2020-04-27
    相关资源
    最近更新 更多