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

Apache Flink — 用真正的引擎跑流处理

改变延迟,数一数被丢弃的行

在 TT Lab 中继续学习

目标

把同样的数据分别以批处理和多种延迟的流式方式运行,用数字确认事件时间窗口会随水位线在何时关闭、丢弃哪些行。用 CURRENT_WATERMARK 亲手找出迟到行,并与窗口丢弃的行进行比较。

为什么重要

水位线关闭窗口之后到达的行,会不报任何错误地消失。所以“数字少了一点”这类问题无法从日志中找出,只有了解水位线是如何移动的,才能解释。本实验的输入是秒级时间戳,所以水位线会随每条记录立即推进,结果与运行速度无关,只由到达顺序决定。评分器不会查询集群——它读取你保存的 sql-client 输出,并把原始 CSV 按到达顺序流过,用与引擎相同的规则(先把行发出,再推进水位线 · 丢弃 window_end ≤ 워터마크 的窗口的行,占位符为水位线)计算出的值进行核对。

步骤

  1. 用 flink-up 启动集群,在 /root/flink/watermark/ddl.sql 中写入设置了 WATERMARK FOR ts AS ts - INTERVAL '5' SECOND 的 events 表和 DESCRIBE events;,并把输出保存到 /root/flink/watermark/ddl.out。
  2. 运行以批处理模式为每个 1 分钟 TUMBLE 窗口输出 cnt(条数)和 total(reading 之和)的 /root/flink/watermark/batch.sql,把输出保存到 /root/flink/watermark/batch.out。
  3. 把以流式模式(延迟 5 秒)运行同一聚合的 /root/flink/watermark/w5.sql 的输出保存到 /root/flink/watermark/w5.out。
  4. 在 /root/flink/watermark/sweep.sql 中创建水位线为 ts 的表 events_0 和水位线为 ts - INTERVAL '30' SECOND 的表 events_30,按 events_0 → events_30 的顺序运行同一聚合,并把输出保存到 /root/flink/watermark/sweep.out。
  5. 把从 5 秒延迟的 events 中取出 CURRENT_WATERMARK(ts) 不为 NULL 且 ts <= CURRENT_WATERMARK(ts) 的行的 event_id, ts, wm 的 /root/flink/watermark/late.sql 的输出保存到 /root/flink/watermark/late.out。
  6. 把在窗口之前过滤掉迟到行(CURRENT_WATERMARK(ts) IS NULL OR ts > CURRENT_WATERMARK(ts))之后再做同样 1 分钟聚合的 /root/flink/watermark/filtered.sql 的输出保存到 /root/flink/watermark/filtered.out。
  7. 把只在第 3 步的聚合上加入 SET 'pipeline.auto-watermark-interval' = '1 h'; 的 /root/flink/watermark/slow.sql 的输出保存到 /root/flink/watermark/slow.out。
  8. 在 /root/flink/watermark/report.json 中写入 total_rows、dropped_0、dropped_5、dropped_30、late_rows_5、dropped_slow。

参考

把 ts 声明为事件时间

用 flink-up 启动集群,在 /root/flink/watermark/ddl.sql 中写入读取源 CSV 的 events 表(列 event_id BIGINT, sensor STRING, reading INT, ts TIMESTAMP(3),WATERMARK FOR ts AS ts - INTERVAL '5' SECOND)和 DESCRIBE events;,并把输出保存到 /root/flink/watermark/ddl.out。

WATERMARK 子句放在列列表之内、最后一列之后。如果 DESCRIBE 的结果中 ts 的类型旁边带有 ROWTIME,并且 watermark 一栏能看到表达式,就说明它已经成为事件时间属性。filesystem 连接器的 path 是 file:///opt/lab/fixtures/data/watermark_events.csv。

用批处理建立基线

在 /root/flink/watermark/batch.sql 中写入 SET 'execution.runtime-mode' = 'batch';、第 1 步的 events 表,以及为每个 1 分钟 TUMBLE 窗口输出 window_start, window_end, COUNT(*) AS cnt, SUM(reading) AS total 的聚合,并把输出保存到 /root/flink/watermark/batch.out。

批处理会在输入全部到齐之后再计算,所以不会因水位线而丢弃行。因此这个结果就是“如果一条迟到行都没有”的基线。在批处理中,TUMBLE 也可以用于 TIMESTAMP 列。把 cnt 全部加起来,必须等于源数据的行数。

延迟 5 秒的流式运行——被丢弃的行

创建用 SET 'execution.runtime-mode' = 'streaming'; 运行与第 2 步相同聚合的 /root/flink/watermark/w5.sql(水位线保持 5 秒延迟),并把输出保存到 /root/flink/watermark/w5.out。

窗口在 window_end 小于等于水位线的那一刻输出一次结果,并清空状态。之后再来应该进入该窗口的行,就会被丢弃。请把结果表中 cnt 的合计与批处理比较。窗口结果只以 +I 的形式输出。

一次比较延迟 0 秒和 30 秒

在 /root/flink/watermark/sweep.sql 中创建水位线为 ts 的表 events_0 和水位线为 ts - INTERVAL '30' SECOND 的表 events_30(列和 Source 与 events 相同),以流式模式依次运行同一个 1 分钟聚合,先 events_0,后 events_30,并把输出保存到 /root/flink/watermark/sweep.out。

水位线表达式是按表附加的,所以要改变延迟,就要另建表。一个文件中的两个 SELECT 会按顺序作为两个作业运行,输出中按那个顺序打印结果表。延迟为 0 时,稍微晚一点就会被丢弃,30 秒的话,大部分都会等到。

用 CURRENT_WATERMARK 找出迟到行

创建在流式模式下对 5 秒延迟的 events 运行 SELECT event_id, ts, CURRENT_WATERMARK(ts) AS wm ... WHERE CURRENT_WATERMARK(ts) IS NOT NULL AND ts <= CURRENT_WATERMARK(ts) 的 /root/flink/watermark/late.sql,并把输出保存到 /root/flink/watermark/late.out。

CURRENT_WATERMARK 是该行所经过算子的当前水位线。水位线生成器先发出行再推进水位线,所以一行看到的值是到它前面那些行为止的最大 ts − 5 秒。请把迟到行的数量与第 3 步中被丢弃的行数比较——迟到了并不等于窗口已经关闭。

在窗口之前过滤,会丢弃更多

创建在 5 秒延迟的 events 中只保留 CURRENT_WATERMARK(ts) IS NULL OR ts > CURRENT_WATERMARK(ts) 的行(视图或子查询),然后以流式方式输出同样 1 分钟 TUMBLE 聚合的 /root/flink/watermark/filtered.sql,并把输出保存到 /root/flink/watermark/filtered.out。

这是文档建议在过滤迟到行时使用的表达式。它按行判断“是否早于水位线”,所以连窗口仍然开着、本来会接收的行也丢弃了。要把视图作为窗口 TVF 的输入,写成 CREATE VIEW 之后的 TUMBLE(TABLE 视图, ...)。

把水位线间隔设为 1 小时——推进停止

运行只在第 3 步的 w5.sql 最前面加上一行 SET 'pipeline.auto-watermark-interval' = '1 h'; 的 /root/flink/watermark/slow.sql,并把输出保存到 /root/flink/watermark/slow.out。

水位线按周期(这个配置)发出,在周期之间,只有当新值比上次发出的值领先超过间隔时,才会在记录处立即发出。间隔为 1 小时的话,第一个水位线(此前没有发出过任何水位线,所以立即发出)之后,直到文件结束都不会推进。那么,会有窗口关闭吗?

报告——延迟与被丢弃的行

在 /root/flink/watermark/report.json 中以整数写入 total_rows(batch.out 的 cnt 之和)、dropped_0、dropped_5、dropped_30(批处理合计减去各延迟的 cnt 之和)、late_rows_5(late.out 的行数)、dropped_slow(批处理合计 − slow.out 的 cnt 之和)。

全都可以从保存的输出中数出来。结果行以“| +I |”开头,批处理的结果行以日期开头。用 awk -F'|' 把 cnt 一栏加起来即可。sweep.out 中有两个结果表,请以表头行(含有 op 的行)为界分开数。