Flink Windows 边界、水印、事件时间戳和处理时间

sam*_*ser 6 watermark stream-processing apache-flink flink-streaming

问题定义和建立概念

比方说,我们有一个TumblingEventTimeWindow与大小5分钟。我们有包含2 条基本信息的事件:

  • 数字
  • 事件时间戳

在这个例子中,我们在工作人员机器的挂钟时间下午 12:00启动我们的Flink拓扑(当然工作人员可能有不同步的时钟,但这超出了本问题的范围)。该拓扑包含一个处理运算符,其职责是汇总属于每个窗口的事件值和一个与此问题无关的 KAFKA Sink。

  • 这个窗口有一个BoundedOutOfOrdernessTimestampExtractor,允许延迟一分钟。
  • 水印:据我所知,Flink 和 Spark Structured Stream 中的水印定义为(max-event-timestamp-seen-so-far - allowed-lateness)。任何事件时间戳小于或等于此水印的事件都将在结果计算中被丢弃和忽略。

第 1 部分(确定窗口的边界)

快乐(实时)路径

在这种情况下,几个事件到达Flink Operator,具有不同的事件时间戳12:01 - 12:09。此外,事件时间戳与我们的处理时间相对一致(如下面的 X 轴所示)。由于我们正在处理EVENT_TIME特性,因此应通过其事件时间戳来确定偶数是否属于特定事件。

在此处输入图片说明

旧数据涌入

在那个流程中,我假设两个翻滚窗口的边界是并且仅仅因为我们在12:00开始执行拓扑。如果这个假设是正确的(我希望不是),那么在回填情况下会发生什么,其中几个旧事件带有更旧的事件时间戳,并且我们在12:00再次启动了拓扑?(足够老,我们的迟到津贴不包括他们)。类似于以下内容:12:00 -- 12:0512:05 -- 12:10

在此处输入图片说明

  1. 如果是这样,那么我们的事件当然不会在任何窗口中捕获,所以再次,我希望这不是行为:)
  2. 另一种选择是通过到达事件的事件时间戳来确定窗口的边界。如果是这样,那将如何运作?注意到的最小事件时间戳成为第一个窗口的开始,然后根据大小(在本例中为5 分钟)确定随后的边界?因为这种方法也会有缺陷和漏洞。你能解释一下这是如何工作的以及如何确定窗口的边界吗?

回填事件蜂拥而至

上一个问题的答案也将解决这个问题,但我认为在这里明确提及它会有所帮助。比方说,我有这个TumblingEventTimeWindow的大小5分钟。然后在12:00我开始回填工作,它在许多事件中冲向时间戳覆盖范围的 Flink 操作员10:02 - 10:59;但由于这是一个回填工作,整个执行大约需要3 分钟才能完成。

作业是否会分配12 个单独的窗口并根据事件的事件时间戳正确填充它们?这12 个窗口的边界是什么?我最终会得到12 个输出事件,每个事件都有每个分配窗口的总和值吗?

第 2 部分(此类有状态运算符的单元/集成测试)

我也对此类逻辑和运算符的自动化测试有一些担忧。操纵处理时间的最佳方式,触发某些行为,从而为测试目的塑造所需的窗口边界。特别是因为到目前为止我读到的关于利用的东西Test Harnesses似乎有点混乱,并且可能会导致一些不太容易阅读的混乱代码:

参考

我在这方面学到的大部分知识以及我的一些困惑的根源都可以在以下地方找到:

  • 时间戳提取器和水印发射器
  • 事件时间处理和水印
  • 在 Spark 中处理延迟数据和水印
    • Spark 文档该部分中的图像非常有用且具有教育意义。但与此同时,窗口边界与这些处理时间而不是事件时间戳对齐的方式给我带来了一些困惑。
    • 此外,在该可视化中,似乎每5 分钟计算一次水印,因为这是窗口的滑动规范。这是应该多久计算一次水印的决定因素吗?这是如何工作弗林克对于不同的窗口(例如,,等等)?TumblingSlidingSession

非常感谢您的帮助,如果您知道有关这些概念及其内部工作的任何更好的参考资料,请告诉我。

以下@snntrable 回答后的更新

如果您运行具有事件时间语义的作业,则窗口运算符的处理时间完全无关

这是正确的,我理解那部分。一旦你处理了EVENT_TIME特征,你就几乎脱离了语义/逻辑中的处理时间。我提出处理时间的原因是我对以下关键问题感到困惑,这对我来说仍然是个谜:

窗户的边界是如何计算的?!

另外,非常感谢澄清之间的区别out-of-orderness和lateness。我正在处理的代码因用词不当而完全让我失望(继承自的类的构造函数参数BoundedOutOfOrdernessTimestampExtractor被命名maxLatency):/

考虑到这一点,让我看看,如果我能得到这个正确的关于如何水印计算和当一个事件将被丢弃(或侧输出):

  • 无序分配器
    • 当前水印 = max-event-time-seen-so-far - max-out-of-orderness-allowed
  • 允许迟到
    • 当前水印 = max-event-time-seen-so-far - allowed-lateness
  • 常规流量
    • 当前水印 = max-event-time-seen-so-far

并且在任何情况下,无论事件,其事件时间戳是小于或等于的current-watermark,将被丢弃(侧输出的),正确?

这带来了一个新问题。你什么时候想使用out of orderness而不是lateness?由于当前的水印计算(数学上)在这些情况下可以相同。当您同时使用两者时会发生什么(甚至有意义)?!

回到 Windows 的边界

这对我来说仍然是主要的谜团。鉴于上面的所有讨论,让我们重新审视我提供的具体示例,看看这里是如何确定窗口的边界的。假设我们有以下场景(事件的形状为(value, timestamp)):

  • 操作员在12:00 PM 开始(这是处理时间)
  • 事件按以下顺序到达操作员
    • (1, 8:29 )
    • (5, 8:26 )
    • (3, 9:48 )
    • (7, 9:46 )
  • 我们有一个大小为5 分钟的TumblingEventTimeWindow
    • 窗口被施加到DataStream与BoundedOutOfOrdernessTimestampExtractor具有2分钟 maxOutOfOrderness
  • 此外,窗口配置allowedLateness为1 分钟

注意:如果你不能同时拥有out of orderness和lateness或根本没有任何意义,请只考虑out of orderness在上面的例子。

最后,请您布置将分配一些事件的窗口,并指定这些窗口的边界(窗口的开始和结束时间戳)。我假设边界也由事件的时间戳决定,但在像这样的具体例子中弄清楚它们有点棘手。

再次,非常感谢,并真正感谢您的帮助:)

小智 7

原答案

Watermark:据我了解,Flink 和 Spark Structured Stream 中的水印定义为(max-event-timestamp-seen-so-far - allowed-lateness)。任何事件时间戳小于或等于此水印的事件都将在结果计算中被丢弃并忽略。

这是不正确的,可能是造成混乱的根源。在 Flink 中,乱序和迟到是不同的概念。带有BoundedOutOfOrdernessTimestampExtractor水印的是max-event-timestamp-seen-so-far - max-out-of-orderness。有关允许延迟的更多信息,请参阅 Flink 文档 [1]。

如果您使用事件时间语义运行作业,则窗口运算符的处理时间完全无关:

  • 事件将根据事件时间时间戳分配给它们的窗口
  • 一旦水印达到最大时间戳 ( window end time -1),时间窗口就会被触发。
  • 时间戳早于的事件current watermark - allowed lateness将被丢弃或发送到后期数据侧输出 [1]

这意味着,如果您在中午 12:00(处理时间)开始工作并开始提取过去的数据,则水印也将是(甚至更远)过去的。因此,配置allowedLateness是无关紧要的,因为数据相对于偶数时间来说并没有迟到。

另一方面,如果您首先从中午 12:00 提取一些数据,然后从晚上 10:00 提取数据,则在您提取旧数据之前,水印将已经提前到 ~12:00pm。在这种情况下,晚上 10:00 的数据将“迟到”。如果它晚于配置的allowedLateness(默认=0),它将被丢弃(默认)或发送到侧面输出(如果配置)[1]。

跟进答案

事件时间窗口的时间线如下:

  1. 窗口内带有时间戳的第一个元素到达 -> 创建该窗口(和键)的状态
  2. watermark >= window_endtime - 1到达 -> 窗口被触发(结果被发出),但状态不会被丢弃
  3. watermark >= window_endtime + allowed_latenes到达 -> 状态被丢弃

此窗口的 2. 和 3. 事件延迟,但在允许的延迟范围内。这些事件将添加到现有状态中,并且默认情况下,会在每条记录上触发窗口,发出精确的结果。

3.之后,该窗口的事件将被丢弃(或发送到后期输出接收器)。

所以,是的,配置两者是有意义的。无序性决定了窗口第一次被触发的时间,而允许的延迟决定了状态保持多长时间以可能更新结果。

关于边界:翻滚事件时间窗口具有固定长度,跨键对齐并从 unix 纪元开始。空窗户,不存在。对于您的示例,这意味着:

  • (1, 8:29) 添加到窗口 (8:25 - 8:29:59:999)
  • (5, 8:26) 添加到窗口 (8:25 - 8:29:59:999)
  • (3, 9:48) 添加到窗口 (9:45 - 9:49:59:999)
  • (8:25 - 8:29:59:999) 被触发,因为水印已前进到 9:48-0:02=9:46,该时间大于窗口的最后一个时间戳。窗口状态也被丢弃,因为水印已经提前到了9:46,这也大于窗口的结束时间+允许的迟到(1分钟)
  • (7, 9:46) 添加到窗口 添加到窗口 (9:45 - 9:49:59:999)

希望这可以帮助。

康斯坦丁

[1] https://ci.apache.org/projects/flink/flink-docs-release-1.8/dev/stream/operators/windows.html#allowed-lateness