
第 8 篇「Streaming Lakehouse 实战」—— 构建实时湖仓统一数据层阅读本文你将了解Lakehouse 架构的实时性困境、Fluss Streaming Lakehouse 架构原理、Tiering Service 的工作机制、Union Read 如何实现流批统一查询、以及与 Iceberg 和 Paimon 的深度集成。8.1 Lakehouse 的实时性困境8.1.1 传统 Lakehouse 的矛盾Lakehouse 存储格式Iceberg/Paimon/Hudi的写入困境 方案 A高频写入追求低延迟 → 每分钟 Commit 一次 → 产生大量小 Parquet 文件 → 读取效率极差打开 1440 个文件/天 方案 B低频写入追求读取效率 → 每小时 Commit 一次 → 文件大且规整 → 数据延迟 60 分钟无法满足实时需求 结论传统 Lakehouse 最佳数据新鲜度 ≈ 5-10 分钟 无法满足亚秒级实时分析需求8.1.2 Fluss 的解决方案Fluss Streaming Lakehouse Hot Tier (Fluss) Cold Tier (Lakehouse) ┌─────────────────────────────────────────────────────────┐ │ 用户查询SELECT * FROM orders │ └─────────────────────┬───────────────────────────────────┘ │ ┌───────────┴───────────────┐ │ Union Read │ │ (自动合并两个数据源) │ └───────────┬───────────────┘ │ ┌────────────┴──────────────┐ │ │ ┌──────▼──────┐ ┌─────────▼────────┐ │ Hot Tier │ │ Cold Tier │ │ (Fluss) │ │ (Iceberg/Paimon)│ │ │ │ │ │ Arrow 格式 │ │ Parquet 格式 │ │ 秒级新鲜度 │ │ 分钟级新鲜度 │ │ 保留 1-3 天 │ │ 保留 1-12 月 │ │ 低延迟查询 │ │ 高效率分析 │ └─────────────┘ └───────────────────┘8.2 Tiering Service 详解8.2.1 工作流程Tiering Service 定期执行 Compaction 任务 ┌─────────────────────────────────────────────────────┐ │ Tiering Service │ │ │ │ 1. 发现需要 Compaction 的 Bucket │ │ └── 检查 Tiering offset 上次 Compaction offset │ │ │ │ 2. 读取 Hot Tier 中的 Arrow 数据 │ │ └── LogScanner 从指定 offset 开始读取 │ │ │ │ 3. 排序、去重、合并 │ │ └── 按主键排序 Deduplicate │ │ │ │ 4. 写入 Cold Tier (Parquet) │ │ └── 批量写入 Iceberg/Paimon 表 │ │ │ │ 5. 更新 Tiering offset 元数据 │ │ └── 记录已 Compaction 的位置 │ │ │ │ 6. 清理 Hot Tier 中已 Compaction 的 Segment │ │ └── 释放本地 SSD 空间 │ └─────────────────────────────────────────────────────┘8.2.2 源码结构org.apache.fluss.server.tiering ├── TieringService.java # Tiering 服务入口 ├── TieringTask.java # 单个 Compaction 任务 ├── lake/ │ ├── LakeTableTierer.java # 湖表 Tiering 逻辑 │ ├── iceberg/ │ │ └── IcebergCommitter.java # Iceberg Commit │ ├── paimon/ │ │ └── PaimonCommitter.java # Paimon Commit │ └── lance/ │ └── LanceCommitter.java # Lance Commit └── compactor/ └── ArrowParquetCompactor.java # Arrow → Parquet 转换8.2.3 Compaction 配置-- 配置 Tiering 行为ALTERTABLEordersSET(table.tiering.enabledtrue,-- 启用分层存储table.tiering.commit-interval5min,-- 每 5 分钟 Compaction 一次table.tiering.compaction.max-memory2gb,-- Compaction 内存限制table.tiering.local-retention3d-- Hot Tier 保留 3 天);8.3 Union Read流批统一查询8.3.1 Union Read 原理查询: SELECT COUNT(*), SUM(amount) FROM orders Union Read 执行计划 ┌────────────────────────────────────────────────────┐ │ Union Read Coordinator │ │ │ │ ┌──────────────────┐ ┌──────────────────────┐ │ │ │ Hot Tier Scan │ │ Cold Tier Scan │ │ │ │ (Fluss, Arrow) │ │ (Iceberg, Parquet) │ │ │ │ dt2026-08-06 ~ │ │ dt2026-01-01 ~ │ │ │ │ 2026-08-08 │ │ 2026-08-05 │ │ │ │ rows: 200M │ │ rows: 5B │ │ │ └────────┬─────────┘ └──────────┬───────────┘ │ │ │ │ │ │ └─────────┬───────────────┘ │ │ │ │ │ ┌─────────▼──────────┐ │ │ │ Merge Aggregate │ │ │ │ (合并两个扫描结果) │ │ │ └────────────────────┘ │ └────────────────────────────────────────────────────┘8.3.2 Union Read 的 Flink 执行-- 用户只需写一个 SQLUnion Read 透明执行SELECTDATE_FORMAT(order_time,yyyy-MM-dd)ASdt,COUNT(*)ASorder_count,SUM(amount)ASrevenueFROMordersWHEREorder_timeTIMESTAMP2026-01-01 00:00:00GROUPBYDATE_FORMAT(order_time,yyyy-MM-dd)ORDERBYdt;Flink 优化器会自动拆分为两个子查询Flink 执行计划 UnionAll ├── FlussTableScan (Hot Tier, 近 3 天数据) │ └── filter: order_time 2026-08-05 └── IcebergTableScan (Cold Tier, 历史数据) └── filter: order_time BETWEEN 2026-01-01 AND 2026-08-058.4 与 Iceberg 集成8.4.1 配置 Iceberg 作为冷存储-- 创建 Fluss 表配置 Iceberg 冷存储CREATETABLEorders_lakehouse(order_idBIGINT,user_idBIGINT,amountDECIMAL(10,2),statusSTRING,order_timeTIMESTAMP(3),dt STRING,PRIMARYKEY(order_id,dt)NOTENFORCED)PARTITIONEDBY(dt)WITH(bucket.num16,-- Iceberg 冷存储配置table.tiering.enabledtrue,table.tiering.lake.formaticeberg,table.tiering.lake.cataloghadoop,table.tiering.lake.warehouses3://my-bucket/warehouse/,-- S3 配置s3.endpointhttps://s3.amazonaws.com,s3.access-keyYOUR_ACCESS_KEY,s3.secret-keyYOUR_SECRET_KEY);8.4.2 外部引擎读取 Iceberg 数据-- Spark 直接读取 Fluss Tiering 到 Iceberg 的数据-- 注意Spark 连接的是 Iceberg Catalog不是 FlussCREATECATALOG iceberg_catalogWITH(typeiceberg,catalog-typehadoop,warehouses3://my-bucket/warehouse/);USECATALOG iceberg_catalog;-- 查询冷数据SELECTdt,COUNT(*),SUM(amount)FROMorders_lakehouseWHEREdtBETWEEN2026-01-01AND2026-07-31GROUPBYdt;8.5 与 Paimon 集成8.5.1 Paimon 配置CREATETABLEorders_paimon(order_idBIGINT,user_idBIGINT,amountDECIMAL(10,2),order_timeTIMESTAMP(3),dt STRING,PRIMARYKEY(order_id,dt)NOTENFORCED)PARTITIONEDBY(dt)WITH(bucket.num16,table.tiering.enabledtrue,table.tiering.lake.formatpaimon,table.tiering.lake.warehouses3://my-bucket/paimon-warehouse/);8.5.2 Iceberg vs Paimon 选型维度Apache IcebergApache Paimon生态成熟度更成熟Spark/Trino/Flink 原生支持快速发展中Flink 原生优化更新性能Merge-on-Read适合批量更新LSM 引擎适合高频更新查询性能优秀Snapshot Isolation优秀主键索引与 Fluss 集成深度Tiering Union ReadTiering Union Read推荐场景注重生态兼容性注重高频更新和 Flink 优化8.6 实战案例实时用户行为分析 Lakehouse8.6.1 表设计USECATALOG fluss_catalog;CREATEDATABASEanalytics;USEanalytics;-- 行为事件表Log Table → 无需主键CREATETABLEuser_events(event_idBIGINT,user_idBIGINT,event_type STRING,-- click, view, purchasepage_url STRING,product_idBIGINT,event_timeTIMESTAMP(3),event_date STRING)PARTITIONEDBY(event_date)WITH(bucket.num32,table.tiering.enabledtrue,table.tiering.lake.formaticeberg,table.tiering.local-retention2d);-- 用户画像表PK TableCREATETABLEuser_profile(user_idBIGINT,name STRING,city STRING,segment STRING,-- new, active, viptotal_visitsINT,total_spentDECIMAL(12,2),PRIMARYKEY(user_id)NOTENFORCED)WITH(bucket.num8,table.tiering.enabledtrue,table.tiering.lake.formatpaimon,table.merge-engineaggregation,fields.total_visits.aggregate-functionsum,fields.total_spent.aggregate-functionsum);-- 实时指标大屏表CREATETABLErealtime_dashboard(event_date STRING,event_type STRING,pvBIGINT,uvBIGINT,revenueDECIMAL(14,2),PRIMARYKEY(event_date,event_type)NOTENFORCED)WITH(bucket.num4,table.merge-enginededuplicate);8.6.2 流处理 PipelineSETexecution.runtime-modestreaming;-- Pipeline 1: 实时指标聚合INSERTINTOrealtime_dashboardSELECTevent_date,event_type,COUNT(*)ASpv,COUNT(DISTINCTuser_id)ASuv,SUM(CASEWHENevent_typepurchaseTHEN1ELSE0END)ASrevenueFROMuser_eventsGROUPBYevent_date,event_type;-- Pipeline 2: 用户画像更新INSERTINTOuser_profile(user_id,total_visits,total_spent)SELECTuser_id,1AStotal_visits,CASEWHENevent_typepurchaseTHEN1ELSE0ENDAStotal_spentFROMuser_events;8.6.3 历史分析查询SETexecution.runtime-modebatch;-- 查询过去 90 天的趋势自动 Union Read 冷热数据SELECTevent_date,SUM(pv)ASdaily_pv,SUM(uv)ASdaily_uv,SUM(revenue)ASdaily_revenueFROMrealtime_dashboardWHEREevent_dateDATE_FORMAT(TIMESTAMPADD(DAY,-90,NOW()),yyyy-MM-dd)GROUPBYevent_dateORDERBYevent_date;8.7 总结与下一篇预告组件作用Hot Tier (Fluss)Arrow 格式秒级新鲜度3 天内数据Cold Tier (Iceberg/Paimon)Parquet 格式分钟级 Compaction长期存储Tiering Service自动将热数据 Compaction 为 Parquet 写入冷存储Union Read用户一个 SQL 自动覆盖冷热数据透明合并外部引擎Spark/Trino 可直接读 Cold Tier Parquet 数据下一篇我们将转向生产运维——如何部署 Fluss 集群、配置 Prometheus Grafana 监控、以及 Top 10 常见故障排查。本文基于 Apache Fluss 0.9.1。项目 GitHub: https://github.com/apache/fluss