- 数据湖
- 大数据
- 数据存储
【免费下载链接】iceberg
Apache Iceberg
本文是一份面向数据工程师与平台架构师的 Apache Iceberg 表迁移技术指南,围绕 Iceberg 官方文档中的 Table Migration 主线,系统讲解"全量数据迁移"与"原地元数据迁移"两类路径的取舍,并深入剖析原地迁移中三大核心动作 ——Snapshot Table(快照表)、Migrate Table(迁移表)与Add Files(追加文件)的原理、Spark SQL 调用方式、参数语义与底层实现。读完本文,你将能够根据业务对停机时间、数据隔离性和历史保留的需求,为 Hive、Delta Lake 等存量表制定并落地一套安全、可回滚的 Iceberg 迁移方案。
一、两种迁移路径:全量数据迁移 vs 原地元数据迁移
Apache Iceberg 支持将其他格式的存量表转换为 Iceberg 表,官方将迁移方式划分为两大类:
全量数据迁移(Full Data Migration):把源表的全部数据文件拷贝到一张全新的 Iceberg 表中。其最大优点是新表与源表完全隔离——后续对源表的任何清理、删除操作都不会影响新表;代价是迁移速度慢,且需要双倍存储空间。在实践中,全量迁移通常借助以下手段完成:
- Create-Table-As-Select(CTAS) 语法;
- INSERT INTO 语句;
- 各类 Change-Data-Capture(CDC)同步管道。
原地元数据迁移(In-Place Metadata Migration):保留源表现有数据文件不动,仅在数据之上叠加 Iceberg 元数据(metadata)。这种方式速度更快且无需复制数据,但新表与源表之间并不完全隔离——如果源表侧有任何进程对数据文件执行了清理(vacuum),新表也会连带受到影响。本指南的核心内容即围绕原地元数据迁移展开。
Iceberg 的原地元数据迁移共包含三个重要动作:Snapshot Table、Migrate Table与Add Files,三者分别应对"无停机探测性迁移""原地接管替换"与"迁移后的增量补齐"三类场景。
二、Snapshot Table:零停机创建独立快照表
Snapshot Table动作会以源表相同的 schema 与分区方式,创建一张名称不同的全新 Iceberg 表;动作执行期间及执行之后,源表都保持不变,源表上的既有读写任务可以继续运行,不受任何影响。
整个流程分为三步:
- 创建新表:以源表的元数据(schema、partition spec 等)为模板,创建一张名称不同、位置独立的 Iceberg 表。源表上的 Readers 与 Writers 可以继续正常工作,无需停机。
提交数据文件:将源表所有分区的数据文件全部提交到新 Iceberg 表中。此时源表仍然不变,读侧(Readers)可以先切换到新 Iceberg 表。
切换写侧:待读侧验证无误后,将全部 Writers 切换到新 Iceberg 表。当所有写任务完成切换后,迁移流程即宣告完成。
从源码实现来看,SnapshotTable是定义在 api 模块的 Action 接口 中的一流动作(一等公民),其方法链包括:as(destTableIdent)指定新表标识、tableLocation(location)指定新表位置、tableProperties(...)/tableProperty(...)设置表属性、executeWith(ExecutorService)指定并行读文件的线程池,以及ignoreMissingFiles()用于跳过已消失的源数据文件;执行结果通过importedDataFilesCount()返回导入的数据文件数。具体的不可变实现由 core 模块的 BaseSnapshotTable 通过 Immutables 框架生成。
在 SnapshotTableSparkAction 中,doExecute()的核心逻辑清晰地呈现了"先暂存、后提交"的安全模式:
- 通过
stageDestTable()创建暂存表(StagedSparkTable); - 强校验源表位置与暂存表位置不得重叠(包括互为前缀的情况),否则会混合两张表的文件而直接报错;
- 调用
SparkTableUtil.importSparkTable(...)为源表生成 Iceberg 元数据; - 成功后
commitStagedChanges()提交;一旦抛错则进入finally分支执行abortStagedChanges()回滚暂存变更。
值得注意的细节是,快照表默认写入gc.enabled=false并标记snapshot=true表属性(见destTableProps()),因为快照表与源表共享数据文件、并非这些文件的唯一所有者,所以禁止对快照表执行会物理删除数据文件的expire_snapshots之类操作;仅影响元数据的 Iceberg 删除(如 DELETE 语句产生的 delete 文件)仍然允许。相应地,对源表执行 DELETE 移除原始数据文件,也会破坏快照表的完整性。
三、Migrate Table:原地接管并替换源表
Migrate Table动作同样会创建一张与源表 schema、分区方式一致的 Iceberg 表,但区别在于:动作执行过程中会锁定并从 catalog 中移除(drop)源表。因此,Migrate Table 要求在执行前停止所有正在操作源表的修改任务;支持 Iceberg 的读者(Readers)可以继续读取。
流程同样分为三步:
- 停止写侧:停止所有与源表交互的 Writers。
备份并建新表:创建一张与源表相同标识和元数据(schema、partition spec 等)的 Iceberg 表,同时把源表重命名为备份表(默认后缀
_BACKUP_),以备失败时回滚。提交并清理:将源表所有分区的数据文件提交到新 Iceberg 表,然后删除源表;此时 Writers 即可开始向新 Iceberg 表写入。迁移完成后,默认保留的备份表(如
db.sample_BACKUP_)可通过drop_backup=true参数选择删除。
Migrate 的实现同样遵循"暂存 + 提交"的安全模式。在 MigrateTableSparkAction 中可以看到:
- 常量
BACKUP_SUFFIX = "_BACKUP_"定义了默认备份名规则,构造函数即生成backupIdent; renameAndBackupSourceTable()先把源表重命名为备份表,从而"冻结"源表、暂停一切修改,并为其后的暂存建表腾出位置;- 若备份名已存在则抛出
AlreadyExistsException,源表不存在则抛出NoSuchTableException; - 后续流程从备份表(而非源表)导入数据文件到暂存的 Iceberg 表;
- 失败时
restoreSourceTable()会把备份表重命名回原标识完成回滚;成功后若开启dropBackup()则删除备份表; destTableProps()会为迁移表写入migrated=true属性,并继承源表位置(putIfAbsent(LOCATION, sourceTableLocation())),确保新表原地接管原数据目录。
从源码还可推断,Migrate 对源 catalog 有较强约束:checkSourceCatalog要求源 catalog 必须是SparkSessionCatalog,即当前实现只支持从 Spark Session Catalog 中的非 Iceberg 表进行迁移。此外,Migrate 会拒绝迁移使用不支持文件格式(仅支持 Avro、Parquet、ORC)的分区表,也会因分桶无法在 Iceberg 中保留而直接失败。
四、Add Files:补齐迁移窗口期的新增数据
在完成初始迁移(无论采用 Snapshot Table 还是 Migrate Table)之后,经常会发现还有部分数据文件未被迁移。这些文件通常来自并发写入者——它们在迁移过程中或迁移结束后仍继续向源表写入数据。具体到不同格式:
- 对于 Hive 表,这些未迁移文件是新增的 Hive 数据文件;
- 对于 Delta Lake 表,这些未迁移文件是新产生的 snapshot(版本)。
Add Files动作正是为将这些遗漏文件纳入 Iceberg 表而设计的。它不创建新表,而是直接向一张已存在的 Iceberg 表追加来自 Hive/文件型表的数据文件,且可以只导入指定分区;Iceberg 会为这些文件生成元数据但不会移动文件本身。从 AddFilesProcedure 的源码看,其源标识还支持以parquet.path、orc.path、avro.path形式直接指向文件型表位置(isFileIdentifier()负责识别这类命名空间)。
使用 Add Files 前必须明确两个重要警告:
- 不校验 schema:该过程不会分析文件 schema 是否与 Iceberg 表匹配,添加 schema 不一致的文件会引发数据问题;
- 文件所有权转移:一旦添加完成,Iceberg 会将这些文件视为自己拥有的文件,后续
expire_snapshot等操作将能够物理删除这些文件;因此只要可能,应优先使用migrate或snapshot而非add_files。
五、实战:从 Hive 迁移到 Iceberg
Hive 的 ORC、Parquet、Avro 三种文件格式均可迁移到 Iceberg。由于 Hive 表没有 snapshot 概念,迁移过程本质上是"用现有 schema 创建一张新的 Iceberg 表,并把所有分区的数据文件一次性提交进去";初始迁移之后的新增数据文件,则通过 Add Files 动作持续补齐。
这些动作由 Spark 集成模块以 Spark Procedure(存储过程)的形式提供,已打包进 Spark runtime jar(见 releases 下载页 中的 Spark runtime 产物)。对应的过程定义与参数细节可参考 Spark Procedures 文档。
5.1 Snapshot Hive 表
CALL catalog_name.system.snapshot('db.source', 'db.dest')snapshot过程的完整参数如下(详见 spark-procedures.md#snapshot):
| 参数 | 是否必填 | 类型 | 说明 |
|---|---|---|---|
source_table | ✔️ | string | 要快照的源表名 |
table | ✔️ | string | 要创建的新 Iceberg 表名 |
location | string | 新表的位置(默认交给 catalog 决定) | |
properties | map<string, string> | 添加到新表的属性 | |
parallelism | int | 文件读取线程数(默认 1) | |
ignore_missing_files | boolean | 为 true 时跳过找不到的源数据文件而不是失败(默认 false) |
输出为imported_files_count(long),即添加到新表的文件数。典型用法:
-- 在 catalog 默认位置创建引用 db.sample 的隔离快照表 db.snap CALL catalog_name.system.snapshot('db.sample', 'db.snap'); -- 在指定位置 /tmp/temptable/ 创建快照表 CALL catalog_name.system.snapshot('db.sample', 'db.snap', '/tmp/temptable/');快照表适合测试场景:测试完成后用DROP TABLE清理即可。在 SnapshotTableProcedure 中可以看到各参数的默认行为——parallelism必须大于 0,ignore_missing_files默认为 false,且源表与目标表名不能相同。
5.2 Migrate Hive 表
CALL catalog_name.system.migrate('db.sample')migrate过程会复制源表的 schema、分区、属性和位置,并用源表数据文件填充新表(详见 spark-procedures.md#migrate):
| 参数 | 是否必填 | 类型 | 说明 |
|---|---|---|---|
table | ✔️ | string | 要迁移的表名 |
properties | map<string, string> | 新 Iceberg 表的属性 | |
drop_backup | boolean | 为 true 时不再保留原表作为备份(默认 false) | |
backup_table_name | string | 备份表名称(默认table_BACKUP_) | |
parallelism | int | 文件读取线程数(默认 1) | |
ignore_missing_files | boolean | 为 true 时跳过找不到的源数据文件而不是失败(默认 false) |
输出为migrated_files_count(long),即追加到 Iceberg 表的文件数。典型用法:
-- 迁移并在新表上添加属性 foo='bar' CALL catalog_name.system.migrate('spark_catalog.db.sample', map('foo', 'bar')); -- 不添加额外属性,直接迁移当前 catalog 中的表 CALL catalog_name.system.migrate('db.sample');5.3 从 Hive 表追加文件到 Iceberg 表
CALL spark_catalog.system.add_files( table => 'db.tbl', source_table => 'db.src_tbl' )add_files过程参数(详见 spark-procedures.md#add_files):
| 参数 | 是否必填 | 类型 | 说明 |
|---|---|---|---|
table | ✔️ | string | 要追加文件的目标 Iceberg 表 |
source_table | ✔️ | string | 文件来源表,也支持file_format`.`path形式的路径 |
partition_filter | map<string, string> | 只导入指定分区的文件 | |
check_duplicate_files | boolean | 是否阻止添加已存在于表中的文件(默认 true) | |
parallelism | int | 文件读取线程数(默认 1) |
输出为added_files_count(long)与changed_partition_count(long,未知时为空)。注意:当表属性compatibility.snapshot-id-inheritance.enabled为 true 或表格式版本大于 1 时,changed_partition_count会返回 NULL。
-- 只导入 part_col_1='A' 分区中的文件 CALL spark_catalog.system.add_files( table => 'db.tbl', source_table => 'db.src_tbl', partition_filter => map('part_col_1', 'A') ); -- 从 parquet 文件型表位置导入全部文件 CALL spark_catalog.system.add_files( table => 'db.tbl', source_table => '`parquet`.`path/to/table`' );六、实战:从 Delta Lake 迁移到 Iceberg
Delta Lake 采用 Parquet 文件格式,并支持时间旅行(time travel)与版本管理。与 Hive 不同,从 Delta Lake 迁移时通常希望保留全部历史,因此常见的做法是把 Delta Lake 的所有 snapshot 都迁移过来,以维持数据历史。
目前 Iceberg 对 Delta Lake 只支持Snapshot Table动作:由于 Delta Lake 表维护事务日志,源表所有可用的事务会被按顺序提交到新 Iceberg 表,作为对应的事务。初始迁移之后 Delta 表新增的数据文件,会包含在其对应事务中,后续通过Add Transaction动作(Add Files 的变体,目前仍在开发中)追加到新表。
6.1 启用迁移能力
iceberg-delta-lake模块不会随 Spark、Flink 引擎运行时打包,需要额外添加以下依赖:
iceberg-delta-lake(Maven 坐标org.apache.iceberg:iceberg-delta-lake)delta-standalone-0.6.0(io.delta:delta-standalone_2.13:0.6.0)delta-storage-2.2.0(io.delta:delta-storage:2.2.0)
该模块基于Delta Standalone 0.6.0构建与测试,支持的 Delta Lake 表协议版本为:minReaderVersion: 1、minWriterVersion: 2(协议版本语义参见 Delta Lake 官方的 Table Protocol Versioning 说明)。
6.2 API 与默认实现
模块提供了DeltaLakeToIcebergMigrationActionsProvider接口,包含动作snapshotDeltaLakeTable:将一张已有 Delta Lake 表快照为 Iceberg 表。接口的默认实现可通过以下方式获取:
DeltaLakeToIcebergMigrationActionsProvider defaultActions = DeltaLakeToIcebergMigrationActionsProvider.defaultActions()snapshotDeltaLakeTable动作会读取 Delta Lake 表的事务,在一个 Iceberg 事务内将其转换为一张具有相同 schema 与分区方式的新 Iceberg 表,源 Delta Lake 表保持原样。新表可以独立读写而不影响源表,但快照使用的是源表的数据文件——通过从源表 schema 生成的 name-to-id 映射(name mapping)来读取。快照表上的 INSERT / OVERWRITE 产生的新文件会写入快照表自身的位置,该位置默认与源 Delta Lake 表相同,也可通过 API 另行指定。
SnapshotDeltaLakeTable动作接口定义在 delta-lake 模块,其必需输入与配置方式如下:
| 必需输入 | 配置方式 | 说明 |
|---|---|---|
| 源表位置 | 参数sourceTableLocation | 源 Delta Lake 表的位置 |
| 新表标识 | APIas(TableIdentifier) | 指定新 Iceberg 表的 namespace 与表名 |
| Iceberg catalog | APIicebergCatalog(Catalog) | 用于创建新表的 catalog |
| Hadoop 配置 | APIdeltaLakeConfiguration(Configuration) | 读取源 Delta Lake 表所需的 Hadoop 配置 |
输出为imported_files_count(long),即添加到新表的文件数。动作执行后还会为新表写入以下默认属性:
| 属性名 | 值 | 说明 |
|---|---|---|
snapshot_source | delta | 标记该表由 Delta Lake 表快照而来 |
original_location | 源 Delta Lake 表位置 | 源表的绝对路径 |
schema.name-mapping.default | 由 schema 推导的 JSON name mapping | 用于读取 Delta Lake 数据文件的 name mapping |
6.3 Java 调用示例
import org.apache.iceberg.catalog.TableIdentifier; import org.apache.iceberg.catalog.Catalog; import org.apache.hadoop.conf.Configuration; import org.apache.iceberg.delta.DeltaLakeToIcebergMigrationActionsProvider; String sourceDeltaLakeTableLocation = "s3://my-bucket/delta-table"; String destTableLocation = "s3://my-bucket/iceberg-table"; TableIdentifier destTableIdentifier = TableIdentifier.of("my_db", "my_table"); Catalog icebergCatalog = ...; // 从 Spark 等引擎获取,或通过 CatalogUtil.loadCatalog 创建 Configuration hadoopConf = ...; // 从引擎获取、且已配置访问 Delta Lake 表所需文件系统的 Hadoop Configuration DeltaLakeToIcebergMigrationActionsProvider.defaultActions() .snapshotDeltaLakeTable(sourceDeltaLakeTableLocation) .as(destTableIdentifier) .icebergCatalog(icebergCatalog) .tableLocation(destTableLocation) .deltaLakeConfiguration(hadoopConf) .tableProperty("my_property", "my_value") .execute();七、源码实现与测试验证
7.1 Action 接口体系
三类迁移动作在 Iceberg 中都被建模为Action接口,位于 api 模块的 actions 包:
- SnapshotTable:提供
as、tableLocation、tableProperties、tableProperty、executeWith、ignoreMissingFiles等方法,结果返回importedDataFilesCount; - MigrateTable:提供
tableProperties、tableProperty、dropBackup、backupTableName、executeWith、ignoreMissingFiles等方法,结果返回migratedDataFilesCount。
它们的不可变实现分别由 BaseSnapshotTable 与 BaseMigrateTable 通过 Immutables 注解生成,保证 Action 配置后不可变、可安全复用。
7.2 Spark 侧的暂存提交机制
在 Spark 集成中,两个核心 Action 实现都依赖StagedSparkTable+StagingTableCatalog的暂存机制:
- SnapshotTableSparkAction 要求源 catalog 必须是 session catalog(
spark_catalog),并校验新旧表位置不得重叠; - MigrateTableSparkAction 要求源 catalog 为
SparkSessionCatalog,通过"重命名备份 → 暂存建表 → 提交 → 失败回滚/成功删备份"的流程实现原子替换。
对应的过程层封装(SnapshotTableProcedure、MigrateTableProcedure、AddFilesProcedure)把 Action 的参数校验与默认值落实为可调用的 Spark Procedure。
7.3 测试覆盖
仓库中的测试为上述行为提供了可复现的验证:
- TestSnapshotTableAction 覆盖了并行任务快照、位置重叠时报错、非重叠位置等场景;
- TestMigrateTableAction 覆盖了并行任务下的迁移流程。
八、方案选型与注意事项
| 场景 | 推荐方案 | 原因 |
|---|---|---|
| 不想停机、先验证再切换 | Snapshot Table | 源表全程可用,读侧先切、写侧后切,迁移失败无影响 |
| 可以接受短暂停机、希望原地接管 | Migrate Table | 原标识原地替换,自动保留备份可回滚,位置继承源表 |
| 迁移后有并发写入者遗留的数据 | Add Files | 按分区精准补齐新文件,无需重建表 |
几点必须牢记的约束:
- 隔离性差异:Snapshot 与 Migrate 都共享源数据文件,源表侧一旦 vacuum/删除数据文件,新表会受影响;也不要对快照表执行
expire_snapshots等物理删除操作; - 格式支持:原地迁移仅支持 Avro、Parquet、ORC 文件格式,分桶表无法迁移(分桶语义无法保留);
- Add Files 的风险:不校验 schema,且添加后的文件归 Iceberg 所有、可被物理删除,能不用就不用;
- Delta Lake 迁移特殊性:需额外引入三个依赖,且当前只支持 Snapshot 路径,增量事务补齐(Add Transaction)仍在开发中。
通过上述三种动作的组合,你可以在控制停机时间与数据冗余成本的前提下,将 Hive、Delta Lake 等存量表平稳迁移到 Apache Iceberg,并获得 Iceberg 的事务、版本管理与 ACID 能力。更细的语法与参数,可继续参阅 Hive 迁移文档、Delta Lake 迁移文档 与 Spark Procedures 参考。
- 数据湖
- 大数据
- 数据存储
【免费下载链接】iceberg
Apache Iceberg
相关推荐
Apache Iceberg Hive表迁移完全指南
Apache Iceberg Hive表迁移完全指南 概述 Apache Iceberg作为新一代数据湖表格式,相比传统Hive表具有诸多优势,包括ACID事务
大数据数据湖OLAP数据存储Apache Iceberg Hive表迁移完全指南
Apache Iceberg Hive表迁移完全指南 概述 在现代数据架构中,将传统Hive表迁移到Apache Iceberg表已成为提升数据管理能力的重要步
数据湖大数据数据存储Apache Iceberg 迁移指南:使用 snapshotDeltaLakeTable 将 Delta Lake 表完整迁移至 Iceberg
Apache Iceberg 迁移指南:使用 snapshotDeltaLakeTable 将 Delta Lake 表完整迁移至 Iceberg Delta
数据湖大数据数据存储
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考