{T}

状态管理:为什么说流计算是有"状态"的?

0. 引言

Flink 官网对自己的第一个描述词就是 Stateful(有状态的)——正是 stateful 让 Flink 从众多流计算框架中脱颖而出。会话窗口、乱序处理、多流关联、聚合寄存器、CEP 状态机、模型参数……这些都是"状态"。可以说,只有理解了"状态",才能理解"流计算"。本文先区分两类状态,再辨析两种窗口,最后讲状态存储方案。

1. 两类状态:流数据状态 vs 流信息状态

状态类型定义典型场景本质
流数据状态为处理流而临时缓存的部分流数据,计算完成即清理事件窗口、时间乱序、多流关联对"时间"维度的管理
流信息状态分析得到的业务信息,被持续查询与更新聚合值、关联图谱节点数、CEP 状态机对"空间"维度的管理

为什么区分如此重要?以"用户过去 7 天交易总金额"为例,用窗口函数实现:

java
userTransactions
    .keyBy(0)
    .timeWindow(Time.days(7), Time.seconds(1))  // 7 天窗口、1 秒滑动
    .sum(1);

这个实现有四个不妥:

  1. 实时性不符:每 1 秒才输出一次,而业务要求"每来一个事件"就计算;
  2. 重复计算:同一数据在窗口滑动中被反复计算约 60 万次(7 天 ÷ 1 秒);
  3. 存储浪费:为算一个 sum 却要保存 7 天全部数据;
  4. 多指标难同步:几十个指标各有不同窗口与步长,同步汇总困难。

根因是混淆了"数据窗口"与"业务窗口"

  • 数据窗口:框架对"流数据"的分块管理(分而治之),保存的是流数据本身 → 流数据状态;
  • 业务窗口:业务意义上的时间范围(如"7 天"),保存的是聚合后的业务信息 → 流信息状态。

两者可以解耦:先用 1 小时的"数据窗口"做批处理,或干脆逐事件处理,都能算出"7 天业务窗口"的总额——不必强行让数据窗口等于业务窗口

2. 流数据状态:窗口、乱序与关联

2.1 事件窗口

"每 30 秒计算一次过去五分钟交易总额""每满 100 个事件计算平均金额"等场景,需要缓冲区临时存储事件,触发条件满足时才处理——缓冲区中的部分流数据就是流数据状态。

2.2 时间乱序与水印(Watermark)

网络传输与并发处理会导致事件乱序到达。解决思路:将事件先保存,等一段时间后乱序事件到达,再按时间排序恢复顺序。但"等多久"是个问题——水印(Watermark)提供了优雅方案:

text
水印 = 事件时间戳 - 等待时间(Δt)
处理单元收到水印时,处理所有时间戳小于水印的事件

给"先到的较大时间戳事件"一个等待"晚到的较小时间戳事件"的机会,同时不会无限等待;代价是处理延迟(等待 Δt)与过期不候(太晚到达的事件被丢弃)。缓存中时间戳大于水印的数据,就是流数据状态。

2.3 流的关联与存储

join/union 等关联操作需要先将各流在窗口内的数据全部缓存,再做类似关系型数据库的 join,最后以流输出。排序、分组等操作同样依赖流数据状态。

存储:流数据状态最理想是全部放内存(接收→处理→删除实时快速),仅在 checkpoint 时写盘;数据量超过内存时,可放文件或外部存储,牺牲部分性能换取容量。

3. 流信息状态:业务信息的持久化

流信息状态主要记录分析出的业务信息,通常依赖数据库存储,原因有三:

  1. 保存时间长、数据量大、需频繁增删查改——数据库天然适合;
  2. 存在"数据变冷"与"过期淘汰"问题——数据库的热数据缓存与 TTL 机制可解;
  3. 数据量大需扩展为集群——数据库便于水平扩展。

3.1 Redis:INCR 原子计数

第 10 课已用 Redis 存储流信息状态:INCR 原子加一维护"过去一周同一设备交易次数"的各窗口寄存器。

3.2 Apache Ignite:分布式内存网格 + CAS

Apache Ignite 是基于内存的数据网格,符合 JCache 标准、支持 SQL。实现同样功能时,因没有 INCR 式原子操作,需用 CAS(Compare And Swap) 保证并发安全:

python
# 更新窗口计数(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 等)。

下一章讲解扩展为集群:如何实现分布式状态存储。