跳转至

第13章 流式计算与实时数据


实时数据的价值来自业务动作,而不是“越快越先进”。支付拦截、设备异常、库存预警和客服积压如果晚几分钟就失去价值,实时链路值得建设;月度经营复盘更需要完整、稳定和可审计的最终口径,批处理通常更合适。

Agent 进入实时链路后,问题不只剩“当前值是多少”。模型还必须知道事件时间、窗口、Watermark、是否最终、迟到修正和数据血缘。否则一个尚未闭合的窗口,很容易被解释成确定事实。

实时系统真正难的不是把延迟压低,而是在低延迟下仍然做到可解释、可恢复、可回放。

13.1 实时链路的边界:先问延迟是否改变业务动作

一家企业的 DataAgent 既要回答“上季度毛利为什么下降”,也可能回答“最近 10 分钟哪些门店支付失败率异常”。前者适合湖仓快照和离线分析,后者需要持续到达的事件和低延迟计算。

图13-1:离线数据与实时数据服务不同业务动作

图13-1:离线数据与实时数据服务不同业务动作。来源:本书自绘。Alt text:左侧离线数据对应报表、复盘等可容忍延迟的动作,右侧实时数据对应告警、风控、大屏等延迟敏感动作,对比两类数据服务的业务场景。

图13-2:实时链路立项判断路径

图13-2:实时链路立项判断路径。来源:本书自绘。Alt text:决策树从"延迟是否影响业务价值"出发,分出需要实时、准实时、离线足够三条路径,提醒先确认实时必要性再投入。

实时链路会引入常驻计算、状态存储、消息积压、乱序、迟到和恢复成本,因此不应成为默认架构。只有当“迟到的数据会导致错误或错过动作”时,实时复杂度才有明确回报。

实时进入 Agent 后还要区分三类结果:

  • 实时信号:适合告警、候选解释和下一步调查;
  • 已发布实时指标:可以用于问数和运营看板;
  • 可审计最终事实:适合正式报告、财务和合规结论。

13.2 流批一体:事件、状态与服务结果是不同对象

流式平台通常从事件生产者进入消息总线,再经过状态计算,输出湖仓明细、实时宽表、指标和告警。

图13-3:从原始事件到实时数据产品的端到端链路

图13-3:从原始事件到实时数据产品的端到端链路。来源:本书自绘。Alt text:横向链路依次为事件采集、事件总线、流计算、状态/窗口、实时存储、数据服务,箭头表示原始事件逐步加工为可消费的实时数据产品。

消息系统和计算系统不能混为一谈。Kafka 提供可重放日志、分区和消费进度;Flink/Spark Streaming 负责窗口、Join、去重和状态;Doris、StarRocks、ClickHouse 等负责低延迟查询。

表13-1:事件流、变更日志、实时宽表、实时指标四个概念的定义与区别。来源:本书整理。

概念 定义 与相邻概念的区别
事件流 按时间持续追加的业务事实 强调“发生过什么”
变更日志 数据库行级变化形成的日志流 强调表状态变化
流式计算 持续消费并过滤、转换、窗口、关联、聚合 强调持续计算
实时宽表 事件与维表、规则、历史状态关联后的可查询状态 强调当前上下文
实时指标 按事件时间和窗口持续更新的指标结果 强调低延迟服务
Watermark 对事件时间进度的估计 不是“之前事件已全部到齐”的证明
Checkpoint 状态和输入位置的一致性快照 用于故障恢复
Savepoint 主动触发的可迁移状态快照 用于升级和迁移
Exactly-once 源、状态、Sink 协作下的一致性语义 不代表业务世界绝对只发生一次
背压 下游不足导致上游被迫降速 是容量和瓶颈信号

流批一体不是取消离线,而是让实时结果和离线事实能够围绕同一业务语义对账。 实时支付成功率可以用于告警,日终仍应从完整明细重算并解释两者差异。

实时层位于采集之后、湖仓和服务层之前。它同时承担两条输出路径:一条把原始或清洗事件沉淀到湖仓,用于回放和审计;另一条把窗口指标、告警和实时宽表写到服务层,用于低延迟消费。

图13-4:流式计算层在企业 Agent 平台中的位置

图13-4:流式计算层在企业 Agent 平台中的位置。来源:本书自绘。Alt text:分层图中流式计算层位于数据采集之上、实时存储与特征服务之下,向 Agent 平台提供实时特征与告警,标出其"实时供给"职责。

表13-2:事件总线、流计算引擎等流式组件的职责、输入输出与失败模式。来源:本书整理。

组件 职责 输入 输出 失败模式
事件生产者 把业务动作、日志或 CDC 写入总线 业务事务、CDC、设备消息 标准事件 重复、乱序、字段漂移
事件总线 保存事件、分区、有序追加和进度 标准事件 分区日志 倾斜、堆积、保留不足
流式计算 过滤、转换、窗口、关联和聚合 流、维表、规则 指标、告警、宽表 状态膨胀、背压、Checkpoint 失败
状态存储 保存窗口、去重和中间状态 key、窗口、事件 可恢复状态 状态过大、恢复慢
服务层 提供查询、告警和在线上下文 实时结果 查询、告警、特征 重复写入、结果不一致
治理观测 Schema、血缘、延迟、权限、审计 作业元数据、运行指标 告警、审计、影响分析 Owner 不清、证据断裂

表13-3:批处理与流处理在延迟、成本、对账上的取舍。来源:本书整理。

方案 优势 代价 适用场景 本书建议
批处理 简单、成本低、易对账 延迟高 日报、财务、离线特征 保留为最终事实底座
微批 复杂度适中、吞吐高 秒到分钟延迟 近实时报表、湖仓增量写入 多数企业准实时场景优先
连续流 延迟低、状态能力强 运维和恢复复杂 风控拦截、设备告警 只用于明确低延迟动作

表13-4:Flink、Spark Streaming 等流计算引擎的优势、代价与适用场景。来源:本书整理。

方案 优势 代价 适用场景 本书建议
Flink 事件时间和状态能力强 调优和运维门槛高 低延迟、复杂窗口和实时关联 企业实时主力候选
Spark Structured Streaming 与 Spark/湖仓生态一致 超低延迟和复杂状态不够自然 增量 ETL、近实时 适合离线团队平滑演进
Kafka Streams 嵌入应用、部署轻 集中治理能力较弱 单服务局部流处理 不作为统一实时平台默认方案

关键实时结果通常应同时写入湖仓与服务层。

表13-6:直接写 OLAP 与保留事件流两种实时服务方式的取舍。来源:本书整理。

方案 优势 代价 适用场景 本书建议
直接写 OLAP 查询快 审计和回放不足 高频实时看板 不作为唯一事实
写湖仓明细 可追溯、可回放 服务延迟较高 原始事件、清洗明细 关键事件必须保留
同时写湖仓和服务层 兼顾追溯和低延迟 双写治理复杂 关键实时指标和告警 默认推荐,配套幂等和对账

13.4 时间语义:事件时间、Watermark 与“是否最终”

实时指标最容易出现的误解,是把处理时间当业务时间。支付 10:00:03 发生,但 10:00:11 才被处理;如果按处理时间计入窗口,就会让 10:00–10:05 的成功率被低估。

图13-5:事件时间、摄入时间、处理时间和 Watermark

图13-5:事件时间、摄入时间、处理时间和 Watermark。来源:本书自绘。Alt text:时间轴上标出同一事件的三个时间戳(发生、进入系统、被处理)及其间隔,Watermark 线表示允许的乱序边界,超过即视为迟到。

Watermark 是“多数事件已经推进到这里”的估计,不是完整性证明。容忍 30 秒迟到可以更快告警,但更容易漏掉网络差的门店;容忍 10 分钟更完整,却会拖慢动作。

表13-5:速度优先与正确性优先两种时间语义策略的取舍。来源:本书整理。

方案 优势 代价 适用场景 本书建议
速度优先 告警快 结果会修正或误报 异常预警、运营监控 返回 Watermark 与是否最终
完整性优先 稳定、少修正 延迟更高 财务、监管、正式复盘 用离线或长 Watermark 确认
双轨输出 快速结果和最终结果并存 治理更复杂 高价值指标、风控、供应链 关键业务优先

实时指标服务不应只返回一个数值:

{
  "metric": "payment_success_rate",
  "window_start": "2026-06-11T10:00:00+08:00",
  "window_end": "2026-06-11T10:05:00+08:00",
  "value": 0.982,
  "watermark": "2026-06-11T10:04:30+08:00",
  "is_final": false,
  "late_event_policy": "update_until_10_minutes",
  "source_lag_seconds": 35
}

is_final=false 是业务语义,不是技术细节。Agent 必须据此调整结论强度。

13.5 状态、一致性与恢复:Exactly-once 不能替代业务幂等

窗口聚合、去重、流式 Join 和规则匹配都依赖状态。Checkpoint 定期保存状态与输入位置,失败后从一致性点恢复;Savepoint 更适合版本升级、迁移和并行度调整。

图13-6:状态、Checkpoint、Savepoint 与失败恢复关系

图13-6:状态、Checkpoint、Savepoint 与失败恢复关系。来源:本书自绘。Alt text:流作业的运行状态周期性写出 Checkpoint,手动触发 Savepoint,故障时从最近 Checkpoint 恢复、升级时从 Savepoint 恢复,箭头标出两类恢复路径。

Exactly-once 通常要求输入可重放、计算状态可恢复、Sink 支持事务或幂等。它只保证特定系统边界内的一致性,不保证工单、支付或告警在业务层绝对只发生一次。

例如告警系统仍应使用稳定业务键:

alert_id = rule_id + store_id + window_start

事件契约也必须贯穿生产、消费和审计:

{
  "event_id": "evt_20260611_000001",
  "event_type": "payment.succeeded",
  "schema_version": "v3",
  "event_time": "2026-06-11T10:00:03+08:00",
  "source": "pos-payment",
  "partition_key": "store_1024",
  "trace_id": "trace_8f4a",
  "producer_time": "2026-06-11T10:00:04+08:00",
  "payload": {
    "order_id": "ord_10086",
    "store_id": "store_1024",
    "amount": 128.50,
    "status": "succeeded"
  }
}

图13-7:事件契约贯穿实时链路

图13-7:事件契约贯穿实时链路。来源:本书自绘。Alt text:同一份事件契约(字段、类型、时间戳、主键)从生产者、事件总线到流计算、消费端逐段标注,表示契约在全链路一致约束。

事件契约至少应明确事件唯一键、事件时间、分区键、Schema 版本、PII、保留期以及迟到/补发/修正语义。

13.6 从流到表:Agent 应消费受控实时产品,而不是原始 Topic

Stream-table Duality 可以把追加事件理解成“发生过什么”,把动态表理解成“当前是什么状态”。订单创建、支付、取消是一组事件;按 order_id 折叠后得到订单 current 表;再按 5 分钟窗口聚合得到实时指标。

图13-8:事件流、变更日志、动态表和物化视图的关系

图13-8:事件流、变更日志、动态表和物化视图的关系。来源:本书自绘。Alt text:事件流经聚合变为变更日志,变更日志物化为动态表,动态表再生成物化视图,箭头双向标注流与表可相互转换(流表二象性)。

DataAgent 不应直接读取任意 Kafka Topic。原始事件可能含敏感字段、坏数据、重复和未稳定口径。更合理的是通过实时 OLAP、特征服务或上下文服务返回受控结果。

表13-7:Agent 直读事件流与经特征服务两种实时供给方式的取舍。来源:本书整理。

方案 优势 代价 适用场景 本书建议
Agent 直接读事件流 延迟低、灵活 权限、脱敏、迟到难治理 调试、内部实验 不作为生产默认路径
Agent 读实时 OLAP 查询表达强 要管理 SQL 安全和窗口完整性 实时指标问答 响应带新鲜度和血缘
Agent 读上下文服务 契约清晰、易限流脱敏 建设成本高 风控、客服、调度等动作型 Agent 关键动作优先

一个上下文响应应至少提供窗口、Watermark、新鲜度、finality、信号和允许动作:

context:
  subject: store_1024
  window: 5m
  generated_at: "2026-06-11T10:05:12+08:00"
  watermark: "2026-06-11T10:04:30+08:00"
  freshness_seconds: 42
  finality: provisional
signals:
  payment_success_rate:
    value: 0.982
    baseline: 0.995
    severity: warning
governance:
  pii_status: masked
  allowed_actions: [explain, create_ticket, request_human_review]

13.7 背压、状态膨胀与发布:恢复能力决定生产质量

背压意味着下游处理能力不足,压力会逐步反向传播到计算和消息总线。排查不能只看“任务是否运行”,还要同时看输入/输出速率、消费延迟、Watermark、Checkpoint、状态大小和 Sink 耗时。

图13-9:背压沿实时链路向上游传播并触发恢复动作

图13-9:背压沿实时链路向上游传播并触发恢复动作。来源:本书自绘。Alt text:下游消费变慢后,背压信号沿链路逐级向上游传递,箭头标出各级触发的限流、扩容、缓冲等恢复动作。

表13-8:背压、重复消费、状态膨胀等流式失败模式的检测与恢复策略。来源:本书整理。

失败模式 影响 检测 恢复
消息堆积 指标延迟、告警滞后 消费和端到端延迟 扩容、优化慢算子、下游限流
分区倾斜 少数任务拖慢全链路 分区吞吐差异 重设计 key、热点拆分
迟到增加 窗口反复修正或漏算 迟到率、Watermark 调整策略、修正流
Checkpoint 失败 恢复点变旧 耗时、失败率 缩减状态、优化后端
状态膨胀 恢复变慢、资源升高 状态大小、恢复耗时 TTL、清理无效 key
Sink 重复 告警重复、指标翻倍 幂等冲突、对账 事务、幂等键、去重
Schema 不兼容 解析失败或数据错误 Schema 校验 灰度、兼容策略、回滚
日志保留不足 无法重算 回放失败 同步写湖仓、提高保留期

流式作业发布不能只替换镜像。涉及状态的升级,应经过回放、影子对比、Savepoint、兼容性检查和灰度。

图13-10:实时作业发布与治理流程

图13-10:实时作业发布与治理流程。来源:本书自绘。Alt text:发布流程含版本打包、Savepoint 触发、状态兼容校验、灰度、回滚等步骤,箭头表示流作业升级须经状态兼容检查而非直接替换。

job:
  name: payment-success-rate-stream
  owner: payment-data-team
  version: 2026.06.11
source:
  type: kafka
  topic: payment-events-v3
  event_time_field: event_time
watermark:
  max_out_of_orderness: 2m
  allowed_lateness: 10m
state:
  checkpoint_interval: 30s
  state_retention: 2h
  savepoint_required_for_upgrade: true
sink:
  lakehouse_table: dwd.payment_events_rt
  olap_table: ads.payment_success_rate_5m
  idempotent_key: window_start,window_end,store_id

窗口计算逻辑可以表达为:

CREATE TABLE payment_events (
  event_id STRING,
  store_id STRING,
  status STRING,
  event_time TIMESTAMP(3),
  WATERMARK FOR event_time AS event_time - INTERVAL '2' MINUTE
);

INSERT INTO payment_success_rate_5m
SELECT
  TUMBLE_START(event_time, INTERVAL '5' MINUTE),
  TUMBLE_END(event_time, INTERVAL '5' MINUTE),
  store_id,
  SUM(CASE WHEN status = 'succeeded' THEN 1 ELSE 0 END) * 1.0 / COUNT(*),
  COUNT(*)
FROM payment_events
GROUP BY TUMBLE(event_time, INTERVAL '5' MINUTE), store_id;

表13-9:实时链路在发布、监控、扩缩容、治理各领域的必备能力。来源:本书整理。

领域 必备能力 说明
发布 版本、配置快照、Savepoint、回滚 状态兼容必须显式处理
监控 输入/输出、消费延迟、Watermark、Checkpoint 共同解释端到端延迟
扩缩容 按吞吐、状态和 Sink 调并行度 先确认瓶颈位置
治理 事件契约、Owner、SLA、血缘、权限 结果进入 Agent 前可解释
恢复 Checkpoint、Savepoint、回放、幂等、对账 恢复需要提前演练
成本 常驻资源、状态、消息保留、重算 实时链路空闲时也有成本

图13-11:实时链路延迟诊断路径

图13-11:实时链路延迟诊断路径。来源:本书自绘。Alt text:诊断流程从端到端延迟升高出发,沿生产、总线、消费、状态算子逐段排查,箭头指向各段对应的延迟来源与处理动作。

图13-12:实时系统的技术取舍需要同时看延迟、正确性、成本和治理

图13-12:实时系统的技术取舍需要同时看延迟、正确性、成本和治理。来源:本书自绘。Alt text:以延迟、正确性、成本、治理为四轴的雷达图,标注"无法四者同时最优",说明实时方案选择是四维平衡而非单点优化。

低延迟、正确性、成本和治理无法同时做到极致,生产设计必须明确优先级。

本章小结

实时链路服务低延迟动作,但不能替代离线湖仓。关键业务应同时保留可查询的实时结果和可回放的事实明细。

事件时间、Watermark、状态、Checkpoint 和业务幂等共同决定实时结果是否可信。Agent 使用实时数据时,必须获得窗口、新鲜度、是否最终、血缘和质量状态;否则模型会把临时信号包装成正式事实。

实时系统的成熟度不取决于“多少毫秒”,而取决于异常、迟到和恢复发生时,系统是否仍能解释自己正在使用什么数据。

参考文献

Apache Flink. (n.d.). Documentation.

Apache Kafka. (n.d.). Documentation.

Apache Spark. (n.d.). Structured Streaming Programming Guide.

Akidau, T. et al. (2015). The Dataflow Model. VLDB.