sam*_*ser 6 watermark stream-processing apache-flink flink-streaming
比方说,我们有一个TumblingEventTimeWindow与大小5分钟。我们有包含2 条基本信息的事件:
在这个例子中,我们在工作人员机器的挂钟时间下午 12:00启动我们的Flink拓扑(当然工作人员可能有不同步的时钟,但这超出了本问题的范围)。该拓扑包含一个处理运算符,其职责是汇总属于每个窗口的事件值和一个与此问题无关的 KAFKA Sink。
在这种情况下,几个事件到达Flink Operator,具有不同的事件时间戳12:01 - 12:09。此外,事件时间戳与我们的处理时间相对一致(如下面的 X 轴所示)。由于我们正在处理EVENT_TIME特性,因此应通过其事件时间戳来确定偶数是否属于特定事件。
在那个流程中,我假设两个翻滚窗口的边界是并且仅仅因为我们在12:00开始执行拓扑。如果这个假设是正确的(我希望不是),那么在回填情况下会发生什么,其中几个旧事件带有更旧的事件时间戳,并且我们在12:00再次启动了拓扑?(足够老,我们的迟到津贴不包括他们)。类似于以下内容:12:00 -- 12:0512:05 -- 12:10
上一个问题的答案也将解决这个问题,但我认为在这里明确提及它会有所帮助。比方说,我有这个TumblingEventTimeWindow的大小5分钟。然后在12:00我开始回填工作,它在许多事件中冲向时间戳覆盖范围的 Flink 操作员10:02 - 10:59;但由于这是一个回填工作,整个执行大约需要3 分钟才能完成。
作业是否会分配12 个单独的窗口并根据事件的事件时间戳正确填充它们?这12 个窗口的边界是什么?我最终会得到12 个输出事件,每个事件都有每个分配窗口的总和值吗?
我也对此类逻辑和运算符的自动化测试有一些担忧。操纵处理时间的最佳方式,触发某些行为,从而为测试目的塑造所需的窗口边界。特别是因为到目前为止我读到的关于利用的东西Test Harnesses似乎有点混乱,并且可能会导致一些不太容易阅读的混乱代码:
我在这方面学到的大部分知识以及我的一些困惑的根源都可以在以下地方找到:
TumblingSlidingSession非常感谢您的帮助,如果您知道有关这些概念及其内部工作的任何更好的参考资料,请告诉我。
如果您运行具有事件时间语义的作业,则窗口运算符的处理时间完全无关
这是正确的,我理解那部分。一旦你处理了EVENT_TIME特征,你就几乎脱离了语义/逻辑中的处理时间。我提出处理时间的原因是我对以下关键问题感到困惑,这对我来说仍然是个谜:
窗户的边界是如何计算的?!
另外,非常感谢澄清之间的区别out-of-orderness和lateness。我正在处理的代码因用词不当而完全让我失望(继承自的类的构造函数参数BoundedOutOfOrdernessTimestampExtractor被命名maxLatency):/
考虑到这一点,让我看看,如果我能得到这个正确的关于如何水印计算和当一个事件将被丢弃(或侧输出):
max-event-time-seen-so-far - max-out-of-orderness-allowedmax-event-time-seen-so-far - allowed-latenessmax-event-time-seen-so-far并且在任何情况下,无论事件,其事件时间戳是小于或等于的current-watermark,将被丢弃(侧输出的),正确?
这带来了一个新问题。你什么时候想使用out of orderness而不是lateness?由于当前的水印计算(数学上)在这些情况下可以相同。当您同时使用两者时会发生什么(甚至有意义)?!
这对我来说仍然是主要的谜团。鉴于上面的所有讨论,让我们重新审视我提供的具体示例,看看这里是如何确定窗口的边界的。假设我们有以下场景(事件的形状为(value, timestamp)):
DataStream与BoundedOutOfOrdernessTimestampExtractor具有2分钟 maxOutOfOrdernessallowedLateness为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]。
事件时间窗口的时间线如下:
watermark >= window_endtime - 1到达 -> 窗口被触发(结果被发出),但状态不会被丢弃watermark >= window_endtime + allowed_latenes到达 -> 状态被丢弃此窗口的 2. 和 3. 事件延迟,但在允许的延迟范围内。这些事件将添加到现有状态中,并且默认情况下,会在每条记录上触发窗口,发出精确的结果。
3.之后,该窗口的事件将被丢弃(或发送到后期输出接收器)。
所以,是的,配置两者是有意义的。无序性决定了窗口第一次被触发的时间,而允许的延迟决定了状态保持多长时间以可能更新结果。
关于边界:翻滚事件时间窗口具有固定长度,跨键对齐并从 unix 纪元开始。空窗户,不存在。对于您的示例,这意味着:
希望这可以帮助。
康斯坦丁
| 归档时间: |
|
| 查看次数: |
1060 次 |
| 最近记录: |