从分钟到秒级:我们用 Fluss + Paimon 替换掉 Kafka + Iceberg,实时宽表终于不用 Flink 死扛了
从分钟到秒级:我们用 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),数据中台彻底变轻了
更多推荐
所有评论(0)