17  批处理、流处理与“及时数据”

本章产出: 从业务决策窗口映射到批次、微批、CDC、事件流和请求时查询的时效策略,并形成可度量的时效 SLO。

“Agent 需要实时数据”是一句常见而昂贵的需求。它常把不同问题混在一起:业务事件是否实时发生、数据是否快速进入平台、派生状态是否及时更新、用户请求时是否需要再次确认。若没有明确决策窗口,团队可能为低频文档建设复杂流系统,也可能用每日批处理支持库存承诺,前者浪费成本,后者制造风险。

及时数据并不等于所有数据零延迟。它意味着在业务需要作出判断的时间之前,数据以足够新鲜、完整和可恢复的方式到达。架构师要把“实时”翻译为时间、损失和服务承诺。

17.1 四个时间与两个延迟

任何记录至少有:

  • 业务事件发生时间;
  • 源系统记录时间;
  • 平台采集与处理时间;
  • Agent 实际读取时间。

管道延迟是事件进入平台所花时间,数据年龄是 Agent 使用时距离业务事件的时间。管道刚刚成功不代表数据新鲜,源系统本身可能停止更新。监控必须同时观察。

对于乱序和补发,还要区分处理时间与事件时间。离线设备今天上传昨天的告警,如果按到达时间排序,会错误理解故障序列。流处理窗口和状态应优先使用事件时间,并定义迟到容忍与更正规则。

17.2 从业务决策窗口倒推

可以用四个问题确定时效:

  1. 业务在事件发生后多快必须判断;
  2. 数据变旧到什么程度会改变判断;
  3. 错误或延迟的损失是什么;
  4. 发现错误后需要多快恢复与更正。
决策窗口 典型场景 可能方案
天至周 历史案例整理、月度分析 批处理
小时 文档发布、客户主数据 增量批/微批
分钟 工单状态、设备配置变化 CDC/事件
秒 告警关联、风险检测 事件流与状态计算
请求时 权限、价格、库存、审批状态 权威 API 查询

请求时查询并不排斥流。系统可以用流维护候选和缓存,在真正承诺前再调用权威服务。这样兼顾响应和一致性。

17.3 批处理的优势与边界

批处理实现清晰、吞吐高、易于重算和对账,适合历史数据、复杂转换和允许较大延迟的任务。它能形成稳定快照,便于评测与审计。很多中小企业的首期 AI 场景,小时或日级数据已经足以验证价值。

批处理的挑战包括窗口内不可见、长任务失败恢复、全量扫描成本和回补。增量批要有水位、删除处理和源端时间可靠性;输出按分区或版本发布,避免消费者读到半成品。重跑应幂等,并区分历史回补与当前结果。

不要因为“批”听起来落后就拒绝它。若维修手册每日正式发布一次,秒级流化不会创造业务价值,反而增加状态和运维复杂度。

17.4 微批与 CDC

微批缩短调度间隔,保留部分批处理简单性,适合分钟级需求。需要注意任务重叠、文件碎片、频繁启动成本和水位。

CDC 从数据库日志获取变化,减少源端扫描并保留操作顺序。它适合交易表变更,但不等于业务事件。一次业务动作可能涉及多表,一条技术更新也可能只是系统维护。下游需要把变化解释为领域状态,或由源应用发布明确事件。

CDC 设计要处理初始快照、切换位点、模式变化、事务、删除、重复和断点。若直接把源表模式作为下游接口,源系统一次字段调整会影响所有消费者。应在产品层隔离。

17.5 事件流与状态

事件流适合需要持续反应和时间窗口的任务。例如设备短时间连续告警的组合比单个告警更有意义。流处理维护近期状态、聚合和规则,并把结果发布为可查询上下文或触发受控流程。

流系统面对三个事实:事件会重复、会乱序、会迟到。每个事件有唯一标识和业务键,消费者幂等;使用事件时间和水位判断窗口;允许迟到数据更正已发布结果;保存可重放日志和状态快照。

“恰好一次”通常是端到端业务语义,而非消息中间件按钮。即使管道技术上一次处理,外部 API 超时重试仍可能重复创建工单。写操作必须使用业务幂等键和结果确认。

状态管理决定恢复能力。窗口、实体当前值和规则状态要有保留、检查点和版本。规则变更后,是否重算历史、从何时生效,应由业务契约说明。

17.6 一致性选择

强一致意味着读取立即看到最新已提交结果,代价通常是延迟、可用性和跨系统复杂度。最终一致允许派生表示稍后收敛,适合索引、分析和多数知识场景。关键是公开一致性模型。

山城精工的手册索引可最终一致,结果显示版本与更新时间;设备最近告警可以分钟级;库存推荐可以使用短缓存,但预留前必须权威确认;工单创建要事务性且幂等。不同对象拥有不同一致性,不必强迫全链路同级。

如果两个来源冲突,要定义权威、有效时间和降级。不要用“最后到达获胜”替业务决策。数据服务可以返回冲突状态,Agent 进行澄清或转人工。

17.7 时效 SLO

完整的时效 SLO 包含:

  • 业务事件到可消费的最大数据年龄;
  • 在统计窗口内达标的比例;
  • 适用实体或记录范围;
  • 测量点与时钟来源;
  • 迟到、更正和计划维护的处理;
  • 超标后的通知、降级和恢复。

例如:“范围内设备配置变化在过去七天百分之九十九于十分钟内进入上下文服务;超过三十分钟时,助手不得给出确定配置结论,并提示现场核实。”这比“准实时同步”可执行得多。

还应定义恢复点与恢复时间。事件日志能重放到什么位置,故障后多快追平积压,恢复期间消费者如何判断状态。管道可用性与数据时效分开衡量。

17.8 成本与复杂度

越低延迟通常意味着常驻计算、消息平台、状态存储、更多监控和更高值班要求。评估成本包括软件、基础设施、源端影响、开发、运维、故障恢复和组织能力。技术许可免费不等于实时系统免费。

可以用“每减少一分钟延迟带来的业务价值”审视。若库存变化会立即改变外部承诺,低延迟值得;若历史案例延迟一小时不会改变决策,就不必流化。分层 SLO 让预算投入关键路径。

流系统还增加认知成本。团队若缺少事件建模与运行经验,应从少量高价值事件开始,保留批量对账和人工降级。架构先进却无人能恢复,是更大的不可靠。

17.9 山城精工策略

数据 时效 技术 超标处理
历史结案工单 日级 批处理 暂不使用新增案例
手册与公告 发布后十分钟 文件事件/微批 显示索引状态,链接正式库
设备主数据 小时级 增量批 要求核实关键身份
配置变更 十分钟级 CDC/业务事件 标记配置可能陈旧
设备告警 分钟级 事件流 转原监控系统
库存与权限 请求时 权威 API 不承诺、转人工
工单写入 即时事务 受控 API 按幂等键查询结果

策略体现“及时而非全部实时”。未来若业务结果证明某条链路是瓶颈,再降低延迟,而不是提前为所有来源建设流平台。

17.10 批流一体的现实含义

批流一体不是必须使用某个统一引擎,而是让两种路径共享模式、实体、规则、质量、血缘和消费契约。实时结果可以快速提供当前状态,批处理定期重算和对账;发现偏差后更正并记录。

同一指标的批流计算要有一致定义和测试。实时近似值与离线最终值应明确标识,不能无提示替换。数据产品端口隐藏内部实现,但公开时间和确定性。

原始事件保留使规则变化后可以重放,批快照支持审计和评测。二者结合比追求某种口号式统一更有价值。

17.11 架构决策步骤

  1. 写出任务和业务决策窗口;
  2. 标记关键数据变化及错误损失;
  3. 定义数据年龄、覆盖、可用和恢复 SLO;
  4. 比较批、微批、CDC、事件、请求时查询;
  5. 设计重复、乱序、迟到、删除、回放与状态;
  6. 定义权威、缓存、一致性与行动前确认;
  7. 估算端到端成本与团队运行能力;
  8. 设定降级和未来复审触发条件。

17.12 本章交付物:时效决策 ADR

ADR 不应只写“采用 Kafka 实现实时”。它要记录为什么业务需要某个窗口,为什么其他方案不足,选择带来什么状态和运维成本,失败时如何降级,以及什么变化会触发复审。

及时数据的核心不是速度竞赛,而是承诺管理。对于每项事实,企业知道它何时发生、何时可见、何时会过期、发生冲突信任谁、失败后如何恢复。Agent 因此不会把昨天的世界当作此刻,也不会因为追求毫秒级响应而让组织承担没有价值的复杂度。

17.13 组织运行同样决定时效

数据能在一分钟内到达,并不意味着业务能在一分钟内反应。如果异常没有所有者、告警没有分级、值班人员无法判断影响,技术上的实时只会更快地产生无人处理的消息。路线图应把生产者、平台、数据产品和业务值班连接起来,规定谁确认源端停止,谁决定降级,谁向用户解释,谁在恢复后复查遗漏动作。

时效目标还要纳入业务日历。夜间设备告警可能需要全天响应,手册发布只在工作日发生,合同规则必须在生效前完成验证。不同时间段可以采用不同资源与告警策略,但不能让 SLO 的统计平均掩盖关键营业窗口的失效。

定期演练是验证承诺的唯一方式。团队可以模拟事件积压、源端时钟错误、重复消息、CDC 位点丢失、权威 API 超时和流状态损坏,观察能否在目标时间发现、降级、恢复和对账。没有经过演练的恢复时间只是一项愿望。

高级架构决策最终要同时说明三条时间线:数据变化的时间线、系统处理的时间线和组织响应的时间线。三者共同落在业务决策窗口内,才是真正的“及时”。