ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

第 7 篇:「Fluss 状态外部化」—— Delta Join 与 Aggregation Merge Engine

2026/8/15 22:17:51 拓冰建站 浏览量
第 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 ← 3min700GB 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_id1, user_id100, amount99.9) ┌──────────────────────────┐ │ Flink Task (无状态!) │ │ │ │ 1. 接收 order 记录 │ │ user_id 100 │ │ │ │ 2. 查询 Fluss 中的 Join 状态 │ │ SELECT name, city │ │ FROM users │ │ WHERE user_id 100 │───────→ ┌──────────────────┐ │ │ │ Fluss TabletServer│ │ 3. 得到结果 │←─────── │ KvStore │ │ nameAlice, │ │ user_id100 → │ │ cityBeijing │ │ {name:Alice, │ │ │ │ city:Beijing} │ │ 4. 组合输出 │ └──────────────────┘ │ (1, 100, Alice, │ │ Beijing, 99.9) │ └──────────────────────────┘ 当 users 更新了 user_id100 的记录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.num16,table.merge-enginededuplicate);-- 右表被驱动表CREATETABLEusers(user_idBIGINT,name STRING,city STRING,PRIMARYKEY(user_id)NOTENFORCED)WITH(bucket.num8);-- Delta Join SQLSETexecution.runtime-modestreaming;INSERTINTOorder_wideSELECTo.order_id,o.user_id,u.name,u.city,o.amountFROMordersASoLEFTJOINusersFORSYSTEM_TIMEASOFo.order_timeASuONo.user_idu.user_id;7.2.4 Delta Join 内部实现// 简化自 org.apache.fluss.flink.delta.DeltaJoinOperatorpublicclassDeltaJoinOperator{/** * 无状态的 Join 处理逻辑 * 每次收到一条左表记录直接向 Fluss 查询右表 */publicvoidprocessLeftRecord(RowDataleftRecord){// 1. 提取 Join Keybyte[]joinKeyextractJoinKey(leftRecord,joinKeyIndexes);// 2. 向 Fluss 点查右表数据状态在 Fluss 端byte[]rightValueflussClient.pointLookup(rightTable,joinKey);// 3. 组合左右表数据输出if(rightValue!null){RowDataoutputcombine(leftRecord,rightValue);collector.collect(output);}// 注意此 Operator 中没有任何状态变量// 所有 Join 状态都在 Fluss TabletServer 的 KvStore 中}/** * Flink Checkpoint 时无需保存任何状态 */OverridepublicvoidsnapshotState(StateSnapshotContextcontext){// 空方法没有本地状态需要保存// Fluss 端的状态由 Fluss 自己的 WAL 和 Checkpoint 机制保证}}7.3 Aggregation Merge Engine7.3.1 聚合状态外部化-- 聚合表用户消费统计CREATETABLEuser_stats(user_idBIGINT,total_spentDECIMAL(12,2),order_countINT,last_orderTIMESTAMP(3),PRIMARYKEY(user_id)NOTENFORCED)WITH(bucket.num16,table.merge-engineaggregation,fields.total_spent.aggregate-functionsum,-- 累加fields.order_count.aggregate-functionsum,-- 计数fields.last_order.aggregate-functionlast_value-- 最新值);-- 写入Fluss 自动聚合INSERTINTOuser_statsVALUES(1,99.90,1,TIMESTAMP2026-08-08 10:00:00);INSERTINTOuser_statsVALUES(1,49.90,1,TIMESTAMP2026-08-08 11:00:00);-- 自动合并结果-- user_id1, total_spent149.80, order_count2, last_order2026-08-08 11:00:007.3.2 合并引擎源码// 简化自 org.apache.fluss.table.merge.AggregationMergeEnginepublicclassAggregationMergeEngineimplementsMergeEngine{privatefinalMapString,AggregateFunctionfieldAggregators;/** * 合并新旧两行数据 * * param oldRow RocksDB 中的旧值 * param newRow 新写入的值 * return 合并后的结果 */publicRowDatamerge(RowDataoldRow,RowDatanewRow){RowData.BuilderbuilderRowData.builder(schema);for(inti0;ischema.getFieldCount();i){StringfieldNameschema.getFieldName(i);AggregateFunctionfuncfieldAggregators.get(fieldName);if(func!null){// 聚合列执行聚合函数ObjectoldValoldRow!null?oldRow.getField(i):null;ObjectnewValnewRow.getField(i);Objectmergedfunc.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/时 ├── 状态存储本地 SSD100TB × $0.08/GB/月 $8000/月 ├── Checkpoint 存储S3700GB × 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 RocksDBFluss 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