从分钟到秒级:我们用 Fluss + Paimon 替换掉 Kafka + Iceberg,实时宽表终于不用 Flink 死扛了

📅 更新于 2026-05-21 | 🏷️ Fluss · Paimon · 湖流一体 · 实时数仓 · 架构升级
摘要:上一代湖仓一体架构中,Kafka + Iceberg 的组合存在数据冗余、实时更新代价高、Lambda 架构维护负担重三大痛点。本文完整复盘我们引入 Apache Fluss + Apache Paimon 构建「湖流一体」新底座的全过程,从组件替换逻辑、核心数据流设计、性能实测收益到迁移方案,毫无保留。如果你也在被实时宽表和流批两套代码折磨,这篇实战经验或许能帮你省下半年折腾。


引言:当 Iceberg 的分钟级延迟成为业务天花板

上一代数据中台,我们基于 Apache Iceberg 构建了经典的湖仓一体架构。存算分离、多引擎共享、ACID 事务——这些能力稳定支撑了两年多的业务增长。

但最近半年,三个问题越来越尖锐:

  • 实时宽表场景下,Merge-On-Read 的写放大让集群不堪重负。
  • Kafka 和 Iceberg 两份存储,数据冗余 + 一致性校验,运维成本居高不下。
  • 流批两套代码的 Lambda 架构,需求一变动,两个地方都要改,开发速度跟不上业务节奏。

我们意识到:湖仓一体的分钟级延迟,已经成为实时风控、秒级报表、AI 特征供给这些新场景的硬瓶颈。是时候在数据湖之上,加一层真正的实时流存储了。

经过调研和 POC,我们选定了 Fluss + Paimon 的组合。这篇文章,就是这次架构升级的完整复盘。


一、旧架构:Iceberg + Kafka 做了很多事,但每件事都留了尾巴

老读者都知道,我们的计算引擎层跑在 Kubernetes 上,核心组件如下:

Spark 批处理 ────┐
Flink 流处理 ────┼──→ Iceberg 表(MinIO 对象存储)
Trino 即席查询 ──┘ ↑
Hive Metastore + MySQL 元数据

实时链路单独一套:Kafka → Flink → Iceberg

这套架构做了很多事,但每件事都留了尾巴:

痛点 具体表现 对团队的影响
数据冗余 实时链路 Kafka 存一份,离线 Iceberg 再存一份,两份数据可能不一致 半夜对账成了常规操作
实时更新代价高 Iceberg 的 Merge-On-Read 在大量 Upsert 场景下写放大严重 宽表拼接任务经常 OOM
Lambda 维护重 实时和离线两套代码,改一个逻辑两个地方都要动 新需求交付周期长
延迟只能到分钟级 Commit 间隔 + 小文件合并,端到端延迟很难压到 30 秒以内 实时风控场景无法接受
元数据压力 HMS 在高频 DDL 下成为瓶颈 偶尔雪崩,影响全局

说白了,Kafka 负责“快”,Iceberg 负责“稳”,但它们之间有一条缝。这条缝靠 Flink 来粘合,而 Flink 的状态越积越重,维护成本越来越高。


二、新架构:Fluss + Paimon,让流和湖长在一起

我们的目标很明确:

  • 用一个可查询的流存储替换 Kafka,既能像消息队列一样快,又能像表一样被 SQL 直接查询。
  • 用一个与流存储深度集成的湖格式替换 Iceberg,让热数据自动归档为冷数据,查询时自动联合。
  • 尽量不改变上层应用:DolphinScheduler 调度、自研 PSC 采集引擎、Trino 即席查询、Kyuubi 批处理入口全部保留。

新架构的核心变化只有两处:

原:Kafka → Flink → Iceberg → MinIO
新:Fluss(热数据)→ 自动归档 → Paimon(冷数据)→ MinIO

但这两处变化,解决了上一代架构的所有尾巴。

新架构分层总览

应用层(数据门户、BI、指标中心)
↓ JDBC
服务层(即席查询、API 网关、结果缓存)
↓ JDBC
治理层(元数据、质量、脱敏、权限、审计)
↓ JDBC
调度开发层(DolphinScheduler + 自研采集引擎 PSC + SQL 编辑器)
↓ JDBC / Flink SQL
计算引擎层(Kyuubi / Trino / Spark / Flink)
├── 实时写入 → Fluss 集群(流存储,秒级新鲜度)
│ ↓ 自动分层归档
└── 批量/即席查询 → Paimon 表(湖格式,冷数据高性能分析)
↓
存储层:MinIO(Parquet 文件)+ MySQL(元数据)


三、Fluss 和 Paimon,到底分别解决了什么?

3.1 Fluss:不只是 Kafka 的替代品

Fluss 是阿里巴巴开源的新一代流存储系统,定位超越消息队列:

Fluss 的能力 对 Kafka 的升级
主键表 + Upsert Kafka 是日志,Fluss 是表。原生支持按主键更新和删除
列式存储 基于 Arrow,查询只读需要的列,网络开销降低 10 倍
自动归档到 Paimon 热数据自动下沉,冷热分层查询自动联合
Flink 一等公民 原生 Flink SQL Connector,读写和 CDC 订阅都比 Kafka 稳定
SQL 可查询 无需额外 OLAP 引擎,Fluss 自身就支持 SQL 点查和即席分析

3.2 Paimon:为流而生的湖格式

Paimon 的前身是 Flink Table Store,与 Flink 生态深度绑定:

Paimon 的能力 对 Iceberg 的升级
CDC 原生支持 Iceberg 消费 CDC 需要额外处理;Paimon 直接消费变更日志
主键表 Compaction 后台自动合并,不需要手动的 Compaction 作业
标签快照 类似 Iceberg 的时间旅行,但创建和管理更轻量
流式读写 同一张表可以被流任务写入,同时被批任务读取,互不阻塞

3.3 组合后的化学反应

Flink SQL (实时写入)
↓
Fluss 主键表 (秒级可见,支持 Upsert 和点查)
↓ 自动归档(按时间或数据量)
Paimon 表 (历史数据,列存高性能分析)
↓
Trino / Spark 查询时自动 Union 两部分数据

一句话:以前 Flink 既要管实时写入,又要管状态拼接,现在 Fluss 管写入和热数据,Paimon 管冷数据和归档,Flink 的压力被大幅卸载。


四、三个关键技术细节:部分列更新、Delta Join、流式裁剪

4.1 部分列更新:实时宽表不用再写复杂 Flink 作业

以前的宽表拼接,需要 Flink 多流 Join,状态膨胀到几百 GB,动不动就 OOM。

现在利用 Fluss 的部分列更新能力:

  • 订单表、用户表、商品表分别写入 Fluss 主键表的不同列
  • Fluss 自动按主键合并,查询时直接得到完整宽表
  • Flink 作业只需简单的单表写入,不再维护大状态

实测下来,同一场景的 Flink 作业,内存消耗从 120GB 降到了 18GB。

4.2 Delta Join:内存和 CPU 消耗下降超 86%

传统双流 Join 需要维护两个流的状态,而利用 Fluss 的实时 KV 点查,一条流当主表,另一条流通过索引直接查询 Fluss 中的维表数据,状态只需维护一份。

某风控宽表场景,Delta Join 替代双流 Join 后,TaskManager 堆内存从 32G 降至 4G,CPU 利用率降低 86%。

4.3 流式列裁剪:即席查询不再扫全表

Kafka 查询本质是消费全量消息,而 Fluss 的 Arrow 列式存储,让 Trino 查询时只读取需要的列。结合主键索引和分区裁剪,百亿级表的单条点查延迟压到了 50 毫秒以内。


五、新旧架构核心对比

维度 旧(Kafka + Iceberg) 新(Fluss + Paimon)
数据新鲜度 分钟级 秒级/亚秒级
存储冗余 Kafka + Iceberg 双份 Fluss 一份,自动归档
实时更新 MERGE INTO 写放大 主键表原生 Upsert
宽表拼接 Flink 多流 Join,状态重 部分列更新,Flink 轻量化
查询效率 Kafka 全行扫描 列式存储 + 列裁剪
架构模式 Lambda(流批独立) Kappa(流批统一)
开发效率 两套代码 统一 SQL,宽表开发效率提升 60%+

六、迁移方案:我们的分步走策略

6.1 数据迁移

历史 Iceberg 表通过 Spark 作业批量转换为 Paimon 表。Paimon 提供从 Iceberg 迁移的兼容工具,格式转换可以离线完成。

6.2 实时链路切换

原来 Kafka → Flink → Iceberg 的链路,改为 Flink → Fluss 写入。Flink 作业需要修改 Sink Connector,但 SQL 逻辑基本不变。切换时采用双跑验证,确认数据一致性后再下线旧链路。

6.3 对上游应用的影响

  • DolphinScheduler 调度任务:仅调整数据源配置,任务编排逻辑不变。
  • 自研 PSC 采集引擎:离线采集链路不受影响,仍写 Iceberg 或 Paimon 均可。
  • Trino 即席查询:增加 Paimon 连接器和 Fluss 连接器,前端 SQL 编辑器无需改动。

七、结语

从 Iceberg 到 Paimon,从 Kafka 到 Fluss,不是追逐新技术,而是业务倒逼的结果。当实时风控要求秒级响应,当宽表拼接把 Flink 集群压到报警边缘,当两份数据的对账脚本越来越长——我们就知道,架构必须往前走一步。

Fluss + Paimon 的组合,让我们的数据新鲜度从分钟级跃升到秒级,同时 Flink 作业的内存占用大幅降低。更重要的是,流和湖的边界终于被打破了。


👍 如果你也在做实时数仓或湖仓升级,点个赞让更多同行看到这套方案。
💬 你们团队目前用 Kafka + Iceberg 还是已经升级了?遇到过什么坑?评论区聊聊。
收藏本文,下次做实时架构选型时直接翻出来对比。

🔗 延伸阅读【运维必备】Docker/K8s/Linux 高频命令速查手册(持续更新)
🔗 延伸阅读从 Hadoop 到湖流一体:数据底座的三次革命与终极选型指南
🔗 延伸阅读揭秘大型数据中台:湖仓一体架构如何支撑万亿级数据流转?
🔗 延伸阅读再见,组件地狱!我们将 7 个开源引擎替换成一个华为云 GaussDB(DWS),数据中台彻底变轻了

更多推荐