第 7 篇:「Fluss 状态外部化」—— Delta Join 与 Aggregation Merge Engine
第 7 篇:「Fluss 状态外部化」—— Delta Join 与 Aggregation Merge Engine
阅读本文你将了解:传统 Flink 状态管理的痛点、Delta Join 如何将 Join 状态外部化、Aggregation Merge Engine 的聚合状态管理、无状态计算架构如何实现秒级故障恢复和 85% 成本优化。
7.1 Flink 状态管理的痛点
7.1.1 传统有状态 Flink 的问题
传统 Flink-on-Kafka 架构的状态模型: Flink Job ┌──────────────────────────────────────────┐ │ Task Manager 1 │ │ ┌────────────────────────────────────┐ │ │ │ Operator: Windowed Join │ │ │ │ ┌──────────────────────────────┐ │ │ │ │ │ RocksDB State Backend │ │ │ │ │ │ ├── Join State: 500GB │ │ │ │ │ │ ├── Aggregate State: 200GB │ │ │ │ │ │ └── Checkpoint: 700GB/snap │ │ │ │ │ └──────────────────────────────┘ │ │ │ └────────────────────────────────────┘ │ └──────────────────────────────────────────┘三大痛点:
| 痛点 | 影响 | 量化数据 |
|---|---|---|
| RocksDB 膨胀 | 状态随数据量线性增长 | 日增 100GB 状态 |
| 慢恢复 | 故障后需从 Checkpoint 重建 | 5-15 分钟 |
| 扩缩容困难 | 状态需重新分布 | 增加节点反而变慢 |
7.1.2 Checkpoint 恢复的代价
故障恢复流程(传统方案): 1. JobManager 检测到 TaskManager 故障 ← 0s 2. 从最新 Checkpoint 读取元数据 ← 5s 3. 下载 RocksDB 状态快照到新 TaskManager ← 3min(700GB) 4. RocksDB 打开并重建 LSM 树 ← 2min 5. 从 Kafka 回放 Checkpoint → 当前 offset 的数据 ← 30s 6. 恢复正常处理 ← 总计 ~6min7.2 Delta Join:状态外部化原理
7.2.1 架构对比
传统(状态在 Flink): Flink Task = 计算逻辑 + RocksDB 状态 问题:状态绑死在计算节点上 Delta Join(状态在 Fluss): Flink Task = 仅计算逻辑(无状态!) Fluss TabletServer = 状态的唯一所有者7.2.2 Delta Join 工作原理
待 Join 的两张表: orders (左表/驱动表): [order_id, user_id, amount] users (右表/被驱动表): [user_id, name, city] Delta Join 执行流程: 当 orders 来了一条新记录: (order_id=1, user_id=100, amount=99.9) ┌──────────────────────────┐ │ Flink Task (无状态!) │ │ │ │ 1. 接收 order 记录 │ │ user_id = 100 │ │ │ │ 2. 查询 Fluss 中的 Join 状态 │ │ SELECT name, city │ │ FROM users │ │ WHERE user_id = 100 │───────→ ┌──────────────────┐ │ │ │ Fluss TabletServer│ │ 3. 得到结果: │←─────── │ KvStore │ │ name='Alice', │ │ user_id=100 → │ │ city='Beijing' │ │ {name:'Alice', │ │ │ │ city:'Beijing'} │ │ 4. 组合输出: │ └──────────────────┘ │ (1, 100, 'Alice', │ │ 'Beijing', 99.9) │ └──────────────────────────┘ 当 users 更新了 user_id=100 的记录,Fluss 自动维护状态: 新的 Flink Task 可以直接查询到最新状态!7.2.3 Delta Join DDL
-- 左表(驱动表)CREATETABLEorders(order_idBIGINT,user_idBIGINT,amountDECIMAL(10,2),order_timeTIMESTAMP(3),PRIMARYKEY(order_id)NOTENFORCED)WITH('bucket.num'='16','table.merge-engine'='deduplicate');-- 右表(被驱动表)CREATETABLEusers(user_idBIGINT,name STRING,city STRING,PRIMARYKEY(user_id)NOTENFORCED)WITH('bucket.num'='8');-- Delta Join SQLSET'execution.runtime-mode'='streaming';INSERTINTOorder_wideSELECTo.order_id,o.user_id,u.name,u.city,o.amountFROMordersASoLEFTJOINusersFORSYSTEM_TIMEASOFo.order_timeASuONo.user_id=u.user_id;7.2.4 Delta Join 内部实现
// 简化自 org.apache.fluss.flink.delta.DeltaJoinOperatorpublicclassDeltaJoinOperator{/** * 无状态的 Join 处理逻辑 * 每次收到一条左表记录,直接向 Fluss 查询右表 */publicvoidprocessLeftRecord(RowDataleftRecord){// 1. 提取 Join Keybyte[]joinKey=extractJoinKey(leftRecord,joinKeyIndexes);// 2. 向 Fluss 点查右表数据(状态在 Fluss 端!)byte[]rightValue=flussClient.pointLookup(rightTable,joinKey);// 3. 组合左右表数据输出if(rightValue!=null){RowDataoutput=combine(leftRecord,rightValue);collector.collect(output);}// 注意:此 Operator 中没有任何状态变量!// 所有 Join 状态都在 Fluss TabletServer 的 KvStore 中}/** * Flink Checkpoint 时,无需保存任何状态 */@OverridepublicvoidsnapshotState(StateSnapshotContextcontext){// 空方法!没有本地状态需要保存// Fluss 端的状态由 Fluss 自己的 WAL 和 Checkpoint 机制保证}}7.3 Aggregation Merge Engine
7.3.1 聚合状态外部化
-- 聚合表:用户消费统计CREATETABLEuser_stats(user_idBIGINT,total_spentDECIMAL(12,2),order_countINT,last_orderTIMESTAMP(3),PRIMARYKEY(user_id)NOTENFORCED)WITH('bucket.num'='16','table.merge-engine'='aggregation','fields.total_spent.aggregate-function'='sum',-- 累加'fields.order_count.aggregate-function'='sum',-- 计数'fields.last_order.aggregate-function'='last_value'-- 最新值);-- 写入:Fluss 自动聚合INSERTINTOuser_statsVALUES(1,99.90,1,TIMESTAMP'2026-08-08 10:00:00');INSERTINTOuser_statsVALUES(1,49.90,1,TIMESTAMP'2026-08-08 11:00:00');-- 自动合并结果:-- user_id=1, total_spent=149.80, order_count=2, last_order='2026-08-08 11:00:00'7.3.2 合并引擎源码
// 简化自 org.apache.fluss.table.merge.AggregationMergeEnginepublicclassAggregationMergeEngineimplementsMergeEngine{privatefinalMap<String,AggregateFunction>fieldAggregators;/** * 合并新旧两行数据 * * @param oldRow RocksDB 中的旧值 * @param newRow 新写入的值 * @return 合并后的结果 */publicRowDatamerge(RowDataoldRow,RowDatanewRow){RowData.Builderbuilder=RowData.builder(schema);for(inti=0;i<schema.getFieldCount();i++){StringfieldName=schema.getFieldName(i);AggregateFunctionfunc=fieldAggregators.get(fieldName);if(func!=null){// 聚合列:执行聚合函数ObjectoldVal=oldRow!=null?oldRow.getField(i):null;ObjectnewVal=newRow.getField(i);Objectmerged=func.aggregate(oldVal,newVal);builder.setField(i,merged);}elseif(newRow.getField(i)!=null){// 非聚合列:使用新值(或合并策略)builder.setField(i,newRow.getField(i));}elseif(oldRow!=null){builder.setField(i,oldRow.getField(i));}}returnbuilder.build();}}7.3.3 支持的聚合函数
| 聚合函数 | 说明 | 使用场景 |
|---|---|---|
sum | 累加求和 | 金额、次数统计 |
max/min | 取最大/最小值 | 峰值、极端值监控 |
last_value/first_value | 取最新/最早值 | 时间戳、状态记录 |
count | 计数 | 事件计数 |
7.4 无状态计算架构的性能收益
7.4.1 故障恢复:分钟级 → 秒级
Fluss 无状态架构故障恢复: 1. JobManager 检测到 TaskManager 故障 ← 0s 2. 调度新 TaskManager(无需下载状态!) ← 3s 3. 新 Flink Task 直接从 Fluss 读取最新状态 ← 2s 4. 恢复正常处理 ← 总计 ~5s 对比:传统方案 ~360s → Fluss 方案 ~5s,快 70 倍+7.4.2 弹性扩缩容
传统方案扩容 10 → 20 节点: → RocksDB 状态需要重新分布到 20 个节点 → 触发 Savepoint → 全量状态传输 → 耗时:10-30 分钟 Fluss 方案扩容 10 → 20 节点: → Flink 任务无状态,直接启动 20 个并行实例 → 所有实例从 Fluss 读取状态(Fluss 自动负载均衡) → 耗时:10-30 秒7.4.3 成本优化
成本对比(100TB 状态、100 核 Flink 集群): 传统方案: ├── 计算:100 核 × $0.1/核时 = $10/时 ├── 状态存储(本地 SSD):100TB × $0.08/GB/月 = $8000/月 ├── Checkpoint 存储(S3):700GB × 10 保留 = $160/月 └── 总计:~$15,560/月 Fluss 方案: ├── Flink 计算(无存储需求):50 核 × $0.1/核时 = $5/时 ← 减半! ├── Fluss 状态存储:统一管理,成本融入 Fluss 集群 └── 总计:~$4,000/月(含 Fluss 集群),约 75% 成本节省7.5 当前限制与路线图
| 限制 | 当前状态 | 路线图 |
|---|---|---|
| Join 类型 | 仅支持 LEFT JOIN | 计划支持 RIGHT/FULL JOIN |
| 多流 Join | 不支持 N 路 Delta Join | 路线图中 |
| 非等值 Join | 不支持ON a.id > b.id | 暂无计划 |
| Aggregation 类型 | sum/max/min/last/first/count | 计划新增自定义聚合 |
7.6 总结与下一篇预告
| 要点 | 传统方案 | Fluss 方案 |
|---|---|---|
| 状态归属 | Flink RocksDB | Fluss TabletServer |
| 故障恢复 | 5-15 分钟 | 3-5 秒 |
| 扩缩容 | 需状态重分布 | 秒级弹性 |
| Checkpoint 大小 | 数百 GB | 接近 0 |
| 成本 | 计算+存储两套 | 统一基座 |
下一篇我们将学习 Fluss 的 Streaming Lakehouse 架构——如何实现流批统一的实时湖仓,让一个 SQL 查询同时覆盖秒级实时数据和月级历史数据。
本文基于 Apache Fluss 0.9.1 + Apache Flink 1.20.3。项目 GitHub: https://github.com/apache/fluss
