Apache Iceberg 表格式数据湖
┌─────────────────────────────────┐ │ 计算引擎 │ │ (Spark, Flink, Trino, ...) │ └──────────────┬──────────────────┘ │ 读写操作 ┌──────────────▼──────────────────┐ │ Apache Iceberg 表格式 │ │ - 元数据层 (metadata layer) │ │ - 数据层 (data layer) │ │ - 快照与事务管理 │ └──────────────┬──────────────────┘ │ 存储 ┌──────────────▼──────────────────┐ │ 对象存储 / HDFS / 本地文件系统 │ │ (Parquet / ORC / Avro 文件) │ └─────────────────────────────────┘
## 2. Iceberg 的核心架构与元数据
### 2.1 表结构的分层设计
Iceberg 表由三层元数据构成,分别存放在数据湖的存储目录中:
- **`metadata/` 目录**:存放全局元数据文件
- **`data/` 目录**:存放实际数据文件(按分区组织)
一张 Iceberg 表的文件布局示例:
demo_db.my_table/ ├── metadata/ │ ├── v1.metadata.json # 版本1的元数据根文件 │ ├── v2.metadata.json # 版本2的元数据根文件 │ ├── snap-1234567890.avro # 快照 manifest list │ └── ... └── data/ ├── event_time_month=2024-01/ │ ├── 00000-0-xxx.parquet │ └── 00001-1-xxx.parquet ├── event_time_month=2024-02/ │ └── 00002-0-yyy.parquet └── ...
### 2.2 元数据体系详述
整个元数据体系由四个关键组件组成,形成一棵稳定的树状结构:
- **元数据文件 (metadata file)**:存储表级别的信息,包括表 Schema、分区规范、当前快照 ID、快照日志等。每次表结构变更都会生成一个新的 `metadata.json`,文件名包含递增版本号。
- **快照 (snapshot)**:代表表在某个时间点的数据状态,是对 Manifest List 的引用,并记录添加和删除的数据文件。
- **Manifest List (清单列表)**:指向一组 Manifest File,每个 Manifest List 文件通常对应一个分区规范下的部分数据。
- **Manifest File (清单文件)**:包含数据文件列表以及每个数据文件的分区值、列统计信息(如列的最小值、最大值、null 计数等)。
这种层级结构使得查询规划极其高效:引擎从 metadata 文件找到当前快照,加载 Manifest List,然后根据查询过滤条件,只读取相关的 Manifest File,并进一步利用文件级别的统计信息跳过不符合条件的数据文件。
### 2.3 快照隔离与乐观并发
Iceberg 的写操作采用**乐观并发控制**(OCC)。每个写事务在提交时会检查当前元数据版本是否与开始写入时一致,如果不一致(即存在冲突),则根据配置进行重试或报错。这种“快照隔离”保证了读操作始终看到一致的快照,不会被未完成的写入污染,而写操作之间互不阻塞,仅在提交时做最终校验。
## 3. Iceberg 的核心特性详解
### 3.1 分区演进 (Partition Evolution)
传统 Hive 表的分区一旦定义就极难修改,如果需要调整分区粒度或逻辑,必须重写整个表。Iceberg 支持 **在不停机、不重写数据的情况下修改表的分区规范**。
- 旧数据保持原有的物理分区布局不变
- 新写入的数据按照新的分区规范进行组织
- 查询时 Iceberg 会根据元数据中记录的分区信息自动适配,确保正确读取
- 分区值的转换基于列的值,而非文件的物理目录路径,从而避免“目录依赖”
示例:将按月分区的表改为按日分区:
```sql
ALTER TABLE prod.orders
SET PARTITION SPEC (days(order_date));
3.2 模式演进 (Schema Evolution)
Iceberg 完整支持安全的 Schema 变更,包括添加、删除、重命名、重新排序列,以及更新列的数据类型(支持安全的类型提升,如 int→long)。
- 所有 Schema 变更都是元数据操作,不会触发数据重写
- Iceberg 使用唯一列 ID来跟踪列的身份,即使重命名列也不会丢失映射
- 新写入的文件使用最新 Schema,旧文件沿用其写入时的 Schema,读取时自动进行 Schema 适应(例如缺失列会用 NULL 填充)
3.3 隐藏分区 (Hidden Partitioning)
用户无需在查询时手动指定分区过滤器,Iceberg 会根据列的值和分区转换自动完成分区剪裁。例如,如果按 month(event_time) 分区,查询 WHERE event_time = '2025-03-01' 时,Iceberg 会推导出它属于 2025-03 分区,从而只扫描该月的数据文件,有效避免全表扫描。
3.4 时间旅行与快照管理
Iceberg 保留每次写操作生成的快照,用户可通过快照 ID 或时间戳查询历史数据状态:
-- 根据时间戳查询过去的数据
SELECT * FROM prod.orders
TIMESTAMP AS OF '2025-03-07 10:00:00';
-- 根据快照 ID 查询
SELECT * FROM prod.orders
VERSION AS OF 1234567890123456789;
并可通过管理命令清理过期快照,或回滚到特定版本:
CALL sys.rollback_to_snapshot('prod.orders', 1234567890123456789);
3.5 文件级索引(列统计)
Manifest File 中存储了每列的最小值、最大值、null 计数等统计信息。对于有序写入的数据,这种文件级索引能极大加速范围查询。例如,如果查询 WHERE user_id > 1000,且某个数据文件的 user_id 最大值小于 1000,则该文件会被自动跳过,无需解压和读取。
4. 写入与读取流程
4.1 写入流程(以 Spark 为例)
- 规划阶段:引擎获取当前表的最新快照和元数据。
- 数据写入:将 DataFrame 按照新的分区规范写入到数据文件(Parquet 等)。
- 构建元数据:为新的数据文件创建 Manifest File,生成 Manifest List,并指向新的快照。
- 提交:通过原子操作(如对象存储的 rename 或 HDFS 的原子重命名)将新的
metadata.json文件写入,并将表的当前快照指针更新。如果出现冲突则根据策略重试。
Iceberg 的提交并未使用外部锁服务,而是依赖存储层的原子文件创建与重命名特性,这对 S3 等对象存储需要通过 HDFS 接口或专门的 metastore 辅助,但 Iceberg 提供了 HadoopCatalog、HiveCatalog、JDBCCatalog、RESTCatalog 等多种 catalog 实现,可兼容不同环境。
4.2 读取流程
- 从 catalog 获取表的当前 metadata 文件位置。
- 加载 metadata,找到目标快照。
- 加载 Manifest List,根据查询谓词筛选 Manifest File。
- 加载相关 Manifest File,利用列统计信息进一步过滤数据文件。
- 根据剩余的数据文件列表,并行读取 Parquet/ORC 等数据,应用查询和过滤条件。
整个过程在协调节点(如 Spark Driver)上完成,计算被推送到执行器,具备极低的规划开销。
5. 与其它数据湖格式的对比
| 特性 | Apache Iceberg | Delta Lake | Apache Hudi |
|---|---|---|---|
| 核心设计思想 | 元数据表格式 | 事务日志 + 文件管理 | 数据湖增量处理 |
| 并发控制 | 乐观并发 | 乐观并发(有检查点) | 乐观并发(支持 OCC/MVCC) |
| 模式演进 | 完整支持且安全 | 支持,但部分操作需重写 | 支持 |
| 分区演进 | ✅ 原生支持 | 不支持(需重写数据) | 不支持 |
| 隐藏分区 | ✅ | ✅ | 部分支持 |
| 时间旅行 | 快照 / 时间戳 | 版本号 / 时间戳 | 提交时间 / 即时时间 |
| 查询引擎支持 | Spark, Flink, Trino, Presto, Hive, Impala | Spark, Flink, Trino, Presto | Spark, Flink, Hive, Presto, Trino |
| 主要适用场景 | 通用数据湖分析,多引擎共享,缓慢变化维度 | Spark 生态紧密,机器学习 | 流式 upsert,增量 ETL |
6. 入门实战:基于 Spark 创建一个 Iceberg 表
这里以本地 Spark 环境为例演示 Iceberg 的基本操作。请确保已配置好 Spark 和 Iceberg 依赖(Iceberg Spark Runtime)。
6.1 环境配置与启动
# 使用 spark-shell 时添加 Iceberg 包
spark-shell \
--packages org.apache.iceberg:iceberg-spark-runtime-3.3_2.12:1.4.3 \
--conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
--conf spark.sql.catalog.local=org.apache.iceberg.spark.SparkCatalog \
--conf spark.sql.catalog.local.type=hadoop \
--conf spark.sql.catalog.local.warehouse=/tmp/iceberg-warehouse
6.2 创建表
-- 使用 local 目录 catalog
CREATE TABLE local.db.sales (
id BIGINT NOT NULL,
product STRING,
amount DECIMAL(10,2),
sold_at TIMESTAMP
) USING iceberg
PARTITIONED BY (months(sold_at));
6.3 插入数据
INSERT INTO local.db.sales VALUES
(1, 'Laptop', 1200.00, TIMESTAMP '2025-01-15 10:00:00'),
(2, 'Mouse', 25.50, TIMESTAMP '2025-02-10 14:30:00'),
(3, 'Keyboard', 79.99, TIMESTAMP '2025-03-05 09:15:00');
6.4 读取与分区剪裁验证
-- 查询二月的销售数据,Iceberg会自动只读取sold_at_month=2025-02分区
SELECT * FROM local.db.sales
WHERE sold_at BETWEEN '2025-02-01' AND '2025-02-28';
可以用 EXPLAIN 指令查看实际的扫描计划,确认分区过滤生效。
6.5 Schema 变更与分区演进
-- 添加新列
ALTER TABLE local.db.sales ADD COLUMNS (discount DECIMAL(3,2));
-- 变更分区:从按月改为按日(新数据按天分区)
ALTER TABLE local.db.sales SET PARTITION SPEC (days(sold_at));
-- 再次插入包含新列的数据
INSERT INTO local.db.sales VALUES
(4, 'Monitor', 300.00, TIMESTAMP '2025-04-20 11:00:00', 0.05);