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

三个现象值得单独看:

官方文档怎么规定(检索日期 2026-10-05)

以下是与本机切片对照的规范陈述,来源为 Apache Flink 与 Apache Spark 官方文档,不属于本机实测:

把两边放在一起看:本机切片给出的是「这套规则会产生什么数字」,官方文档给出的是「这套规则在工程产品里被怎么实现和打折」。

生产边界

动手

  1. 跑 watermark_engine.py,把 LATENESS_SPREAD 里的 12 改成 20,重新跑,看 丢弃 和 总金额 怎么变。
  2. 把窗口宽度 WINDOW 从 10 改成 15,说明为什么真值窗口金额不再是 [55,155,255],并验证你的预期。
  3. 断言题:在 L=0 下打印那 3 条被丢弃事件的 event_time,验证它们都落在「窗口已触发且 watermark >= w_end + 0」之后。

自测

  1. 处理时间和事件时间在什么条件下相等?链路里哪三种情况会让它们分叉?
  2. 水位线是「保证」还是「猜测」?它把什么问题换成了什么问题?
  3. allowed_lateness 从 5 调到 12,最终金额不变而重算次数下降,为什么?
  4. 为什么把容忍度调到远大于数据时间跨度没有收益?多出来的延迟和状态代价花在了哪里?
  5. 如果业务方要求「每个窗口的数据绝对不能丢」,你会怎么改这条链路?代价是什么?

↓ 下一步:04 章 · 状态与容错 —— 窗口状态放在内存里,进程崩了怎么办。

进入 keel 阅读