流湖一体 · Streaming Lakehouse

Flink + Fluss + Paimon
实时数仓分层架构技术文档

从 ODS→DWD→DWS→ADS 分层实时数仓出发,拆解 Flink 计算、Fluss 流存储、Paimon 湖表格式如何协同,覆盖窗口聚合、主键 upsert、写入路径、Union Read、以及 Doris 与 Fluss+Paimon 的集成实现。

Apache Flink · 计算 Apache Fluss · 流存储(热) Apache Paimon · 湖表(冷) Apache Doris · 实时大屏
01

三个组件各自是做什么的

这三者常被混为一谈,但它们处在不同的层次。理解分工是理解整套架构的前提。

计算层 · COMPUTE Flink 流/批计算引擎 清洗 · 打宽 · 窗口聚合 CDC 摄入 · upsert 加工 流存储 · HOT Fluss 实时热数据层 秒级新鲜度 · 可查询 替代 Kafka 中间层 自动分层 湖表 · COLD Paimon 湖表格式(历史) LSM · 主键 upsert 多模态 · 治理/回滚
图 1 · 计算(Flink)→ 热存储(Fluss)→ 冷存储(Paimon)的分层关系
组件层次负责什么类比
Flink计算引擎数据清洗、关联打宽、窗口聚合、CDC 摄入、upsert 加工逻辑"厨师" — 负责加工
Fluss流存储(热)存最近一小段时间的实时热数据,秒级可见、可像表一样查询;替代传统 Kafka 中间层"备菜台" — 现用现取
Paimon湖表格式(冷)存历史/定型数据,主键表 LSM 支持高频 upsert,提供治理、回滚、多模态"冷库" — 长期保存
关键区分 Fluss 不是湖表格式。Paimon / Iceberg 是湖表格式(定义数据如何在湖上组织);Fluss 是流存储层,负责那段"实时热数据",并把冷数据自动分层(tiering)到 Paimon。三者的黄金组合是 Fluss(热/实时) + Paimon(冷/历史) + DLF(统一元数据/治理)。
02

分层实时数仓:有无 Fluss 的两套架构

场景:ODS→DWD→DWS→ADS 四层,每层有 Flink 加工逻辑,ADS 产出 5 分钟粒度统计。层与层之间既要能"流转"驱动下一层,又要能"查询/回溯"。

方案 A · 没有 Fluss(Kafka + 湖/OLAP 双写)

上层:Kafka 逐层流转(不可查) 下层:为"可查/回溯"额外落湖一份(双写冗余) ODSKafka topic DWDKafka topic DWSKafka topic ADS→ OLAP 大屏 Flink Flink 5min聚合 Iceberg 副本 Iceberg 副本 Iceberg 副本 Iceberg 副本 ⚠ 每层存两份:Kafka(能流不能查)+ 湖(能查不实时)
图 2 · 传统架构:Kafka 逐层流转 + 为可查性额外落湖,双写冗余
这套架构的问题
  • 存储双写、严重冗余:每层 Kafka 一份 + 落湖一份,四层下来重复存多遍,运维两套系统。
  • Kafka 不能直接查询/裁剪:中间层排查问题只能另起消费作业,或看有延迟的落湖数据。
  • 中间层不可回溯:Kafka 默认只留几天,改逻辑重跑历史、或作业恢复都很难。
  • 实时与历史割裂:热数据在 Kafka/OLAP、历史在湖,做"最新+历史"完整口径要应用层自己拼。
  • 链路长、组件多:每层维护 topic、schema、consumer group、compaction + 落湖作业。

方案 B · 有 Fluss(流湖一体,每层一张"流湖一体表")

每层 = 一张 Fluss 表(热层,可流转 + 可直接查) ODSFluss 表 DWDFluss 表 DWSFluss 表 · 5min聚合 ADSFluss 表 Flink 读增量 Paimon(冷 / 历史) 自动分层 tiering · 主键表 · 一份存储 Union Read:一条 SQL = 热(Fluss) + 冷(Paimon)
图 3 · 流湖一体:每层是 Fluss 表,冷数据自动分层到 Paimon,Union Read 合并读
带来的好处
  • 一份存储,消除冗余:热在 Fluss、冷在 Paimon,整体一份,没有 Kafka+湖双写。
  • 每层可直接查询/裁剪:中间层是表,排查直接 SELECT,开发体验同离线数仓。
  • 每层可回溯/重算:历史在 Paimon 完整保留,改逻辑重跑不受 Kafka 保留期限制。
  • 实时+历史统一:Union Read 自动合并,一条 SQL 拿"历史到秒级最新",一套口径。
  • 链路短、schema/SQL 不变:兼容 Kafka 协议,DLF 上"湖流一体"开关即用。
03

DWS 主键 Upsert:按天累计、每 5 分钟刷新

场景升级:DWS 产出按天粒度汇总(dt 为主键的订单数/GMV/UV),但要求每 5 分钟刷新当天累计值 —— 同一个 dt 当天会被 upsert 覆盖近 288 次。这是典型的主键 upsert / 可变流,不是 append-only。

难点没有 Fluss有 Fluss + Paimon 主键表
Upsert 载体Kafka 不支持 upsert → 靠 OLAP 主键表兜底;落 Iceberg 则高频行级更新读放大严重Paimon 主键表原生 upsert(LSM 为高频更新而生)
当天最新值新鲜度OLAP 里秒级,但与湖历史割裂Fluss 热层秒级 + Union Read 合并历史,一份数据
Flink 状态压力全压 Flink(UV 去重状态大),checkpoint 重累加类可下推 Paimon 聚合引擎,Flink 状态大幅减轻
中间层可查/回溯Kafka 不可直查;重算要重放+重建状态Paimon 直接 SELECT;按 dt 读历史重算;支持 branch/快照
实时值 vs 历史值两套存储、两套口径同一张表、同一口径
为什么这里 Iceberg 不合适 Iceberg 对高频 upsert 不友好:同一主键一天更新几百次会产生大量 delete file / 小文件,读放大严重、compaction 压力大。这正是"Iceberg 更 focus 在批分析、实时可变能力弱"的地方 —— 高频 upsert 实时流不是 Iceberg 的强项,而是 Paimon 的主场。

UV 精确去重的两种实现

累加类(count/sum) → Paimon 聚合/部分更新引擎,Flink 几乎无状态 精确去重(UV) → Flink 有状态 或 Paimon 侧配合,历史从 Paimon 回补

对 count/sum 这类可累加指标,可用 Paimon 主键表的 aggregation / partial-update 引擎,增量直接进表按主键聚合,Flink 侧几乎不用维护大状态,重启也不怕丢。对 UV 这类需要精确去重的,仍建议 Flink 有状态聚合,但历史回补时从 Paimon 读,不必全靠 Kafka 重放。

04

写入路径与 Union Read:先写谁、查询发给谁

这是最容易误解的地方。核心结论:你只写一个入口(Fluss),Paimon 那份是系统自动分层出来的,不是你再 upsert 一次。

Flink upsert dt 主键 · 每5min一次 ① 只写这里 Fluss 热层 主键 upsert · 秒级可见 存最近一小段时间热数据 当天 dt 最新累计值在此 ② 自动分层 tiering service(非你写) Paimon 冷层 历史/较早版本沉淀 主键表 · LSM 合并 Union Read 按主键(dt)合并 Fluss 最新版覆盖 Paimon 旧版 = 历史 + 秒级最新 ③ 查询读"逻辑表",引擎去两边取
图 4 · 写入只入 Fluss ① → 自动分层到 Paimon ② → 查询由 Union Read 合并 ③

三个高频疑问的直接回答

Q1 · 是"先 upsert Fluss 再 upsert Paimon"吗?

不是。你只 upsert Fluss(热层入口);Paimon 那份由后台分层服务(table.datalake.enabled='true')自动搬下去,你不用写第二条链路。

Q2 · 这一刻查询,SQL 发给 Fluss 还是 Paimon?

你查的是逻辑表,不用自己指定。实际读谁取决于引擎与表名后缀(见下表)。

Q3 · "秒级可见 + 一套口径"怎么来的?

最新更新在 Fluss 秒级可见;Union Read 按主键把 Fluss 最新版覆盖合并 Paimon 历史版,返回一张表 —— 同主键、同语义,所以是一套口径,不存在"实时库一个口径、离线库另一个"的割裂。

查询写法实际读取结果
SELECT … FROM t(无后缀 · Flink)Fluss 热 + Paimon 冷(Union Read)历史 → 秒级最新,完整
SELECT … FROM t$lake只读 Paimon只到上次分层的历史(不含最新热数据)
StarRocks / Trino默认湖侧读(Paimon)取决于引擎对实时层的支持
05

多引擎新鲜度矩阵:谁能读到 Fluss 热层

这是选型的硬约束:并非所有引擎都能读到 Fluss 热层。"秒级 Union Read" 只在部分引擎上完整成立。

引擎对 Fluss 的读取能力读到热层最新 upsert?典型定位
Flink完整 Union Read✅ 能(秒级)数据加工 / 实时读
Spark批量 Union Read + 增量(Fluss 1.0 起)✅ 能批处理 / 回补
Doris通过 Fluss Catalog 做 Union Read✅ 能实时大屏(本文推荐)
StarRocks湖侧读(lake-only)❌ 只到 Paimon湖分析(非实时热)
Trino湖侧读(lake-only)❌ 只到 Paimon湖分析 / 联邦查询
对"大屏用 Trino/StarRocks 承载"的直接影响 如果大屏是 Trino / StarRocks,它查询时走的是 Paimon(湖侧读),读不到 Fluss 热层的当天最新 upsert。此时大屏新鲜度 = Fluss→Paimon 分层间隔(通常分钟级),不是秒级。

好在你的需求是 5 分钟粒度 —— 只要把分层间隔配到 ≤ 5min,湖侧读够用;接受"大屏新鲜度 = 分层间隔"即可。若确需接近实时,则换 Doris(见下一节)。
06

Doris + Fluss+Paimon 集成:大屏近实时的实现

当大屏需要"当天累计值接近实时可见"时,Doris 是 Trino/StarRocks 之外少数能对 Fluss 直接做 Union Read 的引擎。

Doris 凭什么能读到热层(而 StarRocks/Trino 不能)

维度StarRocks / Trino(湖侧读)Doris(Fluss Catalog + Union Read)
集成方式Paimon 湖连接器 —— "从湖的角度看表"Fluss Catalog 直连 —— "从 Fluss 的角度看表"
读取路径只读 PaimonFluss 热层 + Paimon 冷层合并
当天最新 upsert读不到(要等分层)✅ 读得到(秒级)
大屏新鲜度= 分层间隔(分钟级)秒级 / 近实时
主键合并口径只有 Paimon 里的版本Fluss 最新版覆盖 Paimon 旧版,一套口径
Flink按 dt upsert Fluss 热层当天最新累计 · 秒级 自动分层 Paimon 冷层历史定型 Doris · Fluss Catalog Union Read(热+冷按主键合并) 📊 实时大屏(近实时)
图 5 · Doris 通过 Fluss Catalog 直读热+冷,大屏拿到当天秒级累计值
落地要注意的坑
  • 版本与集成成熟度:Doris 对 Fluss 的 Union Read 是较新能力,上线前用实际版本压测 upsert 合并语义(同 dt 取最新)与热层读稳定性。Fluss 1.0 于 2026-09 才发布,配套生态仍在快速演进。
  • 热层查询的并发/性能:大屏高并发高频刷新,读热层比纯读 Paimon 压力大。建议对"当天"这一行做结果缓存/限流(5min 粒度本就不必每秒打),并评估读放大。
  • "秒级"是否真需要:你的原始需求是 5min 粒度。若分钟级够,StarRocks 湖侧读 + 分层间隔 ≤5min 更省事、并发压力更小。换 Doris 的价值在于"当天值要接近实时"。
  • 引擎统一性:若现有大屏已是 StarRocks,为一个场景引入 Doris 会多一套运维。权衡:Doris 只承载少数"要秒级"的实时大屏,其余仍走 StarRocks。
决策一句话 当天值必须接近实时 → 换 Doris 值得(能力上支持 Fluss Catalog Union Read,直读热层);分钟级够 → StarRocks/Trino 湖侧读 + 压缩分层间隔更省事。真正的"秒级 Union Read"目前在 Flink / Doris / Spark 上成立,StarRocks/Trino 不在其中。

文档说明 · 本文基于我们此前的技术讨论整理,聚焦 Flink + Fluss + Paimon 的实时数仓分层架构与场景实现,以及 Doris 与 Fluss+Paimon 的集成。技术定位说明:Fluss 由阿里 Flink 团队于 2023-07 发起,2026-08 从 Apache 孵化器毕业为顶级项目(TLP),Fluss 1.0 于 2026-09 发布;Union Read 为 Fluss 开源核心能力,非单一厂商私有。

各引擎对 Fluss 的支持程度、分层间隔配置、以及 Doris×Fluss 的具体版本兼容性,请以你目标环境的实际版本压测结果为准。