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 堆积;
- 重建投影时新事件持续进入;
- 热点聚合键压垮单分区。
能从这些场景恢复,事件驱动才不只是“把同步错误换成异步沉默”。
参考资料
- Apache Kafka, Kafka Streams Architecture
- Apache Flink, Timely Stream Processing
- OpenTelemetry, Messaging semantic conventions