【问题标题】:Too many timers cost too much time when checkpointing in Flink在 Flink 中进行检查点时,计时器过多会花费太多时间
【发布时间】:2018-11-10 06:55:33
【问题描述】:

我有一种情况需要使用StateTimeService 对大量消息进行滑动计数。滑动大小为 1,窗口大小大于 10 小时。我遇到的问题是检查点需要很多时间。为了提高性能,我们使用增量检查点。但是当系统做检查点时它仍然很慢。我们发现大部分时间用于序列化用于清理数据的计时器。我们为每个键设置了一个计时器,总共有大约 3 亿个计时器。

任何解决此问题的建议将不胜感激。或者我们可以用另一种方式进行计数? ——————————————————————————————————————————————— 我想为这种情况添加一些细节。滑动大小是一个事件,窗口大小超过10小时(每秒大约有300个事件),我们需要对每个事件做出反应。所以在这种情况下我们没有使用 Flink 提供的 windows。我们使用keyed state 来存储以前的信息。 timers 用于ProcessFunction 触发旧数据的清理工作。最后 dinstinct 键的数量非常大。

【问题讨论】:

  • 您能否提供更详细的说明?我试图回答你,但没有更多细节很难
  • 请说明情况。滑动尺寸是“一”什么?一小时、一分钟还是一个事件?每个事件分配到多少个不同的窗口?窗口化与所讨论的计时器有何关系(您是在谈论 flink 用于 timeWindows 的计时器,还是 ProcessFunction 中的某些东西)?实际上有 300M 不同的键吗?
  • 感谢您的关注。我在情况中添加了一些细节。我希望这可以澄清这个问题。

标签: apache-flink flink-streaming


【解决方案1】:

我认为这应该可行:

通过有效地执行 keyBy(key mod 100000) 之类的操作,将 Flink 使用的键数量从 300M 大幅减少到 100K(例如)。然后,您的 ProcessFunction 可以使用 MapState(其中键是原始键)来存储它需要的任何内容。

MapStates 具有迭代器,您可以使用它来定期抓取这些地图中的每一个以使旧项目过期。坚持每个键只有一个定时器的原则(如果你愿意的话,每个 uberkey),这样你就只有 10 万个定时器。

更新:

Flink 1.6 包含FLINK-9485,它允许定时器被异步检查点,并存储在 RocksDB 中。这使得 Flink 应用程序拥有大量定时器更加实用。

【讨论】:

  • 非常感谢。我会试一试。我认为它可能有效!
【解决方案2】:

如果不使用计时器,而是向流的每个元素添加一个额外的字段来存储当前处理时间或到达时间,那会怎样?因此,一旦您想从流中清除旧数据,您只需要使用过滤器运算符并检查是否可以删除旧数据。

【讨论】:

  • 感谢您的关注。在某些情况下,这可能是一个解决方案,但不是我的。当该键下没有事件时,我无法清除旧数据。
【解决方案3】:

与其在每个事件上注册一个清除计时器,不如在某个时间段内仅注册一次计时器,例如每1分钟一次?您只能在第一次看到密钥时注册它,并在onTimer 中刷新它。某样东西:

new ProcessFunction<SongEvent, Object>() {

  ...

  @Override
  public void processElement(
      SongEvent songEvent,
      Context context,
      Collector<Object> collector) throws Exception {

    Boolean isTimerRegistered = state.value();
    if (isTimerRegistered != null && !isTimerRegistered) {
      context.timerService().registerProcessingTimeTimer(time);
      state.update(true);
    }

    // Standard processing


  }

  @Override
  public void onTimer(long timestamp, OnTimerContext ctx, Collector<Object> out)
      throws Exception {
    pruneElements(timestamp);

    if (!elements.isEmpty()) {
      ctx.timerService().registerProcessingTimeTimer(time);
    } else {
      state.clear();
    }
  }
}

Flink SQL Over 子句实现了类似的东西。你可以看看here

【讨论】:

  • 我就是这样使用定时器的。一键(或状态)一个定时器。问题是键太多了。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2018-01-08
  • 2019-11-14
  • 2012-09-03
  • 2013-07-11
  • 2017-09-24
  • 1970-01-01
相关资源
最近更新 更多