状态管理:为什么说流计算是有"状态"的?
0. 引言
Flink 官网对自己的第一个描述词就是 Stateful(有状态的)——正是 stateful 让 Flink 从众多流计算框架中脱颖而出。会话窗口、乱序处理、多流关联、聚合寄存器、CEP 状态机、模型参数……这些都是"状态"。可以说,只有理解了"状态",才能理解"流计算"。本文先区分两类状态,再辨析两种窗口,最后讲状态存储方案。
1. 两类状态:流数据状态 vs 流信息状态
| 状态类型 | 定义 | 典型场景 | 本质 |
|---|---|---|---|
| 流数据状态 | 为处理流而临时缓存的部分流数据,计算完成即清理 | 事件窗口、时间乱序、多流关联 | 对"时间"维度的管理 |
| 流信息状态 | 分析得到的业务信息,被持续查询与更新 | 聚合值、关联图谱节点数、CEP 状态机 | 对"空间"维度的管理 |
为什么区分如此重要?以"用户过去 7 天交易总金额"为例,用窗口函数实现:
userTransactions
.keyBy(0)
.timeWindow(Time.days(7), Time.seconds(1)) // 7 天窗口、1 秒滑动
.sum(1);这个实现有四个不妥:
- 实时性不符:每 1 秒才输出一次,而业务要求"每来一个事件"就计算;
- 重复计算:同一数据在窗口滑动中被反复计算约 60 万次(7 天 ÷ 1 秒);
- 存储浪费:为算一个 sum 却要保存 7 天全部数据;
- 多指标难同步:几十个指标各有不同窗口与步长,同步汇总困难。
根因是混淆了"数据窗口"与"业务窗口":
- 数据窗口:框架对"流数据"的分块管理(分而治之),保存的是流数据本身 → 流数据状态;
- 业务窗口:业务意义上的时间范围(如"7 天"),保存的是聚合后的业务信息 → 流信息状态。
两者可以解耦:先用 1 小时的"数据窗口"做批处理,或干脆逐事件处理,都能算出"7 天业务窗口"的总额——不必强行让数据窗口等于业务窗口。
2. 流数据状态:窗口、乱序与关联
2.1 事件窗口
"每 30 秒计算一次过去五分钟交易总额""每满 100 个事件计算平均金额"等场景,需要缓冲区临时存储事件,触发条件满足时才处理——缓冲区中的部分流数据就是流数据状态。
2.2 时间乱序与水印(Watermark)
网络传输与并发处理会导致事件乱序到达。解决思路:将事件先保存,等一段时间后乱序事件到达,再按时间排序恢复顺序。但"等多久"是个问题——水印(Watermark)提供了优雅方案:
水印 = 事件时间戳 - 等待时间(Δt)
处理单元收到水印时,处理所有时间戳小于水印的事件给"先到的较大时间戳事件"一个等待"晚到的较小时间戳事件"的机会,同时不会无限等待;代价是处理延迟(等待 Δt)与过期不候(太晚到达的事件被丢弃)。缓存中时间戳大于水印的数据,就是流数据状态。
2.3 流的关联与存储
join/union 等关联操作需要先将各流在窗口内的数据全部缓存,再做类似关系型数据库的 join,最后以流输出。排序、分组等操作同样依赖流数据状态。
存储:流数据状态最理想是全部放内存(接收→处理→删除实时快速),仅在 checkpoint 时写盘;数据量超过内存时,可放文件或外部存储,牺牲部分性能换取容量。
3. 流信息状态:业务信息的持久化
流信息状态主要记录分析出的业务信息,通常依赖数据库存储,原因有三:
- 保存时间长、数据量大、需频繁增删查改——数据库天然适合;
- 存在"数据变冷"与"过期淘汰"问题——数据库的热数据缓存与 TTL 机制可解;
- 数据量大需扩展为集群——数据库便于水平扩展。
3.1 Redis:INCR 原子计数
第 10 课已用 Redis 存储流信息状态:INCR 原子加一维护"过去一周同一设备交易次数"的各窗口寄存器。
3.2 Apache Ignite:分布式内存网格 + CAS
Apache Ignite 是基于内存的数据网格,符合 JCache 标准、支持 SQL。实现同样功能时,因没有 INCR 式原子操作,需用 CAS(Compare And Swap) 保证并发安全:
# 更新窗口计数(CAS 无锁循环)
$id = md5($name)
$newRecord = new CountTable($name, $atTime, 1)
do {
$oldRecord = $cache.get($id)
if ($oldRecord != null):
$newRecord.amount = $oldRecord.amount + 1
else:
$cache.putIfAbsent($id, $newRecord)
$succeed = $cache.replace($id, $oldRecord, $newRecord) # 原子比较替换
} while (!$succeed)
# 查询汇总(Ignite 支持 SQL)
SELECT sum(amount) FROM CountTable
WHERE name = $name AND timestamp > $startTime AND timestamp <= $atTime注意:用 Ignite 存对象时需重写
equals/hashCode,否则序列化反序列化后同一条记录会被视为不同记录,导致replace等操作出错。
4. 小结
- 两类状态:流数据状态(时间维度管理流本身)与流信息状态(空间维度保存业务信息),区分二者使"执行过程"与"信息管理"解耦;
- 数据窗口 ≠ 业务窗口:前者是流的分块管理,后者是业务时间范围,强行耦合会带来重复计算与存储浪费;
- 流数据状态三用途:事件窗口、时间乱序(Watermark)、流的关联;
- 流信息状态存储:高性能、可集群的数据库(Redis、Apache Ignite、RocksDB 等)。
下一章讲解扩展为集群:如何实现分布式状态存储。