KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
03 · 事件时间与水位线 — keel 龙骨
前两章把数据搬到了 topic 里。这一章开始算。算的第一件事不是聚合,是用哪个时间轴来切窗口——这个选择错了,后面所有数字都会对不上。
前两章把数据搬到了 topic 里。这一章开始算。算的第一件事不是聚合,是用哪个时间轴来切窗口——这个选择错了,后面所有数字都会对不上。
现场
运营对两个口径的数字吵起来了。技术侧说「12:00 到 12:10 这一档的订单额是 40 万」,运营侧拿着业务系统里的截图说「明明是 55 万」。
两边都没算错。技术侧用的是处理时间:事件到达流处理引擎的先后顺序,按每 10 条切一档。运营侧用的是事件时间:订单在自己的系统里产生的时间。一笔 12:03 下单的记录,因为上游重试,12:11 才到引擎——它被算进了后一档。
只要链路里有任何一处会产生延迟(网络、消费者重启、上游批量提交),处理时间和事件时间就会分叉。分叉的大小取决于延迟的分布,而不是数据量。
全链路
flowchart TD
E["事件到达<br/>(event_time, key, amount)"] --> AS["时间戳分配<br/>取 event_time"]
AS --> MW["跟踪已见最大事件时间<br/>max_ts = max(max_ts, ts)"]
MW --> WM["水位线<br/>watermark = max_ts - allowed_lateness"]
AS --> WA["按 event_time 归入滚动窗口<br/>w_end = ceil(ts/10)*10"]
WM --> FW{"watermark >= 窗口结束时间?"}
WA --> FW
FW -->|"否"| KEEP["窗口状态保留<br/>继续累计"]
FW -->|"是"| FIRE["触发:输出窗口结果"]
FIRE --> SINK["下游 sink"]
WA --> LATE{"事件落在已触发窗口?"}
LATE -->|"是,且 watermark < w_end + L"| UPD["并入并重新输出<br/>(late-update)"]
LATE -->|"是,且 watermark >= w_end + L"| DROP["丢弃<br/>(窗口已最终确定)"]
UPD --> FIRE
DROP -. "这批数据永久丢失" .-> LOSS["下游金额偏小"]
KEEP -. "状态一直涨:要等水位线推进才能清" .-> MEM["状态膨胀"]
style E fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
style AS fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style MW fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style WM fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style WA fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style FW fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style KEEP fill:#f3f0ff,stroke:#7048e8,color:#1a1a2e
style FIRE fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style SINK fill:#f3f0ff,stroke:#7048e8,color:#1a1a2e
style LATE fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style UPD fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style DROP fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
style LOSS fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
style MEM fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
精确定义
事件时间是事件里自带的时间戳,由产生它的系统写入,永远不随传输改变。处理时间是引擎处理这条事件的墙上时钟。两者只有在「零延迟」时才相等。
水位线(watermark)是一句承诺:「我认为不会再有 event_time 小于 W 的事件了。」 它不是「数据的最大时间」,而是「最大时间减去一个容忍度」——一句话就把不确定性换成了延迟。
**允许迟到(allowed lateness)**是水位线之后的第二道缓冲:窗口触发之后,状态不立刻扔,继续留一段时间接受迟到事件并重算。
**触发(trigger)**是「什么时候把结果发出去」。触发时机和窗口是否关闭是两件事:一个窗口可以触发多次(迟到更新),只关闭一次。
本地切片
本机跑不了 Flink / Spark(没有离线依赖,下载整套 JVM 生态不现实),所以用一段约 150 行的 Python 把同一套语义跑出来,代码在 lab/exp/stream-data-pipeline/watermark_engine.py。它做的事和图一致:跟踪 max_ts、算水位线、按 event_time 归窗、水位线越过窗口结束就触发、触发后但在最终确定前接受迟到并重算、再晚就丢。
输入是故意乱序的 30 条事件(每秒一条,到达顺序被 0~12 秒的随机延迟打乱),真值窗口金额应该是 [55,155,255],总计 465。
一次完整运行
原始输出见 lab/evidence/stream-data-pipeline/04-watermark.txt。
A. 事件时间 + 水位线(allowed_lateness=0):
触发/更新记录(窗口结束, 累计, 触发时水位线, 类型):
window_end=10 sum=40 watermark=10 fire
window_end=20 sum=138 watermark=20 fire
window_end=30 sum=255 watermark=29 flush
最终窗口:[0,10)=40 [10,20)=138 [20,30)=255
迟到被丢弃=3,迟到被并入=0
和真值 [55,155,255] 比,第一个窗口少了 15。这 3 条被丢的事件迟到超过了容忍度——它们到达时窗口 [0,10) 已经触发且水位线已越过它的最终确定线。这就是水位线的代价:它把「等到齐」换成了「等有限时间」。
B. 同样数据用处理时间切窗:
到达第 0~9 条: sum=63 覆盖的事件时间 0..11
到达第 10~19 条:sum=151 覆盖的事件时间 6..20
到达第 20~29 条:sum=251 覆盖的事件时间 16..29
两种口径的窗口金额对比:
事件时间窗口金额: [40, 138, 255]
处理时间窗口金额: [63, 151, 251]
两组数字完全对不上。处理时间窗口的边界和业务时间没有任何关系,它只是「到达顺序的切片」——每个窗口都混进了上一档和下一档的事件。这是本章最该记住的一组数字:同一批数据,换个时间口径,结果从 [40,138,255] 变成 [63,151,251]。
C. allowed_lateness 取不同值:
allowed_lateness |
丢弃 | 迟到并入 | 窗口金额 | 总金额 |
|---|---|---|---|---|
| 0 | 3 | 0 | [40, 138, 255] |
433 |
| 2 | 1 | 1 | [55, 138, 255] |
448 |
| 5 | 0 | 1 | [55, 155, 255] |
465 |
| 12 | 0 | 0 | [55, 155, 255] |
465 |
| 30 | 0 | 0 | [55, 155, 255] |
465 |
三个现象值得单独看:
- L 从 0 调到 5,第一个窗口从 40 回到 55,总金额从 433 回到 465(真值)。 迟到容忍度确实买到了正确性。
- L=12 时
迟到并入=0,结果却和 L=5 一样。 因为水位线被压低了 12 秒,窗口触发得更晚,那几条「迟到」事件到达时窗口还没触发——它们根本没被归类为迟到。调大 L 不只是更宽容,它让窗口更晚触发、更少重算,代价是结果出得更慢。 - L 调到 30(比数据跨度还大)不再有收益。 总金额停在 465,说明正确性已经在 L=5 就买到了;再往上加,只增加延迟和状态保留时间,不改变任何数字。
官方文档怎么规定(检索日期 2026-10-05)
以下是与本机切片对照的规范陈述,来源为 Apache Flink 与 Apache Spark 官方文档,不属于本机实测:
- Flink 的
WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(20))对应「有界乱序」策略;其内置BoundedOutOfOrdernessGenerator发出的水位线是currentMaxTimestamp - maxOutOfOrderness - 1(Flink 1.20 文档)。本机切片用的是max_ts - L,语义相同,只是没有-1这个边界处理。 - Flink 的 checkpoint 与窗口状态、时间戳、水位线一起被快照;窗口状态的清理依赖水位线推进(Flink 1.20 Checkpointing 文档)。
- Spark Structured Streaming 的
withWatermark(eventTime, delayThreshold):对一个结束于 T 的窗口,引擎保留状态并允许迟到更新,直到max event time seen - late threshold > T,超过阈值的数据会被丢弃(Spark 3.5 / 4.0 文档)。这与本机切片里「窗口最终确定的判据watermark >= w_end + L」是同一个不等式。 - Spark 文档还特别声明:由于跨分区协调
max(eventTime)有成本,实际使用的水位线只保证至少落后delayThreshold,仍可能处理到更晚的记录(withWatermarkAPI 文档)。也就是说「超过阈值一定丢」是对外的一种保守承诺,不是严格保证。
把两边放在一起看:本机切片给出的是「这套规则会产生什么数字」,官方文档给出的是「这套规则在工程产品里被怎么实现和打折」。
生产边界
- 教学替身 vs 真实依赖:本机是单进程、单分区、内存里的 dict 状态。真实引擎的状态在状态后端里(内存 / RocksDB),水位线在并行算子间按各分区最小水位线合并——并行度一高,「全局水位线」会被最慢的那个分区拖住,本机看不到这个效应。
- 要盯的指标:水位线相对最新事件时间的滞后(等价于「结果延迟」)、被丢弃的迟到事件数、活跃窗口数与状态大小。只看吞吐会漏掉「状态一直在涨」。
- 失败策略:丢弃迟到数据要落一条可观测的计数(否则没人知道口径偏了);对账敏感的场景,宁可把 L 调大、把结果延迟拉长,也要让
丢弃=0。 - 口径要写进契约:窗口宽度、水位线容忍度、迟到是否计入,这三件事必须和业务方签字确认。它们不是技术参数,是报表口径。
动手
- 跑
watermark_engine.py,把LATENESS_SPREAD里的12改成20,重新跑,看丢弃和总金额怎么变。 - 把窗口宽度
WINDOW从 10 改成 15,说明为什么真值窗口金额不再是[55,155,255],并验证你的预期。 - 断言题:在
L=0下打印那 3 条被丢弃事件的event_time,验证它们都落在「窗口已触发且watermark >= w_end + 0」之后。
自测
- 处理时间和事件时间在什么条件下相等?链路里哪三种情况会让它们分叉?
- 水位线是「保证」还是「猜测」?它把什么问题换成了什么问题?
allowed_lateness从 5 调到 12,最终金额不变而重算次数下降,为什么?- 为什么把容忍度调到远大于数据时间跨度没有收益?多出来的延迟和状态代价花在了哪里?
- 如果业务方要求「每个窗口的数据绝对不能丢」,你会怎么改这条链路?代价是什么?
↓ 下一步:04 章 · 状态与容错 —— 窗口状态放在内存里,进程崩了怎么办。