昨天的数字今天变了 — 水位线与迟到数据
目标
创建按事件时间汇总的工具 wm.py。以分布的形式测量处理时间和事件时间之间的差距,按到达顺序推进水位线,把窗口关闭之后才到的数据区分出来,比较丢弃策略和反映策略的结果,并把关闭那一刻的值和之后的更正分开保存。
为什么重要
本课程前面的实验中已经讲过一次水位线。那是把已加载数据的最大时间写在表里,让下一次运行只读取它之后的数据,而这个水位线回答的问题只有一个:读到哪里了。这里要讲的是再往后的内容。 移动应用丢了信号,几分钟后才发送事件;处于飞行模式的终端,要过几个小时才发送事件。这些事件发生的时间早已过去,而我们已经把那个时间段的汇总结果发出去了。水位线只是“认为那些数据不会再来”的一个承诺,而不是事实。 所以需要决定三件事。允许延迟设为多少;对于违背承诺才到的数据,是丢弃还是作为迟到更新加以反映;以及,如果选择了反映,该如何解释昨天给出的数字和今天给出的数字之间的差异。 这不是一个测量时间的实验。它是用事件中记录的 event_time 和 ingest_time 这两个字段来计算的实验,程序实际运行了几秒完全无关。 评分器不会相信你写出来的文字。它会把评分器生成的事件流放在临时文件里,真正运行你的工具,并一边改变窗口大小和允许延迟,一边核对答案。记录数和延迟分布每次运行都会变化。
步骤
- 创建并运行 /root/wmark/gen_stream.py,生成 /root/wmark/work/stream.jsonl。
- 在 /root/wmark/wm.py 中实现
skew,输出延迟分布和乱序的记录数。 - 加上
watermark,按到达顺序推进水位线。 - 加上
windows,按事件时间划分窗口并汇总。 - 加上
late,区分窗口关闭之后才到的数据。 - 加上
agg,分别给出丢弃策略和反映策略的结果。 - 加上
close,分别输出关闭那一刻的值和之后的更正。 - 写出 /root/wmark/work/watermark_report.json 和 /root/wmark/work/watermark_report.md。
参考
- 所有工作都在
/root/wmark下进行。数据是/root/wmark/work/stream.jsonl。 - 数据的每一行是一个 JSON 对象,包含
id、key、event_time、ingest_time、amount。两个时间都是 epoch 秒整数。也可以有其他字段。 - 文件中记录的顺序就是到达顺序。如果按事件时间重新排序,本实验的所有判定都会失效。
- 运行契约:
python3 /root/wmark/wm.py <명령> <파일> [--size=초] [--lateness=초] [--policy=drop|update](占位符依次为命令、文件、秒数、秒数)。结果以一个 JSON 对象输出到标准输出。成功时退出码为 0,文件不存在时为 3,用法或策略名称错误时为 2。 - 窗口是固定大小、互不重叠的。事件的窗口起点是
event_time - (event_time % size),窗口的范围是从起点到“起点加上大小”之前的一刻。JSON 的键是把窗口起点写成字符串。 - 延迟是
ingest_time - event_time。分位数按最近秩取——把值按升序排列,取第ceil(건수 * p / 100)(占位符为记录数)个(从 1 开始)。不做插值。 skew <파일>的响应是{"events": 정수, "out_of_order": 정수, "min_lag": 정수, "max_lag": 정수, "p50_lag": 정수, "p95_lag": 정수}(占位符依次为文件;整数,共六个)。out_of_order 是事件时间早于在它之前到达的所有事件的最大事件时间的记录数。watermark <파일> --lateness=<초>的响应是{"lateness": 정수, "max_event_time": 정수, "advances": 정수, "final_watermark": 정수}(占位符依次为文件、秒数;整数,共四个)。advances 是使最大事件时间被刷新的到达记录数。windows <파일> --size=<초>的响应是{"size": 정수, "count": 정수, "windows": {"창시작": {"events": 정수, "amount": 정수}}}(占位符依次为文件、秒数;整数、整数、窗口起点、整数、整数)。不区分迟到与否,全部计入。late <파일> --size=<초> --lateness=<초>的响应是{"size": 정수, "lateness": 정수, "on_time": 정수, "late": 정수, "late_by_window": {"창시작": 정수}}(占位符依次为文件、秒数、秒数;整数、整数、整数、整数、窗口起点、整数)。判断某个事件是否迟到,只用在它之前到达的那些数据算出的水位线。如果水位线大于或等于该窗口的末尾,就是迟到。第一个到达的事件没有可比较的前序数据,所以不算迟到。agg <파일> --size --lateness --policy=drop|update的响应是{"policy": 문자열, "size": 정수, "lateness": 정수, "windows": {...}, "dropped": 정수, "restated": [창시작 문자열 오름차순]}(占位符依次为文件;字符串、整数、整数、按升序排列的窗口起点字符串)。如果是 drop,就去掉迟到数据并用 dropped 计数,restated 为空列表。如果是 update,迟到数据也放进去,dropped 为 0,并把放入了迟到数据的窗口记入 restated。close <파일> --size --lateness的响应是{"size": 정수, "lateness": 정수, "sealed": {...}, "corrections": [{"window": 창시작, "delta_events": 정수, "delta_amount": 정수}], "final": {...}}(占位符依次为文件;整数、整数、窗口起点、整数、整数)。sealed 是只包含关闭之前到达的数据的值(只有迟到数据的窗口置为 0),corrections 按窗口起点升序排列,final 是两者相加的值。- 报告 JSON 中包含 size、lateness、events、windows、on_time、late、dropped、restated、max_lag、p95_lag、policy。lateness 取 1 以上、小于 max_lag 的值,并且在该值下迟到的数据必须至少有 1 条。如果没有,就减小允许延迟。
- 报告 MD 的各节标题是
## 무엇을 재었나(韩文,意为“测量了什么”)、## 허용 지연을 얼마로 잡았나(韩文,意为“允许延迟设为多少”)、## 늦게 온 자료를 어떻게 했나(韩文,意为“如何处理迟到的数据”)、## 어제 숫자가 바뀐 이유(韩文,意为“昨天的数字为什么变了”)。 - 官方文档:Flink Generating Watermarks · Flink Windows · python json
- 常见错误:把数据按事件时间重新排序(迟到数据会显示为 0 条)、让水位线倒退、不统计丢弃的记录数、把窗口起点按字符串排序。
准备好会迟到的事件
创建并运行 /root/wmark/gen_stream.py,生成 /root/wmark/work/stream.jsonl。要求至少 60 条,文件中的顺序按 ingest_time 升序排列,每一行的 ingest_time 都大于或等于 event_time,延迟达到 120 秒以上的行至少有 5 条,事件时间的跨度至少为 1200 秒,key 至少有 3 种。
如果延迟只按一种分布生成,后面就没有可看的东西。要分成三路散开:大部分在几秒内到达,一部分晚几分钟,极少数在几十分钟之后才涌进来。生成完之后按 ingest_time 排序写入文件,这个顺序就是到达顺序。要固定随机数种子,这样在改变允许延迟来比较的过程中,数据才不会变动。
先量一量差多少
在 /root/wmark/wm.py 中实现 skew <파일>(占位符为文件),以 JSON 输出记录数、乱序的记录数,以及延迟的最小值、最大值、中位数和 95 分位数。
延迟是 ingest_time - event_time。分位数按最近秩取,不要插值——结果必须是整数,评分和开会时才不会出现分歧。乱序的记录数,按到达顺序遍历,数一数事件时间早于当前已见最大事件时间的那些就行。
按到达顺序推进水位线
加上 watermark <파일> --lateness=<초>(占位符依次为文件、秒数),输出 lateness、max_event_time、advances、final_watermark。advances 是使最大事件时间被刷新的到达记录数。
水位线不会倒退。如果让迟到的事件把最大事件时间拉低,已经关闭的窗口就会反复重新打开又关闭,任何值都无法确定。第一个到达的事件没有可比较的前序数据,所以它本身就算一次前进。
按事件时间划分窗口
加上 windows <파일> --size=<초>(占位符依次为文件、秒数),输出每个窗口的记录数和金额。窗口起点是 event_time - (event_time % size),这一步不区分迟到与否,全部计入。
JSON 的键必须是字符串,所以把窗口起点写成字符串。之后排序的时候,不要按字符串,而要按整数来比较——位数不同时,按字符串排序会得出错误的顺序。
把窗口关闭之后到达的数据挑出来
加上 late <파일> --size=<초> --lateness=<초>(占位符依次为文件、秒数、秒数),输出 on_time、late、late_by_window。判断某个事件是否迟到,只用在它之前已经到达的那些数据算出的水位线;如果水位线大于或等于该窗口的末尾,就是迟到。
“迟到”不是数据的属性,而是到达顺序的属性。为了方便处理而按事件时间排了一次序的话,到达顺序就消失了,迟到数据会显示为 0 条——不是因为没有问题,而是因为把能看见问题的眼睛去掉了。要按文件中记录的顺序遍历,而且最大事件时间的刷新要放在判定完成之后。第一个到达的事件没有可比较的前序数据,所以不算迟到。
是丢弃还是反映
加上 agg <파일> --size --lateness --policy=drop|update(占位符为文件)。如果是 drop,就去掉迟到数据并用 dropped 计数,restated 为空列表;如果是 update,迟到数据也放进去,dropped 为 0,并把放入了迟到数据的窗口记入 restated。
即使选择丢弃,也一定要统计丢弃的记录数。如果不统计就丢弃,以后就没有依据来解释合计为什么对不上。restated 按窗口起点升序排列,比较时按整数比较。遇到不认识的策略名称,就以退出码 2 结束。
把关闭那一刻的值和更正分开保存
加上 close <파일> --size --lateness(占位符为文件),输出 sealed(只包含关闭之前到达的数据)、corrections(按窗口起点升序排列的更正列表)和 final(两者相加的值)。只有迟到数据的窗口,其 sealed 为 0。
如果把两者合在一起只输出一个,就没有办法解释昨天的数字今天为什么变了。分开保存的话,变化的量本身就是答案。corrections 里只放真正有更正的窗口,final 则放所有窗口。
用一页纸解释变掉的数字
在 /root/wmark/work/watermark_report.json 中写入 size、lateness、events、windows、on_time、late、dropped、restated、max_lag、p95_lag、policy,并在 /root/wmark/work/watermark_report.md 中写成四节,标题分别是 ## 무엇을 재었나(韩文,意为“测量了什么”)、## 허용 지연을 얼마로 잡았나(韩文,意为“允许延迟设为多少”)、## 늦게 온 자료를 어떻게 했나(韩文,意为“如何处理迟到的数据”)、## 어제 숫자가 바뀐 이유(韩文,意为“昨天的数字为什么变了”)。
允许延迟取 1 以上、小于最大延迟的值,并且在该值下迟到的数据必须至少有 1 条。如果没有,就减小这个值。报告里要用数字写出迟到的记录数——这一个数字,就是下一次会议上最先被问到的那个问题的答案。直接调用前面步骤中写好的函数就行。