TT Lab
开始
学习 学习路径 课程

数据流水线

水位线是承诺,不是事实

在 TT Lab 中继续学习

一句话总结

水位线(watermark)是“这个时刻之前的数据都已经到齐了”的承诺,而不是观测,如果事先没有规定好遇到违背这个承诺的数据该怎么办,昨天给出的数字就会变得无法解释。

为什么需要它

本课程前面的实验中已经讲过一次水位线。那是把已加载数据的最大时间写在表里,让下一次运行只读取它之后的数据。这个水位线回答的问题只有一个——读到哪里了?

但光靠这个问题,有些情况解决不了。移动应用在地铁里丢了信号,几分钟后才发送事件;处于飞行模式的终端,要过几个小时才发送事件。这些事件的发生时间早已过去。以时间为基准的水位线永远看不到这些数据,而我们早已把那个时间段的汇总结果发出去了。

于是早会上就会有人说:“昨天在仪表板上看到的数字和今天的数字不一样。”如果回答不了这个问题,这个仪表板从那天起就没人相信了。

工作原理

首先要分清时间有三种。事件时间是事件实际发生的时间,摄入时间是它进入我们系统的时间,处理时间是我们进行计算的时间。分析的基准始终是事件时间。如果以处理时间为基准,每次重处理答案都会不同。

按事件时间划分窗口(window)之后,马上就会冒出一个问题:这个窗口什么时候关闭? 永远开着,就出不了结果;关得太早,又会漏掉还没到的数据。

水位线替我们做出这个判断。Flink 文档把水位线描述为“告知系统事件时间进展情况的机制”。常见的计算式如下。

워터마크 = 지금까지 본 최대 이벤트 시간 − 허용 지연
창이 닫힌다 = 워터마크가 그 창의 끝을 지났다

在事件时间轴上标出到达顺序的示意图。水位线位于已见最大事件时间向左偏移允许延迟的位置,越过窗口末尾的瞬间窗口关闭。之后到达的第 9 个事件,从事件时间看在窗口之内,却成为迟到数据,水位线不会因为它而倒退

这里有两点容易被忽略。

第一,水位线不会倒退。最大事件时间是单调递增的,所以水位线也单调递增。迟到的事件不会让水位线倒退。如果让它能够倒退,已经关闭的窗口就会反复地重新打开又关闭,任何值都无法确定下来。

第二,“迟到”不是数据的属性,而是到达顺序的属性。同一个事件,取决于它何时到达,既可能算迟到,也可能不算。所以判定要用在该事件之前到达的那些数据算出的水位线来做。这里有个常见的事故——为了方便处理数据,先按事件时间排了一次序。这一刻到达顺序就消失了,迟到数据会显示为 0 条。实际测量一下,在有二十多条乱序到达的数据中,结果也恰好是 0。不是因为没有问题,而是因为把能看见问题的眼睛去掉了。

允许延迟(allowed lateness)是在准确度和延迟之间做取舍的旋钮。设得大,能接纳更多迟到数据,但窗口也会相应更晚关闭;设得小,出结果快,但漏掉的更多。这个值不能凭感觉定,而要从实际的延迟分布中读出来。测一测延迟的中位数、95 分位数和最大值,通常会在 95 分位数附近出现拐点,再往上尾巴拖得很长。如果按最大值来设,就会因为那一个点让所有人都在等。

如何处理迟到的数据

对于越过水位线才到达的数据,有两种选择。

丢弃。已确定的数字绝对不会改变。但相应的部分会悄悄消失,所以一定要单独统计丢弃的记录数和金额。如果不统计就丢弃,以后就没有依据来解释“为什么我们的合计对不上”。

作为迟到更新加以反映。同一份 Flink 文档中关于窗口的说明指出:设置允许延迟之后,迟到的元素可以让窗口再次触发,而此时输出的值应当被当作前一次计算的更新结果来处理。如果下游不能把它当作更新来接收,只会追加,就会产生重复。所以这个选择不只是我们这一方的决定。

无论选哪一种,关键都在于把关闭那一刻的值和之后的更正分开保存。合并成一个的话,就回答不了“昨天的数字今天变了”,而分开保存的话,变化的量本身就是答案。

在现场相遇的样子

第一,窗口不关闭。数据一旦中断,最大事件时间就停住,水位线也跟着停住,窗口就永远开着。同一份 Flink 文档讲到的空闲(idleness)问题就是这个。一个安静的分区就能拖住整体。

第二,把允许延迟调大之后就忘了。出了事故,就先“宽松一点”地调大,然后就这样固化下来。六个小时之后才出数的仪表板算不上实时。调大的值要连同恢复原值的日期一起记下来。

第三,重处理和迟到数据混在一起。重新运行过去的区间,和反映迟到的数据,这两件事看上去都是“旧数字变了”。两者原因不同,所以也要分开记录。

第四,迟到数据只集中在一边。按地区或终端机型拆开看,只有特定群体的延迟特别大。只看整体平均值的话,这个群体的数据总是被丢弃,而没有人知道这件事。

实际工作中真正重要的事

下一项实验要做什么

生成经过地铁和飞行模式的订单事件,再一步步扩充工具 wm.py。把延迟作为分布来测量,按到达顺序推进水位线,按事件时间划分窗口,把窗口关闭之后才到的数据区分出来。然后对同一份数据分别用丢弃策略和反映策略进行汇总,看哪个窗口变化了多少,并把关闭那一刻的值和更正分开输出。最后用这些数字写一份报告。评分器每次都会用不同的窗口大小和允许延迟真正运行你的工具,并核对答案。