第13章 流式计算与实时数据¶
实时数据的价值来自业务动作,而不是“越快越先进”。支付拦截、设备异常、库存预警和客服积压如果晚几分钟就失去价值,实时链路值得建设;月度经营复盘更需要完整、稳定和可审计的最终口径,批处理通常更合适。
Agent 进入实时链路后,问题不只剩“当前值是多少”。模型还必须知道事件时间、窗口、Watermark、是否最终、迟到修正和数据血缘。否则一个尚未闭合的窗口,很容易被解释成确定事实。
实时系统真正难的不是把延迟压低,而是在低延迟下仍然做到可解释、可恢复、可回放。
13.1 实时链路的边界:先问延迟是否改变业务动作¶
一家企业的 DataAgent 既要回答“上季度毛利为什么下降”,也可能回答“最近 10 分钟哪些门店支付失败率异常”。前者适合湖仓快照和离线分析,后者需要持续到达的事件和低延迟计算。
图13-1:离线数据与实时数据服务不同业务动作。来源:本书自绘。Alt text:左侧离线数据对应报表、复盘等可容忍延迟的动作,右侧实时数据对应告警、风控、大屏等延迟敏感动作,对比两类数据服务的业务场景。
图13-2:实时链路立项判断路径。来源:本书自绘。Alt text:决策树从"延迟是否影响业务价值"出发,分出需要实时、准实时、离线足够三条路径,提醒先确认实时必要性再投入。
实时链路会引入常驻计算、状态存储、消息积压、乱序、迟到和恢复成本,因此不应成为默认架构。只有当“迟到的数据会导致错误或错过动作”时,实时复杂度才有明确回报。
实时进入 Agent 后还要区分三类结果:
- 实时信号:适合告警、候选解释和下一步调查;
- 已发布实时指标:可以用于问数和运营看板;
- 可审计最终事实:适合正式报告、财务和合规结论。
13.2 流批一体:事件、状态与服务结果是不同对象¶
流式平台通常从事件生产者进入消息总线,再经过状态计算,输出湖仓明细、实时宽表、指标和告警。
图13-3:从原始事件到实时数据产品的端到端链路。来源:本书自绘。Alt text:横向链路依次为事件采集、事件总线、流计算、状态/窗口、实时存储、数据服务,箭头表示原始事件逐步加工为可消费的实时数据产品。
消息系统和计算系统不能混为一谈。Kafka 提供可重放日志、分区和消费进度;Flink/Spark Streaming 负责窗口、Join、去重和状态;Doris、StarRocks、ClickHouse 等负责低延迟查询。
表13-1:事件流、变更日志、实时宽表、实时指标四个概念的定义与区别。来源:本书整理。
| 概念 | 定义 | 与相邻概念的区别 |
|---|---|---|
| 事件流 | 按时间持续追加的业务事实 | 强调“发生过什么” |
| 变更日志 | 数据库行级变化形成的日志流 | 强调表状态变化 |
| 流式计算 | 持续消费并过滤、转换、窗口、关联、聚合 | 强调持续计算 |
| 实时宽表 | 事件与维表、规则、历史状态关联后的可查询状态 | 强调当前上下文 |
| 实时指标 | 按事件时间和窗口持续更新的指标结果 | 强调低延迟服务 |
| Watermark | 对事件时间进度的估计 | 不是“之前事件已全部到齐”的证明 |
| Checkpoint | 状态和输入位置的一致性快照 | 用于故障恢复 |
| Savepoint | 主动触发的可迁移状态快照 | 用于升级和迁移 |
| Exactly-once | 源、状态、Sink 协作下的一致性语义 | 不代表业务世界绝对只发生一次 |
| 背压 | 下游不足导致上游被迫降速 | 是容量和瓶颈信号 |
流批一体不是取消离线,而是让实时结果和离线事实能够围绕同一业务语义对账。 实时支付成功率可以用于告警,日终仍应从完整明细重算并解释两者差异。
13.3 Kafka、Flink 与实时服务层:事实沉淀和低延迟服务要双轨输出¶
实时层位于采集之后、湖仓和服务层之前。它同时承担两条输出路径:一条把原始或清洗事件沉淀到湖仓,用于回放和审计;另一条把窗口指标、告警和实时宽表写到服务层,用于低延迟消费。
图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。来源:本书自绘。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 与失败恢复关系。来源:本书自绘。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:事件契约贯穿实时链路。来源:本书自绘。Alt text:同一份事件契约(字段、类型、时间戳、主键)从生产者、事件总线到流计算、消费端逐段标注,表示契约在全链路一致约束。
事件契约至少应明确事件唯一键、事件时间、分区键、Schema 版本、PII、保留期以及迟到/补发/修正语义。
13.6 从流到表:Agent 应消费受控实时产品,而不是原始 Topic¶
Stream-table Duality 可以把追加事件理解成“发生过什么”,把动态表理解成“当前是什么状态”。订单创建、支付、取消是一组事件;按 order_id 折叠后得到订单 current 表;再按 5 分钟窗口聚合得到实时指标。
图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:背压沿实时链路向上游传播并触发恢复动作。来源:本书自绘。Alt text:下游消费变慢后,背压信号沿链路逐级向上游传递,箭头标出各级触发的限流、扩容、缓冲等恢复动作。
表13-8:背压、重复消费、状态膨胀等流式失败模式的检测与恢复策略。来源:本书整理。
| 失败模式 | 影响 | 检测 | 恢复 |
|---|---|---|---|
| 消息堆积 | 指标延迟、告警滞后 | 消费和端到端延迟 | 扩容、优化慢算子、下游限流 |
| 分区倾斜 | 少数任务拖慢全链路 | 分区吞吐差异 | 重设计 key、热点拆分 |
| 迟到增加 | 窗口反复修正或漏算 | 迟到率、Watermark | 调整策略、修正流 |
| Checkpoint 失败 | 恢复点变旧 | 耗时、失败率 | 缩减状态、优化后端 |
| 状态膨胀 | 恢复变慢、资源升高 | 状态大小、恢复耗时 | TTL、清理无效 key |
| Sink 重复 | 告警重复、指标翻倍 | 幂等冲突、对账 | 事务、幂等键、去重 |
| Schema 不兼容 | 解析失败或数据错误 | Schema 校验 | 灰度、兼容策略、回滚 |
| 日志保留不足 | 无法重算 | 回放失败 | 同步写湖仓、提高保留期 |
流式作业发布不能只替换镜像。涉及状态的升级,应经过回放、影子对比、Savepoint、兼容性检查和灰度。
图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:实时链路延迟诊断路径。来源:本书自绘。Alt text:诊断流程从端到端延迟升高出发,沿生产、总线、消费、状态算子逐段排查,箭头指向各段对应的延迟来源与处理动作。
图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.











