当前位置: 首页 > news >正文

第 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. 恢复正常处理 ← 总计 ~6min

7.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 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

http://www.jsqmd.com/news/1400815/

相关文章:

  • NVIDIA Nemotron 3.5 Lightning:专为高效推理优化的开源大语言模型部署实战
  • Git代码合并与冲突解决实战技巧
  • Dify 中级实验(07):子工作流——如何把公共逻辑做成可复用积木?
  • 基于Human Behavior分析的AI原生应用设计:从行为流协同到工程实践
  • Python开发轻量级员工管理系统的实践与优化
  • 最长公共前缀算法精解:横向与纵向扫描的实战剖析
  • AtCoder Beginner Contest 471 ABCDE
  • 中美Robotaxi“四国杀”:2026,谁在领跑万亿出行终局?
  • Gitee高人气开源项目深度解析:从筛选到源码学习的全链路指南
  • AOSP-- 第 2 章:源码与构建系统
  • IDEA中git stash可视化操作指南:提升多任务开发效率
  • 学习25
  • 从2020亚太赛A题看数学建模:数据驱动与机理分析的融合实践
  • OpenSpec自定义步骤实战:从数据清洗到复杂工作流构建
  • 《ETA1476FT2G 使用笔记|SOT23‑6 2A 同步降压 DC‑DC 芯片》#ETA1476FT2G #钰泰 ETA #DC‑DC 降压 #电源设计 #国产电源芯片
  • SQL GROUP BY 详解:从分类汇总到多维度聚合实战
  • 物流行业GEO优化公司哪家好?2026物流企业AI搜索获客选型全指南 - 科技前沿信息
  • 免费离线OCR软件Umi-OCR完整上手指南:截图、批量、PDF识别一次讲透
  • 无人机维修入门必修课|看懂无人机结构,认识每一个零件 - 湖南阳光技术
  • C++时间处理实战:从time_t到localtime,掌握时间戳与日期转换
  • OpenClaw AI智能体框架:从部署到十大核心Skills实战指南
  • 泰迪狗场选购测评|正规犬舍筛选干货指南 - Full19
  • ArcGIS Pro要素裁剪进阶:从基础操作到自动化工程实践
  • 电动汽车集群优化:Matlab与Yalmip建模实践
  • 小白吃透DeepSeek Harness:新一代AI Agent框架保姆级实战教程
  • 系统重构与现代化改造:从概念到实践的避坑指南
  • 《文明6》DLC解锁补丁技术解析:原理、安装与问题排查
  • 免费Windows内存清理工具Mem Reduct完整上手指南
  • Go缓存策略实战从本地缓存到Redis多级缓存
  • XML Schema 实例详解:从入门到实战