【问题标题】:flink streaming window triggerflink流窗口触发
【发布时间】:2018-12-01 13:17:10
【问题描述】:

我有 flink 流,我在某个时间窗口(比如 30 秒)上计算了一些东西。

这里发生的事情也给了我聚合以前的窗口的结果。

说前 30 秒我得到结果 10。

接下来的 3 秒我想要新的结果,而不是我得到最后一个窗口结果 + 新的 等等。

所以我的问题是如何为每个窗口获得新的结果。

【问题讨论】:

  • 您可以告诉我们您的代码。

标签: scala apache-flink flink-streaming


【解决方案1】:

您需要使用purging 触发器。你想要的是 FIRE_AND_PURGE(发射和删除窗口内容),默认的 flink 触发器所做的是 FIRE(发射和保留窗口内容)。

input
    .keyBy(...)
    .timeWindow(Time.seconds(30))

    // The important part: Replace the default non-purging ProcessingTimeTrigger
    .trigger(new PurgingTrigger[..., TimeWindow](ProcessingTimeTrigger))

    .reduce(...)

如需更深入的解释,请查看Triggers 和FIRE vs FIRE_AND_PURGE。

触发器确定窗口(由窗口分配器形成)何时准备好由窗口函数处理。每个 WindowAssigner 都带有一个默认触发器。如果默认触发器不符合您的需求,您可以使用 trigger(...) 指定自定义触发器。

当触发器触发时,它可以是 FIRE 或 FIRE_AND_PURGE。 FIRE 保留窗口的内容,而 FIRE_AND_PURGE 删除其内容。默认情况下,预实现的触发器只是 FIRE 而不清除窗口状态。

【讨论】:

  • 假设我使用 30 秒翻转窗口,我想在每个窗口中进行一些计算并将其传递给后续窗口。我需要使用ProcessingTimeTrigger.create() 还是PurgingTimeTrgger.create(ProcessingTimeTrigger.create())。如果是30秒的翻滚窗口……不是30秒后窗口内容已经消失了吗?如果是这样,我们可以使用ProcessingTimeTrigger.create()吗?
  • 我无法理解为什么窗口内容仍然存在?另外,当您说窗口内容时,您是指属于窗口的事件还是窗口事件的计算结果?如果我使用 PurgingTrigger 会清除哪一个?
【解决方案2】:

您描述的功能可以在 Tumbling Windows 中找到:https://ci.apache.org/projects/flink/flink-docs-release-1.2/dev/windows.html#tumbling-windows

更多细节和/或代码会有所帮助:)

【讨论】:

    【解决方案3】:

    我在这个问题上有点晚了,但我遇到了与 OP 相同的问题。后来我发现是我自己的代码中的一个错误。仅供参考,我的错误可以很好地解决您的问题。

    // Old code (modified to be an example):
    val tenSecondGrouping: DataStream[MyCustomGrouping] = userIdsStream
          .keyBy(_.somePartitionedKey)
          .window(TumblingProcessingTimeWindows.of(Time.of(10, TimeUnit.SECONDS)))
          .trigger(ProcessingTimeTrigger.create())
          .aggregate(new MyCustomAggregateFunc(new MyCustomGrouping()))
    

    new MyCustomGrouping 发生错误:我无意中创建了一个单独的 MyCustomGrouping 对象并在 MyCustomAggregateFunc 中重用它。随着更多的翻滚窗口创建,后面的聚合结果变得疯狂!修复方法是在每次触发 MyCustomAggregateFunc 时创建新的 MyCustomGrouping。所以:

    // New code, problem solved
              ...
              .aggregate(new MyCustomAggregateFunc(() => new MyCustomGrouping())) 
    // passing in a func to create new object per trigger
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-04-09
      • 2020-07-16
      • 1970-01-01
      • 2019-08-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多