水位线——引擎如何判定“这个窗口现在关闭”
一句话总结
事件时间窗口只有在水位线越过窗口末端之后才会关闭。Flink SQL 的水位线像 WATERMARK FOR ts AS ts - INTERVAL '5' SECOND 这样声明,其值是“到目前为止见到的最大 ts − 延迟”。窗口关闭之后到达的行会被悄悄丢弃。延迟是一个用准确性换取延迟时间的旋钮。
为什么需要它
上一个模块的 GROUP BY user_id 只需要不断改写结果就行。但“每分钟把传感器读数合计后只输出一次”这个需求就不同了。要只输出一次,就必须知道那一分钟已经结束。如果按墙上时钟(处理时间)判断,很简单,但正如官方文档(Timely Stream Processing)所指出的,处理时间会受到达速度、故障、重新处理的影响,是不确定的。把昨天的数据重新跑一遍,结果就会不同。
所以要按记录中写明的事件时间来划分窗口。问题是事件不会按顺序到达。12:00:59 的传感器读数,可能比 12:01:10 的读数更晚到。12:00 这个窗口应该什么时候关闭?不可能永远等下去。水位线就是对这个问题的一个承诺——Watermark(t) 是一份声明:“从现在起,我认为 ts ≤ t 的事件不会再来了。”
工作原理
声明。 根据文档(CREATE Statements),WATERMARK 子句会把一个 TIMESTAMP(3) 列变成事件时间属性。用 DESCRIBE 查看,该列的类型上会带有 *ROWTIME*。常见的策略有三种。
| 表达式 | 含义 |
|---|---|
ts |
严格升序——见到的最大 ts 就是水位线 |
ts - INTERVAL '0.001' SECOND |
升序——与最大 ts 同一时刻的行不算迟到 |
ts - INTERVAL '5' SECOND |
输入顺序被打乱——最多等 5 秒 |
何时发出。 表达式针对每条记录计算,但文档指出,水位线是按 pipeline.auto-watermark-interval(默认 200ms)的周期发出的。不过看一下附加水位线的算子的源码,还有一点——新水位线比上次发出的值领先超过间隔时,不等周期,在那条记录处立即发出。而且顺序很重要:先把记录向下发出,然后再推进水位线。所以某一行看到的水位线,是“到它前面那些行为止的最大 ts − 延迟”。如果是秒级时间戳,推进幅度总是大于 200ms,每条记录都会发出水位线,结果与运行速度无关,完全确定。这就是实验评分器仅凭到达顺序就能复现结果的原因。
窗口关闭的条件。 窗口聚合一旦 window_end ≤ 워터마크(占位符为水位线),就会把该窗口的结果输出一次,并清空状态。之后再来应该进入该窗口的行,因为没有地方接收而被丢弃。既没有错误,也没有警告。读完有终点的文件后,引擎会发送最大水位线,关闭所有剩余的窗口。
行迟到与窗口关闭是两回事。 CURRENT_WATERMARK(ts) 返回该行所经过算子的当前水位线(文档的 Built-in Functions,尚未产生时为 NULL)。ts <= CURRENT_WATERMARK(ts) 的行,以水位线衡量就是迟到行。但那一行的窗口可能仍然开着——12:00:58 的行即使在水位线 12:00:59 之后到来,12:00 窗口的末端(12:01:00)仍在水位线右侧。所以如果按文档的过滤表达式(CURRENT_WATERMARK(ts) IS NULL OR ts > CURRENT_WATERMARK(ts))在窗口之前过滤,丢弃的行会比交给窗口处理时更多。用实验数据测量,5 秒延迟下窗口丢弃的行是 15 条,而按水位线算迟到的行是 98 条。
并行时。 接收多个输入的算子,其事件时间是各输入水位线中的最小值(文档的 Watermarks in Parallel Streams)。只要有一个分区安静下来,整体水位线就会停住。为了解决这个问题,有一个用 table.exec.source.idle-timeout 暂时把安静的 Source 排除在外的配置。本实验的并行度是 1,所以看不到这个效果。
在现场相遇的样子
“聚合数字比原始数据少一点”的反馈,有相当一部分是迟到行。被丢弃的行连日志里都不会留下,所以第一步诊断,是把同样的数据用批处理跑一遍,测量差距。看差距是否会随着延迟增大而减小,就能分清原因。增大延迟的代价是窗口结果相应晚出——仪表板晚 5 秒可不可以、晚 30 秒行不行,是业务要决定的问题。
第二种是“结果根本不出来”。水位线不推进,窗口就永远不会关闭。可能是某个分区空着,可能是 Source 的 ts 为 NULL,也可能是测试数据的时间集中在一个点上。把水位线间隔设得很大,也会发生类似的事——实验中把间隔设为 1 小时,水位线在第一个值之后就停止推进,由于是有终点的文件,最后所有窗口一起关闭,得到的结果里没有一行被丢弃。如果是无限流,窗口会一个小时都关不上。
第三种是要求不要丢弃迟到行,而是单独收集起来。在 SQL 中,常见的做法是用 CURRENT_WATERMARK 标记迟到行,再发送到另一个 Sink。不过正如前面所见,“早于水位线的行”和“窗口已关闭因而会被丢弃的行”是不同的集合。必须先决定收集哪一个。
下一项实验要做什么
把 600 条传感器事件(文件顺序就是到达顺序)声明为 5 秒延迟水位线,并用 DESCRIBE 确认。以批处理方式运行 1 分钟 TUMBLE 聚合,建立基线,再与延迟 5 秒、0 秒、30 秒的流式结果比较,数出被丢弃的行。用 CURRENT_WATERMARK 找出迟到行,看在窗口之前把这些行过滤掉的结果有何不同,再把水位线间隔改成 1 小时,确认推进停止后有什么变化,最后把数字整理成报告。