改变延迟,数一数被丢弃的行
目标
把同样的数据分别以批处理和多种延迟的流式方式运行,用数字确认事件时间窗口会随水位线在何时关闭、丢弃哪些行。用 CURRENT_WATERMARK 亲手找出迟到行,并与窗口丢弃的行进行比较。
为什么重要
水位线关闭窗口之后到达的行,会不报任何错误地消失。所以“数字少了一点”这类问题无法从日志中找出,只有了解水位线是如何移动的,才能解释。本实验的输入是秒级时间戳,所以水位线会随每条记录立即推进,结果与运行速度无关,只由到达顺序决定。评分器不会查询集群——它读取你保存的 sql-client 输出,并把原始 CSV 按到达顺序流过,用与引擎相同的规则(先把行发出,再推进水位线 · 丢弃 window_end ≤ 워터마크 的窗口的行,占位符为水位线)计算出的值进行核对。
步骤
- 用
flink-up启动集群,在 /root/flink/watermark/ddl.sql 中写入设置了WATERMARK FOR ts AS ts - INTERVAL '5' SECOND的events表和DESCRIBE events;,并把输出保存到 /root/flink/watermark/ddl.out。 - 运行以批处理模式为每个 1 分钟
TUMBLE窗口输出cnt(条数)和total(reading 之和)的 /root/flink/watermark/batch.sql,把输出保存到 /root/flink/watermark/batch.out。 - 把以流式模式(延迟 5 秒)运行同一聚合的 /root/flink/watermark/w5.sql 的输出保存到 /root/flink/watermark/w5.out。
- 在 /root/flink/watermark/sweep.sql 中创建水位线为
ts的表events_0和水位线为ts - INTERVAL '30' SECOND的表events_30,按events_0→events_30的顺序运行同一聚合,并把输出保存到 /root/flink/watermark/sweep.out。 - 把从 5 秒延迟的
events中取出CURRENT_WATERMARK(ts)不为 NULL 且ts <= CURRENT_WATERMARK(ts)的行的event_id, ts, wm的 /root/flink/watermark/late.sql 的输出保存到 /root/flink/watermark/late.out。 - 把在窗口之前过滤掉迟到行(
CURRENT_WATERMARK(ts) IS NULL OR ts > CURRENT_WATERMARK(ts))之后再做同样 1 分钟聚合的 /root/flink/watermark/filtered.sql 的输出保存到 /root/flink/watermark/filtered.out。 - 把只在第 3 步的聚合上加入
SET 'pipeline.auto-watermark-interval' = '1 h';的 /root/flink/watermark/slow.sql 的输出保存到 /root/flink/watermark/slow.out。 - 在 /root/flink/watermark/report.json 中写入
total_rows、dropped_0、dropped_5、dropped_30、late_rows_5、dropped_slow。
参考
- 源列:
event_id BIGINT, sensor STRING, reading INT, ts TIMESTAMP(3)(没有表头的 CSV,文件/opt/lab/fixtures/data/watermark_events.csv,文件顺序 = 到达顺序)。 - 窗口聚合的形态:
SELECT window_start, window_end, COUNT(*) AS cnt, SUM(reading) AS total FROM TUMBLE(TABLE 표, DESCRIPTOR(ts), INTERVAL '1' MINUTE) GROUP BY window_start, window_end;(占位符为表名) - 一个 SQL 文件中有多个 SELECT 时,会依次运行作业,输出中按顺序打印结果表。
- 常见错误:批处理模式不使用水位线。要查看迟到行,必须以流式运行。并行度请保持默认值 1——如果有多个输入,水位线取其中的最小值。
- 常见错误:第一行时还没有水位线,所以
CURRENT_WATERMARK(ts)是 NULL。与 NULL 的比较不为真,通不过 WHERE,所以在第 6 步的过滤表达式中去掉IS NULL条件,连第一行也会被丢弃。 - 官方文档:Timely Stream Processing · CREATE — WATERMARK · Time Attributes · Windowing TVF · Built-in Functions · Configuration
把 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 的行)为界分开数。