Apache Iceberg 深入
概述
Iceberg 是一种开放表格式:在对象存储/文件系统之上提供事务、快照、Schema 演化等能力,让数据湖具备数仓特性。本文讲透 Iceberg 的元数据结构、快照隔离、时间旅行、Schema/分区演化与 ACID 实现。
一、Iceberg 是什么
| 特性 | 说明 |
|---|---|
| 表格式 | 定义表如何组织在存储上 |
| 开放 | 支持 Spark/Flink/Presto 等 |
| 事务 | ACID(乐观并发) |
| 快照 | 每次写产生不可变快照 |
| 位置 | 独立于引擎 |
Iceberg 表 = 数据文件(Parquet/ORC)
+ 元数据(清单、快照、表元数据)
+ Catalog(表注册)二、元数据层级
2.1 层级结构
Catalog(表目录)
└── Table Metadata(表元数据,含当前快照)
└── Snapshot(快照)
└── Manifest List(清单列表)
└── Manifest(清单)
└── Data File(数据文件)| 层级 | 作用 |
|---|---|
| Table Metadata | 表结构、Schema、当前快照指针 |
| Snapshot | 某一时刻表的完整状态 |
| Manifest List | 该快照包含的清单列表 |
| Manifest | 一组数据文件的统计信息 |
| Data File | 实际数据(Parquet 等) |
2.2 写流程
1. 写新数据文件
2. 生成 Manifest(描述新文件)
3. 更新 Manifest List
4. 生成新 Snapshot(指向新状态)
5. 原子更新表元数据的当前快照指针三、快照与时间旅行
3.1 快照隔离
每次提交生成新快照,旧快照不可变
并发读基于各自快照 → 读不阻塞写| 特性 | 说明 |
|---|---|
| 快照隔离 | 读写互不阻塞 |
| 一致性 | 读固定快照 |
| 回滚 | 切回旧快照 |
3.2 时间旅行
sql
-- 查某快照的数据
SELECT * FROM orders
FOR SYSTEM_TIME AS OF '2026-08-04 10:00:00';
-- 按快照 ID 查询
SELECT * FROM orders VERSION AS OF 123456;| 能力 | 场景 |
|---|---|
| 历史查询 | 回看数据 |
| 审计 | 数据演化追踪 |
| 对比 | 快照差异 |
3.3 快照保留
| 参数 | 说明 |
|---|---|
| 保留数 | 最近 N 个快照 |
| 保留时间 | 快照有效时长 |
| 过期清理 | 删除旧快照与孤儿文件 |
四、Schema 演化
4.1 特性
| 能力 | 说明 |
|---|---|
| 加列 | 向后兼容 |
| 删列 | 历史快照仍可读 |
| 改名 | 元数据重映射 |
| 类型升级 | 受限升级 |
sql
ALTER TABLE orders ADD COLUMN discount DOUBLE;4.2 原则
| 规则 | 说明 |
|---|---|
| 向后兼容 | 新 Schema 可读旧数据 |
| 无破坏 | 删除/改名不破坏历史 |
| 类型安全 | 只能安全升级 |
五、分区演化
5.1 特性
分区策略是表的属性,随数据演化:
初始:按天分区
后期:按天 + 按城市分区
新数据用新分区,旧数据保留旧分区| 对比 | Hive 分区 | Iceberg 分区 |
|---|---|---|
| 分区变更 | 需重建表 | 元数据变更 |
| 隐藏分区 | 无 | 有(无需知道分区列) |
| 查询过滤 | 手动指定 | 自动裁剪 |
5.2 隐藏分区
sql
-- 按 ts 字段自动推断分区
CREATE TABLE events (
ts TIMESTAMP, city STRING, ...
) USING iceberg
PARTITIONED BY (days(ts), city);查询 WHERE ts = '2026-08-04'
→ 自动定位分区(无需知道分区名)六、ACID 语义
6.1 事务支持
| 特性 | 实现 |
|---|---|
| 原子性 | 快照指针原子切换 |
| 一致性 | 元数据一致 |
| 隔离性 | 乐观并发(CAS) |
| 持久性 | 元数据落存储 |
6.2 并发控制
乐观并发:
提交前检查当前快照与预期一致
冲突 → 重试或失败| 机制 | 说明 |
|---|---|
| CAS | 比较并切换快照 |
| 冲突处理 | 失败重试 |
| 多写 | 通过乐观锁协调 |
6.3 与并发写
多个写者并发 → 冲突检测 → 重试
适合:离线批写 + 实时小批量七、性能特性
| 特性 | 说明 |
|---|---|
| 查询裁剪 | 元数据统计跳过文件 |
| 分区裁剪 | 隐藏分区自动过滤 |
| 文件布局 | 小文件合并优化 |
| 向量化 | 列式读取 |
八、生态集成
| 引擎 | 支持 |
|---|---|
| Spark | Spark SQL + Iceberg |
| Flink | 流写/批读 |
| Presto/Trino | 查询 |
| Hive | 部分支持 |
scala
// Spark 使用 Iceberg
spark.table("lake.orders").write.mode("append")
.saveAsTable("lake.orders_v2")常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 快照爆炸 | 设置保留策略清理 |
| 小文件多 | 定期合并(Rewrite) |
| 并发写冲突 | 降低写入频率或改合并写 |
| 时间旅行不生效 | 快照被过期清理 |
| 分区未裁剪 | 用隐藏分区语法 |