引言
流处理与批处理的本质差异不是"实时"两字,而是对时间的理解。批处理假设"数据都齐了再算",流处理却要面对永恒的问题:数据什么时候来?数据可能乱序、迟到、重复——你永远无法确定"这一刻的数据是不是全部"。窗口(Window)、水印(Watermark)、状态(State)就是流处理回答"何时算、算哪些、坏了怎么恢复"的三件核心武器。
流处理进阶的精髓:用事件时间理解业务,用水印容忍乱序,用状态保存进度,用 Exactly-Once 保证正确。
本文系统拆解这三大机制:窗口类型与选择、水印与迟到处理、状态与检查点、精确一次语义与背压——把"写流式作业"提升为"设计可靠流式系统"。
一、事件时间 vs 处理时间
1.1 两种时钟的差异
| 时钟 | 定义 | 优点 | 缺点 |
|---|---|---|---|
| 处理时间 | 算子处理数据那一刻的墙钟 | 简单、低延迟 | 与实际发生时间可能偏离大 |
| 事件时间 | 数据里记录的发生时间 | 反映业务真实 | 需处理乱序与迟到 |
// 事件时间示例:日志里的 ts 是事件时间
{"user_id": 1001, "action": "purchase", "ts": "2026-09-27T10:15:30Z", "amount": 99}
// 处理时间: 流引擎收到这条记录的时刻(可能晚 10 秒、几分钟)
为什么必须用事件时间:业务指标(如"9 点整的用户数")必须以"事件发生时间"聚合,否则网络延迟、重放、移动端离线上报都会把数据算错窗口。
1.2 事件时间带来的难题
- 乱序:数据到达顺序 ≠ 发生顺序。
- 迟到:早发生的数据可能很晚才到(离线补传、长链路)。
- 窗口关闭时机:窗口"该关了",但迟到的数据还来不来?
水印就是回答"什么时候算"的机制。
二、窗口类型与选择
2.1 四种窗口
| 窗口 | 语义 | 例子 | 适用 |
|---|---|---|---|
| 滚动窗口(Tumbling) | 固定长度,不重叠 | 每分钟统计 | 简单固定周期 |
| 滑动窗口(Sliding) | 固定长度,滑动间隔 | 每 10s 算近 5 分钟 | 平滑曲线 |
| 会话窗口(Session) | 按活动间隙切分 | 30 分钟无操作即会话结束 | 用户行为会话 |
| 计数窗口(Count) | 按条数切分 | 每 1000 条聚合 | 无需时间语义 |
# 选型决策
# 固定周期报表/监控 → 滚动窗口
# 趋势平滑/滚动指标 → 滑动窗口
# 用户行为聚类(购物会话)→ 会话窗口
# 与时间无关的批量 → 计数窗口
2.2 窗口的三种实现语义(关键)
# 1) 全局窗口 + 增量聚合: 状态里存所有,窗口到期输出(内存压力大)
# 2) 预聚合窗口: 每个窗口一个聚合器,增量更新(标准做法)
# 3) 两阶段聚合: 细粒度预聚合 → 粗粒度再聚合(去倾斜)
# 工程: 窗口 + 增量状态 + 定时触发输出,避免"窗口到期才全量算"
三、水印(Watermark)与乱序处理
3.1 水印是什么
水印是一个单调递增的时间戳,表示"早于该时间的事件应已全部到达(或到达概率足够低)":
# 水印推进
# t=10:00:00 → 事件 ts <= 10:00:00 的处理可以触发
# 实现: 取每个分区当前处理的最大事件时间 - 允许乱序延迟
# 乱序容忍 = 延迟策略(如 5s/60s)
# 水印推进速率 = min(观察到的最大事件时间 - 容忍度)
水印决定窗口何时输出:窗口的结束时间早于水印时,窗口判定"到齐",触发计算并输出。
3.2 乱序容忍的代价
# 容忍度越大 → 越能接住迟到数据,但输出越晚
# 容忍度 = 业务权衡
# 实时监控: 秒级容忍(接受少量迟到重算)
# 报表/计费: 大容忍(宁可晚也要准)
# 做法: 水印延迟可配,别"一竿子定死"
3.3 迟到数据的三种策略
窗口触发输出后,迟到的数据才到——三种处理:
# 1) 丢弃: 迟到数据直接忽略(延迟敏感,可接受误差)
# 2) 更新: 触发窗口"再次输出"(迟到数据到达时重算窗口,更新下游)
# —— 输出两次(首次 + 更新),下游需处理"更新"语义
# 3) 旁路: 迟到数据路由到旁路流(延迟队列),人工/批处理补救
# 工程选择: 监控/实时 → 丢弃或更新;计费/聚合精确 → 旁路补救
四、状态管理与状态后端
4.1 为什么流处理要有状态
流式算子是有状态的:窗口聚合、去重、Join、计数都要在算子内存里"记住"进展。
# 有状态算子
# 窗口聚合: 当前窗口累计值
# 去重: 已见过的 key 集合
# Join: 一侧的缓冲数据
# 用户画像: 用户最新状态
# 无状态算子(简单): map/filter —— 无状态则无需恢复
4.2 状态后端选择
| 后端 | 存储位置 | 规模 | 特点 |
|---|---|---|---|
| 内存 | JVM 堆 | 小 | 最快,重启丢 |
| RocksDB | 本地磁盘 | 大(TB 级) | 高吞吐持久,序列化开销 |
| 外部 KV(Redis/HBase) | 外部系统 | 超大 | 可共享,延迟/依赖 |
# 选型
# 小状态、低延迟 → 内存
# 大状态(TB 级、长窗口)→ RocksDB
# 跨作业共享/超大状态 → 外部 KV
# 注意: 状态是"作业的私有记忆",尽量算子内管理,别过度外置
4.3 状态规模与 TTL
- 状态会无限增长(key 持续进入)→ 必须设置状态 TTL,过期自动清理。
- 长窗口 + 海量 key:用 RocksDB + 调优(block 大小、缓存)。
- 大状态恢复慢 → 增量检查点/异步快照。
五、Exactly-Once 与检查点
5.1 三种交付语义
| 语义 | 含义 | 代价 |
|---|---|---|
| At-Most-Once | 最多一次(可能丢) | 最简单 |
| At-Least-Once | 至少一次(可能重复) | 需去重 |
| Exactly-Once | 精确一次 | 最贵 |
# 现实: Exactly-Once 是"端到端"还是"引擎内"?
# 引擎内(Flink checkpoint + 状态): 算子级精确一次,易实现
# 端到端(含外部 Sink): 需要外部系统幂等/事务配合
# 幂等写入(key 唯一 upsert)是最务实的端到端 Exactly-Once
5.2 检查点(Checkpoint)
检查点把算子状态快照到可靠存储(周期触发),崩溃后从最近检查点恢复:
# 检查点机制
# 1) 周期性触发(如 60s),异步对齐各算子状态
# 2) 失败 → 作业从最近检查点重启,重放数据
# 3) 与 exactly-once 配合: 状态恢复 + 数据重放 = 无丢失无重复
# 检查点设置
# 间隔太小 → 开销大;太大 → 恢复丢失多
# 大状态 → 增量检查点
# 数据量巨大 → 检查点存储与作业解耦(可迁移)
5.3 从检查点恢复的排障
- 检查点失败是流式作业最常见故障 → 检查点监控是运维核心。
- 检查点持久化失败(存储不可用)→ 隔离检查点存储 + 重试。
- 大状态恢复时间长 → 提前规划恢复窗口。
六、流表 Join 与维表关联
6.1 流-流 Join
两条流的 Join 需要缓存一侧数据(有状态):
# 流-流 Join 的挑战
# 1) 数据永远"不完全到达" → 靠窗口限定 Join 范围
# 2) 延迟导致一侧先到 → 状态缓冲等另一侧
# 3) 用事件时间 + 水印控制 Join 窗口
# 典型: 订单流 × 支付流 Join(限定 5 分钟窗口)
6.2 流-维表 Join(维表关联)
维表(用户信息、产品属性)是静态/准静态,用"维表缓存"关联:
# 维表关联模式
# 1) 广播维表: 小维表广播到所有算子(免查询)
# 2) 异步 IO: 每事件查外部(慢,需缓存)
# 3) 周期加载: 维表定时全量/增量加载进本地缓存
# 工程: 小维表广播;大维表用 异步 IO + 本地 LRU 缓存 + 定期刷新
七、背压(Backpressure)
7.1 背压是什么
下游处理速度跟不上上游时,压力向后传导:
# 背压表现
# 上游产出 > 下游消费 → 队列/内存积压 → OOM 或延迟飙升
# 处理背压
# 1) 引擎自动背压: Flink 反压信号传递,上游放慢/缓冲
# 2) 手动: 调并行度、优化算子(避免单点慢算子)
# 3) 检测: 背压指标监控(算子处理延迟、队列积压)
# 常见诱因
# 慢算子(大 join/热点 key)、外部系统慢(Sink 阻塞)、资源不足
7.2 热点与倾斜的流式处理
# 流式热点 key
# 单个 key 事件密集(头部用户)→ 单算子吃满,其余空闲
# 缓解
# 1) key 加盐分摊(临时 key)→ 聚合后合并
# 2) 两阶段聚合(局部 → 全局)
# 3) 广播 + 局部聚合降低热点
# 注意: 状态型算子(Join/窗口)加盐后"状态归并"变复杂,权衡使用
八、工程实践清单
8.1 流式作业设计检查
# [ ] 用事件时间做窗口聚合,水印容忍度按业务配
# [ ] 迟到数据策略明确(丢弃/更新/旁路)并配了监控
# [ ] 状态设置了 TTL,不会无限增长
# [ ] 检查点开启,间隔与状态规模匹配,检查点失败会告警
# [ ] 端到端 Exactly-Once 落地(幂等 Sink 或事务)
# [ ] 背压有监控,热点 key 有缓解预案
# [ ] 维表关联有缓存策略,异步 IO 不拖慢
8.2 常见陷阱
- 窗口用手处理时间:网络抖动一次,整批数据错窗。
- 水印容忍设 0:乱序数据永远被误判迟到。
- 状态无 TTL:跑一周后状态爆炸,作业撑不住。
- 检查点没监控:连续检查点失败,一次故障恢复成灾难。
- Exactly-Once 只做引擎内:Sink 重复写,下游照样重复。
8.3 与批处理的配合
# 流批一体趋势
# 同一套语义(SQL/Table API)跑批与流
# 流提供实时,批提供回算(补历史/纠错)
# 关键: 流批口径一致(同一窗口语义、同一计算逻辑)
# 流式输出 + 批处理对账 → 发现流式计算漂移
总结
| 机制 | 解决什么 | 关键实践 |
|---|---|---|
| 事件时间 | 业务真实口径 | 别用处理时间 |
| 窗口 | 何时算什么 | 按业务选类型 |
| 水印 | 乱序容忍 | 容忍度业务权衡 |
| 状态 | 保存进展 | 设 TTL、选后端 |
| Exactly-Once | 不丢不重 | 幂等 Sink 落地 |
| 检查点 | 崩溃恢复 | 监控检查点 |
| 背压 | 速度失衡 | 监控 + 倾斜缓解 |
流处理进阶的本质是用机制换取确定性:事件时间换口径正确,水印换乱序容忍,状态换可恢复,检查点换不丢失。落地的核心原则:事件时间为准、水印按业务配、状态必设 TTL、端到端幂等、检查点必监控。 把这些机制用透,流式系统才能从"能跑"变成"跑得对、跑得稳、坏了能救"。
参考与延伸阅读
- Flink 官方文档:Windows、Watermarks 与 Checkpointing
- 《Streaming Systems》(O’Reilly)——水印与窗口的理论基础
- Google Dataflow 论文:Beam 模型(窗口/触发/水印/累积)
- Confluent 文档:Kafka Streams 状态与精确一次
- Apache Flink 流处理 — Flink 实战基础
- Kafka Connect 与 CDC — 流式数据入口
- 实时数据仓库 — 流批一体落地
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。