跳到内容

17.2 发布订阅、流处理与事件链路观测

事件跨过多个服务后,值守员只看到最终数字不对,却说不清哪一段丢失、重复或延迟。

事件系统不只有一种形态。通知型发布订阅关心“事实发生后谁要响应”,流处理关心“持续事件序列如何被变换、聚合和物化”。两者都使用 broker,不代表状态和时间语义相同。

发布订阅与持久流

维度通知型 Pub/Sub持久事件流
主要目标广播事实、触发反应保存有序记录、持续计算
消费位置常由队列确认控制常由 offset/checkpoint 控制
重放可能有限通常是核心能力
状态多在消费者外部处理器常维护状态存储
典型用途邮件通知、缓存失效排名聚合、实时风控、指标管道

具体产品可能同时支持多种语义,设计应围绕业务需要,而不是围绕产品标签。

流处理必须区分时间

  • Event time:事实在源系统发生的时间;
  • Processing time:处理器实际处理的时间;
  • Ingestion time:平台接收事件的时间。

网络延迟和离线客户端会让事件乱序。按 event time 计算“每五分钟得分”时,需要 watermark 表示系统认为数据大致推进到哪里,并定义迟到事件策略:

  • 允许一定迟到并更新旧窗口;
  • 把过晚事件送入修正流;
  • 只追加更正结果,不静默丢弃;
  • 对外展示结果的最终化程度。

窗口结果何时算“最终”是业务决定,不只是框架参数。

有状态处理需要可恢复 checkpoint

一个排名聚合器维护玩家分数时,状态、输入 offset 和输出必须协调:

text
读取事件 → 更新本地状态 → 写输出 → 提交 checkpoint

任一步崩溃都可能导致重放。处理框架可能提供事务或 checkpoint 机制,但外部数据库、HTTP 副作用仍需幂等键、Outbox 或专门 sink 保证。

升级处理逻辑前要判断:旧状态能否兼容新 schema,是否需要从原始流重建,以及重建会不会压垮下游。

Choreography 不等于没有流程责任

多个消费者围绕事件协作时,仍需有人拥有端到端结果。至少维护:

  • 流程图和事件所有者;
  • 每一步 SLA 与失败语义;
  • 相关 ID 和业务状态查询;
  • 告警、补偿和人工操作入口。

若一个关键流程需要阅读十个消费者源码才能理解,应该引入流程视图、状态投影,或改为显式编排。

事件链路如何追踪

同步调用的 trace 通过调用栈传播;异步事件跨越时间和进程,需要同时保留:

  • traceId:当前技术调用链;
  • correlationId:同一业务流程;
  • causationId:哪个命令或事件直接导致当前事件;
  • eventId:当前事件的唯一身份。

不要把所有事件强行塞进一个永不结束的 trace。长流程可以用 correlation 连接多个 trace,并用业务时间线展示状态。

关键指标

基础设施在线不代表事件系统健康。至少观察:

  • 生产速率、消费速率和 lag;
  • 最老未处理事件年龄;
  • 端到端处理延迟,而不只是 broker 延迟;
  • 重试、死信、重复和乱序数量;
  • 分区倾斜与热点键;
  • 投影版本、重建进度和业务对账差异。

告警应关联用户影响。例如 lag 增长但仍在业务新鲜度目标内,可以是容量预警;关键结算事件超过截止时间则需要立即响应。

上线前故障演练

  • 消费者在业务提交后、确认前崩溃;
  • 单个分区出现永久坏消息;
  • 事件晚到、重复和跨版本混合;
  • broker 暂时不可用导致 Outbox 堆积;
  • 重建投影时新事件持续进入;
  • 热点聚合键压垮单分区。

能从这些场景恢复,事件驱动才不只是“把同步错误换成异步沉默”。

参考资料

Built with VitePress | Software Systems Atlas