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

数据流水线

昨天的数字今天变了 — 水位线与迟到数据

在 TT Lab 中继续学习

目标

创建按事件时间汇总的工具 wm.py。以分布的形式测量处理时间和事件时间之间的差距,按到达顺序推进水位线,把窗口关闭之后才到的数据区分出来,比较丢弃策略和反映策略的结果,并把关闭那一刻的值和之后的更正分开保存。

为什么重要

本课程前面的实验中已经讲过一次水位线。那是把已加载数据的最大时间写在表里,让下一次运行只读取它之后的数据,而这个水位线回答的问题只有一个:读到哪里了。这里要讲的是再往后的内容。 移动应用丢了信号,几分钟后才发送事件;处于飞行模式的终端,要过几个小时才发送事件。这些事件发生的时间早已过去,而我们已经把那个时间段的汇总结果发出去了。水位线只是“认为那些数据不会再来”的一个承诺,而不是事实。 所以需要决定三件事。允许延迟设为多少;对于违背承诺才到的数据,是丢弃还是作为迟到更新加以反映;以及,如果选择了反映,该如何解释昨天给出的数字和今天给出的数字之间的差异。 这不是一个测量时间的实验。它是用事件中记录的 event_time 和 ingest_time 这两个字段来计算的实验,程序实际运行了几秒完全无关。 评分器不会相信你写出来的文字。它会把评分器生成的事件流放在临时文件里,真正运行你的工具,并一边改变窗口大小和允许延迟,一边核对答案。记录数和延迟分布每次运行都会变化。

步骤

  1. 创建并运行 /root/wmark/gen_stream.py,生成 /root/wmark/work/stream.jsonl。
  2. 在 /root/wmark/wm.py 中实现 skew,输出延迟分布和乱序的记录数。
  3. 加上 watermark,按到达顺序推进水位线。
  4. 加上 windows,按事件时间划分窗口并汇总。
  5. 加上 late,区分窗口关闭之后才到的数据。
  6. 加上 agg,分别给出丢弃策略和反映策略的结果。
  7. 加上 close,分别输出关闭那一刻的值和之后的更正。
  8. 写出 /root/wmark/work/watermark_report.json 和 /root/wmark/work/watermark_report.md。

参考

准备好会迟到的事件

创建并运行 /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 条。如果没有,就减小这个值。报告里要用数字写出迟到的记录数——这一个数字,就是下一次会议上最先被问到的那个问题的答案。直接调用前面步骤中写好的函数就行。