很多团队在构建湖仓一体的时候,都会遇到同一个问题:离线计算选了 Spark,实时计算选了 Flink,数据湖格式却迟迟定不下来。不同格式对双引擎的支持深度不一样,这种差异不是简单看文档就能发现的,往往要等到真正在生产环境跑起来,才意识到兼容性问题有多棘手。

数据湖格式听起来是存储层的事,但它的设计目标比 HDFS 上的裸文件要复杂得多。它要管事务、管元数据、管增量读取、管 schema 演化。这些能力必须被 Spark 和 Flink 同时理解,才能在批流一体场景里真正落地。这也是为什么“哪个格式更好”很难一句话回答——脱离了引擎兼容性去谈格式特性,都是空谈。
为什么双引擎兼容性决定数据湖的成败
先说一个基础事实:在多数公司里,Spark 和 Flink 是同时存在、各司其职的。Spark 擅长大规模批处理,适合做离线清洗、汇聚、建宽表;Flink 擅长流式处理,适合做实时链路、监控告警、分钟级更新。湖仓一体的初衷,就是让这两类作业操作同一份数据,而不是像早期数仓那样,离线一套表,实时一套 Kafka 加 Redis,两套数据互相映照,维护成本极高。
要让 Spark 和 Flink 操作同一份数据,表面上看只需要它们都能读写相同的文件格式。但数据湖不只是文件,还包括一层“会话语义”:哪些文件属于当前快照?脏写入如何处理?事务提交后如何让所有引擎立即看到?如果 Spark 写入了一个事务,Flink 需要读到一致视图,而不是写到一半的数据。这些语义必须由数据湖格式和引擎的连接器共同实现。
当元数据存储、文件命名规则、事务协议不互通时,业务上就会出现麻烦:Flink 正在流式写入的数据,Spark 读出来的视图可能落后几分钟甚至完全缺失。这不是延迟问题,而是格式层面隔离级别和可见性控制导致的。
三大数据湖格式的兼容性现状
目前社区里最常被一起比较的是 Delta Lake、Apache Iceberg 和 Apache Hudi。它们都声称支持 Spark 和 Flink,但支持深度有明显差异。先给出核心结论对比。
| 格式 | Spark 支持 | Flink 支持 | 元数据设计 | 增量读取 | 生态成熟度 |
|---|---|---|---|---|---|
| Delta Lake | 原生,历史最久,功能最全 | 支持读写,社区维护,功能滞后 | 事务日志在表目录内 | 依赖版本信息,支持有限 | Spark 生态强 |
| Apache Iceberg | 原生支持,标准 SQL 较完善 | 原生支持,读写与 Delete 均有实现 | 独立元数据层,支持目录抽象 | Snapshot 流式读取支持良好 | 多引擎厂商标配 |
| Apache Hudi | 原生支持 | 原生支持较早,API 变化大 | Timeline 与文件分组 | COW/MOR 提供增量 | 湖仓场景支持较好 |
从这个表能看出,Iceberg 在设计上对“多引擎”最为严肃。它的元数据和数据文件完全分离,元数据通过目录独立管理,这让 Spark 和 Flink 的实现可以几乎对齐。Delta Lake 早期是 Databricks 内部项目,Spark 自然是主平台,Flink 连接器是后加的,功能和对齐程度都稍弱。Hudi 在实时写入场景很强,但它的 Timeline 和文件分组抽象相对复杂,Flink 和 Spark 的实现细节因为这个抽象而略微分化。
不过这些差异也在快速变化。数据湖格式每年都有大版本更新,今天的连接器支持情况可能半年后就过时了。所以更重要的事情是理解设计逻辑:格式是把自己的元数据作为一等公民暴露给引擎,还是只通过特定引擎的 SDK 暴露?前者更容易保持跨引擎兼容性。
兼容性背后的关键机制
要理解为什么兼容性有深有浅,得先意识到数据湖不是“文件夹加文件”。它是一层事务型存储语义。以 Iceberg 为例,它的元数据分为三层:Metadata 文件描述表 schema、分区和快照;Manifest list 列出一个快照内包含的 Manifest 文件;Manifest 文件记录数据文件的位置和统计信息。每次提交都会创建新的 Metadata 文件,并更新 metadata pointer。
这种设计让所有引擎可以独立实现同一套规范,不依赖某一个引擎的特殊优化。Flink 写入 Iceberg 时,如果连接器版本一致,写出的快照结构 Spark 可以正确解析,因为两边遵守同一个格式规范。
-- Flink SQL 侧
CREATE CATALOG my_iceberg WITH (
'type' = 'iceberg',
'catalog-type' = 'hive'
);
CREATE TABLE flink_sink (id INT, name STRING) WITH (
'connector' = 'iceberg',
'catalog-name' = 'my_iceberg',
'catalog-database' = 'default',
'catalog-table' = 'iceberg_table'
);
-- Spark 侧直接读取
SELECT * FROM my_db.iceberg_table;
这里的核心是 Catalog 配置一致。如果 Flink 和 Spark 的 Catalog 指向同一个 Hive Metastore,Iceberg 的 metadata pointer 会让两边读到同一个快照。如果你的表用输出模式落盘后 Spark 读不到,通常要先检查提交是否成功、Catalog 是否共享。
事务协议是另一个关键点。Iceberg 和 Delta Lake 都实现了乐观并发控制,写入前读取当前版本,提交时检查版本冲突。但在双引擎场景,同一个表可能同时被 Spark 和 Flink 写,如果两个引擎对锁的粒度、重试策略不一致,就会出现明明没有业务冲突,却因为提交策略不同而不停失败的现象。Hudi 的并发控制更复杂,它的 Timeline 要处理 commit、compaction 和 clean,Spark 和 Flink 的提交流程可能产生非预期交错。
双引擎场景下最常见的几个坑
纸上谈兵没用,很多问题是切到生产环境才暴露的。我这里列几个非常典型的坑。
- Flink 持续写入后,Spark 查询看不到最新数据。 这通常是快照隔离导致的。Iceberg 读的是快照 ID,Flink 如果没做 checkpoint,写入不会形成可读快照;而 Delta Lake 的 Flink 连接器把提交频率绑定在 checkpoint 上,checkpoint 时间设太大,Spark 端的数据延迟就会很高。
- 清理老文件导致下游 Flink 任务读文件失败。 当 Spark 端对表执行 expire_snapshots 或 remove_orphan_files 时,Flink 正在读取的文件可能被标记为过期,作业就会抛 FileNotFoundException。这种问题需要双方约定清理窗口,而不是无脑保留。
- Catalog 不一致导致双引擎各读各的。 在同一个 Hive Metastore 上,如果 Spark 用 HiveCatalog,Flink 用 HadoopCatalog,表元数据可能指向不同的 location,表面上是一张表,实际是两个独立的数据。排查这类问题时,先统一 catalog 类型。
- Schema 演化在双引擎下并不总是互相可见。 Spark 的 ALTER TABLE 加列,Flink 端写入可能还按旧 schema 写,导致列顺序或类型对不上。Iceberg 对此处理相对好,Hudi 则常常需要 Flink 作业重启才能感知。
- Hudi 的 compaction 和 Flink 写入互相干扰。 Flink 持续产生小文件,Spark 端触发 compaction,如果策略设置激进,写入 job 可能出现插入阻塞或桶冲突。更平滑的方式是在 Flink 侧开启异步 compaction,让 Spark 只做查询。
这些坑有一个共同点:它们不是某个引擎单独的问题,而是跨引擎协作时的语义协调失败。选型时不能只跑一遍官方 quickstart,必须把双引擎并发读写、schema 演化、文件生命周期这些场景纳入验证范围。
不同团队怎么选:一个务实对比
选型没有绝对正确,只有是否符合现状。下面从几个典型团队形态出发讨论。
如果你们公司重 Spark、轻 Flink,实时场景只是对 Kafka 做一层简单 ETL,Delta Lake 会很顺手。因为它 Spark 端体验最好,Flink 端虽然弱一点,但你的 Flink 作业大概率只是写 append 流,不涉及复杂 upsert,连接器能力不足不会被触发。维护成本也最低。
如果你们的核心诉求是让 Spark、Flink、以及未来的 Trino、StarRocks 都去读同一份数据,Iceberg 会更合适。它的元数据抽象和统一规范是跨引擎自然对齐的桥梁。尤其是当你未来可能引入数据分析团队用 Trino 查询时,Iceberg 的社区支持最广泛。
如果你们希望 Flink 能做到实时 upsert,让下游 Spark SQL 直接读取最新状态,Hudi 的 Merge on Read 会很合适。但你需要接受更高的运维复杂度——Timeline 和文件分组的概念需要团队真正吃透,否则一个文件冲突就能让整个链路阻塞。
| 团队特征 | 推荐格式 | 核心原因 |
|---|---|---|
| Spark 为主,Flink append 流 | Delta Lake | Spark 生态最成熟,Flink 基本需求够用 |
| 多引擎(Spark/Flink/Trino)共享数据 | Iceberg | 元数据规范统一,多引擎兼容性最好 |
| 高并发实时 upsert,Flink 双写 | Hudi | Copy on Write 和 Merge on Read 支持丰富 |
| 团队小,基建能力有限 | Iceberg | 部署简单,兼容性坑最少 |
但是要注意,以上只是起点。已有 Hive 数据迁移也很关键。Iceberg 和 Hudi 能基于 Hive 表格式迁移,Delta Lake 也有 migration 工具,不过都需要测试分区布局和 schema 差异。
落地建议:从验证到演进
如果你现在正要确定数据湖格式,我建议按以下步骤推进。
第一步,搭建一个最小验证环境:一个 Spark 集群,一个 Flink 集群,同一个 Hive Metastore。准备三张表,分别用 Delta、Iceberg、Hudi 创建。执行一个固定流程:Flink 流式写入 1 分钟,Spark 批读校验数据量;Spark 插入数据后,Flink 用增量方式读取;同时再跑几个并发任务测试事务冲突。
第二步,测试恢复场景。让 Flink 作业在写入一半时失败,然后重启,再让 Spark 查询表,检查是否出现不可恢复的错误。这比跑通基本功能更能反映兼容性。还要检查 Flink 的 checkpoint 配置,确认提交频率。
如果你们已经定下格式,就要建立双引擎的规范使用文档。比如:
- Flink 写入必须开启 checkpoint,并设定合理的提交间隔;
- Spark 读表时使用固定的 HiveCatalog 或 REST Catalog,不要混合 catalog 类型;
- 清理过期快照必须有全局窗口,并通知所有实时作业;
- Spark 和 Flink 连接器版本必须统一,建议在 CI 中同时跑两个引擎的读写测试。
最后强调一句,数据湖格式的兼容性不是静态的。从长期演进看,Iceberg 这类以规范为中心的格式,更容易在快速变化的引擎生态中生存。但不是每个人都必须追新,重要的是你清楚自己的场景需要什么能力,以及这个能力的代价是什么。
回到开头的问题,数据湖格式对 Spark 和 Flink 双引擎兼容性的影响,本质是存储层对计算层承诺的一致性边界。选对格式,批流一体是加分项;选错格式,双引擎协作可能反过来成为负担。希望这份对比能帮你形成自己的判断框架。
原创文章,作者:fudengji,如若转载,请注明出处:https://fudengji.cn/article/1044/