跳到内容

9.3 事件时间、Watermark 与端到端语义:流处理怎样面对迟到与重放

批处理报告每天清晨生成,但前线希望在装备温度连续十分钟异常时立即告警。网络中断后,昨晚的传感器事件今天才到;作业重启又重放了一批消息。若按到达时间统计,历史会漂移;若只求“每条恰好处理一次”,下游告警接口仍可能收到重复请求。

流处理的难点不是循环读取消息,而是对无界、乱序、可重放数据维持明确的时间和状态语义。

本课目标

  • 区分事件时间、摄取时间和处理时间;
  • 用窗口与 watermark 管理乱序和状态;
  • 解释 at-most-once、at-least-once 与端到端 exactly-once 的边界;
  • 设计可重放、幂等和可观测的流式管道。

1. 无界数据没有“全部到齐”

批表有一个可枚举快照,流则持续产生。系统必须决定:

  • 何时触发计算;
  • 一个窗口何时可输出初步结果;
  • 迟到数据到来时更新、撤回还是丢弃;
  • 状态保留多久;
  • 失败后从哪里恢复。

流处理常以 micro-batch 或逐记录方式执行,但两者都要回答这些语义问题。

2. 三个时钟

  • 事件时间:事件在来源系统中发生的时间;
  • 摄取时间:平台首次接收到事件的时间;
  • 处理时间:某个算子实际处理事件的时间。

按处理时间做十分钟窗口容易实现,但重放历史时结果会改变。业务统计通常使用事件时间,同时保留摄取时间以测量延迟。来源时钟漂移、时区和设备重启也必须进入质量监控。

3. 窗口定义是业务规则

  • tumbling window:不重叠固定窗口;
  • sliding window:按滑动步长重叠;
  • session window:按活动间隔动态闭合。

“过去十分钟”要说明边界、时区、触发频率和结果更新方式。一个事件可进入多个 sliding windows,状态与计算成本随重叠程度增加。

Structured Streaming 风格示例:

python
from pyspark.sql import functions as F

counts = (
    events
    .withWatermark("event_time", "30 minutes")
    .groupBy(
        F.window("event_time", "10 minutes"),
        F.col("equipment_id"),
    )
    .agg(F.avg("temperature").alias("avg_temperature"))
)

具体输出模式和 sink 是否允许更新,要与查询类型一起确认。

4. Watermark 不是“当前时间减 30 分钟”

watermark 向引擎表达对事件时间进度和允许迟到程度的策略,使系统可以清理旧窗口状态。它通常与已观察到的事件时间进展有关,不应想当然地等于墙钟时间。

设置过短:更多迟到事件无法更新结果;设置过长:状态、checkpoint 和恢复成本增加。应根据真实延迟分布、业务修正窗口和资源约束选择,并监控 watermark 进度与被拒绝/遗漏的迟到事件。

watermark 也不是完整性证明。来源可能停滞、某分区可能延迟,多个输入流还需要定义全局进度策略。

5. 去重需要稳定事件身份

消息系统重投、producer 重试和回放都会产生重复。应优先使用来源生成的稳定 event_id,并定义身份作用域和保留期限。

仅按整行值去重会误删内容相同但真实发生两次的事件;仅在无限历史上保存所有 ID 又会让状态无界。结合事件 ID、事件时间和 watermark 管理去重状态,并记录超出保留窗口的重复怎样处理。

业务更新还要区分:重复、修订、撤销和晚到新事件。它们不是同一种操作。

6. 处理次数不是端到端结果语义

常见术语:

  • at-most-once:可能丢失,但不重试造成重复;
  • at-least-once:不丢为目标,失败重试可能重复;
  • exactly-once:在明确系统边界内,每个输入对状态/结果只生效一次。

端到端保证取决于整个链条:可重放 source、进度记录、确定性转换、checkpoint、sink 的幂等或事务提交。引擎内部宣称 exactly-once,不代表任意外部 HTTP API 也自动恰好执行一次。

一个实用写出模式是使用批次/事件幂等键:

text
begin transaction
  if output_id not committed:
      write result rows
      record output_id
commit

结果与提交标记必须原子写入同一事务边界;先写结果再单独记标记仍有崩溃窗口。

7. Checkpoint 记录执行进度

checkpoint 常包含 source offsets、状态和查询元数据,使作业失败后继续。它不是长期业务数据备份,也不保证能跨任意代码和 schema 变更恢复。

部署变更前要确认:

  • 新代码是否兼容旧状态 schema;
  • checkpoint 路径是否唯一且持久;
  • source 数据保留期能否覆盖最长故障;
  • sink 可否安全接受重放;
  • 无法原地升级时怎样从原始日志重建状态。

保留可重放的原始事件日志,是修复逻辑错误和回填结果的关键能力。

8. Backpressure 与运行健康

输入速率长期高于处理速率时,lag 会增长。监控至少包括:

  • source 最新 offset 与已处理 offset;
  • 每批输入行数、处理时长和调度间隔;
  • watermark/event-time lag;
  • state rows、state bytes 和 checkpoint 时长;
  • 迟到、重复、解析失败与死信数量;
  • sink 延迟、错误和幂等冲突。

扩容前先判断瓶颈是 source 限流、shuffle、热点键、状态膨胀还是 sink 吞吐。无限增加消费并行度可能只把压力推给下游。

9. 迟到结果怎样对外发布

可选择:

  • 追加:窗口足够确定后只写一次;延迟更高;
  • 更新/upsert:迟到数据到来时更新同一结果键;
  • 撤回/变更日志:发出旧值撤回与新值;
  • 初值加终值:低延迟初步结果,随后封账并标注版本。

下游必须理解所选语义。把会更新的流写进只支持 append 的表,会产生多个互相冲突的“最终值”。

10. 测试与演练

构造包含以下情况的确定性事件序列:

text
乱序但在 watermark 内
超过允许迟到范围
相同 event_id 重投
同内容不同 event_id
作业在写出前后崩溃
source 分区暂时停滞
schema 增列与不兼容改型
从旧 offset 全量重放

验证的不只是单次输出,还包括重启后状态、重复副作用和最终收敛结果。

常见误区

  • 流处理就是更快的批处理:无界输入需要进度、状态和修正语义。
  • watermark 前的数据一定完整:它是迟到与状态的工程策略,不是事实证明。
  • 引擎 exactly-once 等于端到端 exactly-once:sink 与外部副作用决定最终边界。
  • checkpoint 可以替代原始日志:逻辑错误和不兼容升级仍需要重放源。

练习

  1. 为十分钟温度告警定义事件时间、窗口、watermark 和更新语义。
  2. 生成乱序与重复事件,验证窗口结果最终是否收敛。
  3. 设计一个对作业重试安全的数据库 upsert 键。
  4. 为状态持续增长的作业画出 lag、watermark 和 state size 排障路径。

小结

可靠流处理要把时间、状态和失败当成正常输入。watermark 决定迟到与资源的权衡,checkpoint 支持恢复,而端到端交付保证必须跨过 source、引擎和 sink 的完整边界。

下一章进入数据伦理:系统能够收集、推断和自动决策,并不等于它应该这样做。

Built with VitePress | Software Systems Atlas