【问题标题】:How to "rate-limit" a PCollection in Apache Beam?如何在 Apache Beam 中“限制”PCollection?
【发布时间】:2017-04-06 15:35:26
【问题描述】:

我有一个似乎很常见的问题,但我不知道 Beam 推荐的解决方案是什么。

我有一个原始事件流,我正在寻找两个单独的事件来满足一个滑动窗口(60 分钟)内的条件,以便它“触发”警报。

使用SlidingWindows 很容易做到这一点,但是问题在于它的滑动特性,我可能在多个窗口中有效地获得该警报。我如何最终获得只输出一次此类警报的 PCollection(在特定时间范围/冷却持续时间内)?

我最初认为最近的状态处理功能会是我的解决方案,但后来意识到它只能在窗口内工作。侧面输入也是如此。所以在我看来,我需要一种打破窗户并在一个(可能的会话)窗口中处理警报“触发”的方法。但是文档没有提到任何有效地将元素重新分配给新窗口的方法

【问题讨论】:

    标签: google-cloud-dataflow apache-beam stream-processing


    【解决方案1】:

    有趣的应用程序!

    总结一下:

    • 对于您的用例来说,听起来“滑动窗口”意味着连续滑动。您可以选择最小粒度,但不一定是自然的。
    • 您感兴趣的每组事件应该只产生一个输出。

    有几种方法可以解决这个问题,具体取决于应用程序的其余部分。

    一种方法是将数据保留在全局窗口中并使用状态。您必须自己管理延迟 - 删除太迟的元素,考虑数据乱序等,并通常保持您的状态有界。

    另一种方法是使用带有Combine 或状态的滑动窗口(基本上,您已经尝试过),然后重新设置警报窗口并进行重复数据删除。您可以为此使用固定窗口,因为警报应该具有确定的时间戳;窗口的末尾将控制何时自动收集状态,这样很方便。

    【讨论】:

    • 肯恩,感谢您的回复。我实际上忘记用我的发现来更新这个问题。关于您的建议: > 一种方法是将您的数据留在全局窗口中并使用状态。这确实是一种选择,但我不会错过使用光束的意义吗?
    • > 另一种方法是使用带有组合或状态的滑动窗口(基本上,您已经尝试过),然后重新窗口化警报并进行重复数据删除。我有效地采用了这个解决方案。但是,固定窗口不起作用,因为我可以在窗口的后半部分发出警报,因此重置“冷却时间”,因此不应触发下一个窗口中的任何早期警报。
    【解决方案2】:

    我最终采用了一个重新窗口化策略,类似于 @Kenn 建议的。

    所以我有来自滑动窗口集合的警报,我将其重新窗口化到会话窗口中

    .apply(Window.remerge())            
    .apply(Window.into(Sessions.withGapDuration(Duration.standardHours(1))))
    

    在那个窗口集合中,我可以只做一个groupBy,从而获得一个会话的所有Alerts,在其中我可以应用我的冷却逻辑,即每小时只发出一个警报。

    【讨论】:

      猜你喜欢
      • 2022-12-31
      • 2023-02-03
      • 1970-01-01
      • 1970-01-01
      • 2023-04-10
      • 2018-05-16
      • 1970-01-01
      • 1970-01-01
      • 2023-04-10
      相关资源
      最近更新 更多