Iceberg 跨环境迁移
功能
本文档演示把一张 Iceberg 从源集群迁移到目标集群。迁移前和迁移后文件都存储在 hdfs 上。
迁移前提
- 创建两个集群: 一个集群做为源集群,一个集群作为目标集群。
- 迁移期间,source table 必须停止写入并进入只读窗口,这是离线迁移的硬性前提。窗口内不得再发生
INSERT、UPDATE、DELETE、MERGE、compaction、snapshot 清理或 orphan 文件清理等会改变表文件集合的操作。
版本策略
工具端版本
本项目迁移工具侧统一使用 Iceberg 1.11.0,原因是:
- Iceberg
1.11.0包含了对rewrite_table_path的较完整修复,并原生支持带分区统计文件的表。 rewrite_table_path过程始终保留format-version,因此用1.11.0改写出的 v2 元数据,对1.6.1的兼容性与更低版本没有差异。- 统一固定一个工具端版本,可以避免“按情况切换”带来的执行不确定性。
源集群创建表和样例数据
在源集群执行下面 SQL。创建源表时还不涉及迁移任务,使用源集群日常读写 Iceberg 表的 Spark 环境进入 spark-sql:
1spark-sql --master local[2] \
2 --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
3 --conf spark.sql.catalog.spark_catalog=org.apache.iceberg.spark.SparkSessionCatalog \
4 --conf spark.sql.catalog.spark_catalog.type=hive \
5 --conf spark.sql.storeAssignmentPolicy=ANSI
1drop namespace if EXISTS spark_catalog.demo CASCADE;
2CREATE NAMESPACE IF NOT EXISTS spark_catalog.demo;
3
4CREATE TABLE spark_catalog.demo.orders (
5 id BIGINT,
6 user_name STRING,
7 amount DECIMAL(10, 2),
8 dt STRING
9)
10USING iceberg
11PARTITIONED BY (dt)
12LOCATION '/warehouse/demo.db/orders';
13
14INSERT INTO spark_catalog.demo.orders VALUES
15 (1, 'alice', 10.00, '2026-06-29'),
16 (2, 'bob', 20.00, '2026-06-29'),
17 (3, 'carol', 30.00, '2026-06-30');
18
19SELECT dt, count(*) AS record_count
20FROM spark_catalog.demo.orders
21GROUP BY dt
22ORDER BY dt;


目标集群准备环境
下载 iceberg-spark-runtime-3.5_2.12-1.11.0.jar,然后把此文件上传到服务器的 /mnt 目录下
1curl https://bmr-build-archives.bj.bcebos.com/toolchain/iceberg/iceberg-spark-runtime-3.5_2.12-1.11.0.jar \
2 > iceberg-spark-runtime-3.5_2.12-1.11.0.jar
在目标集群准备迁移专用 Spark 环境:
1cp -r /opt/bmr/spark3 /opt/bmr/spark3-iceberg-governor
2
3rm -rf /opt/bmr/spark3-iceberg-governor/jars/iceberg-spark-runtime-*.jar
4rm -rf /opt/bmr/spark3-iceberg-governor/jars/iceberg-bundled-guava-*.jar
5
6cp /mnt/iceberg-spark-runtime-3.5_2.12-1.11.0.jar \
7 /opt/bmr/spark3-iceberg-governor/jars/
预期迁移专用环境中只保留 Iceberg 1.11.0 runtime。iceberg-spark-runtime-3.5_2.12-1.11.0.jar 已包含 Iceberg relocated Guava 类,通常不需要再额外放 iceberg-bundled-guava-1.11.0.jar。
目标集群执行迁移
生成迁移文件 file-list
在目标集群执行以下命令。
使用迁移专用 Spark SQL 启动。启动前先清理旧 Spark 环境变量,避免 /opt/bmr/spark3-iceberg-governor/bin/spark-sql 被已有 SPARK_HOME 或 SPARK_CONF_DIR 指回原 Spark 环境,导致 Iceberg 1.11.0 runtime 未生效:
1## 切换到 spark 账户
2## su - spark
3unset SPARK_HOME
4unset SPARK_CONF_DIR
5export SOURCE_METASTORE_URI=thrift://master-76f5e19:9083
6export TARGET_METASTORE_URI=thrift://master-b0f35cf:9083
7
8/opt/bmr/spark3-iceberg-governor/bin/spark-sql \
9 --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
10 --conf spark.sql.catalog.source_catalog=org.apache.iceberg.spark.SparkCatalog \
11 --conf spark.sql.catalog.source_catalog.type=hive \
12 --conf spark.sql.catalog.source_catalog.uri=${SOURCE_METASTORE_URI} \
13 --conf spark.sql.catalog.target_catalog=org.apache.iceberg.spark.SparkCatalog \
14 --conf spark.sql.catalog.target_catalog.type=hive \
15 --conf spark.sql.catalog.target_catalog.uri=${TARGET_METASTORE_URI}
注意:source_catalog、target_catalog 这类自定义 catalog 名称应使用 org.apache.iceberg.spark.SparkCatalog。org.apache.iceberg.spark.SparkSessionCatalog 只适合替换 Spark 默认 catalog,catalog 名必须是 spark_catalog;如果把 SparkSessionCatalog 配成 source_catalog,执行 CREATE NAMESPACE source_catalog.demo 时会报 Delegated SessionCatalog is missing。
rewrite_table_path 会把表 metadata 中的源路径前缀改写为目标路径前缀,并在 staging 目录生成重写后的 metadata 文件和 copy plan。
1CALL source_catalog.system.rewrite_table_path(
2 table => 'demo.orders',
3 source_prefix => 'hdfs://master-76f5e19:8020/warehouse/demo.db/orders',
4 target_prefix => 'hdfs://master-b0f35cf:8020/warehouse/target-demo.db/orders',
5 staging_location => 'hdfs://master-b0f35cf:8020/tmp/iceberg-migration/demo.orders/staging'
6);
返回示例:
100001-961e7b9a-a9af-4276-9b24-adde82a7e7c0.metadata.json hdfs://master-b0f35cf:8020/tmp/iceberg-migration/demo.orders/staging/file-list 1 0
字段含义:
| 字段 | 说明 |
|---|---|
latest_version |
重写后的最新 metadata json 文件名。它位于 staging_location 下,后续用于目标 catalog 注册。 |
file_list_location |
copy plan 文件路径。按 apache-iceberg-1.11.0 的实现,它通常是 staging_location/file-list。 |
rewritten_manifest_file_paths_count |
重写后的 manifest file path 数量。 |
rewritten_delete_file_paths_count |
重写后的 delete file path 数量。 |


目标集群执行文件迁移
file_list_location 是 copy plan 输出文件路径。按 apache-iceberg-1.11.0 的实现,rewrite_table_path 直接返回一个以 file-list 结尾的单文件路径,文件内容是逗号分隔的 copy plan,每行通常包含 source_path,target_path。迁移执行器应逐行复制,既包括源表已有 data/delete/metadata 文件,也包括 staging 目录下重写后的 metadata/manifest 文件。
建议先按 20 个 executor 起步,每个 executor 2 cores、8g memory、2g memory overhead;如果 copy plan 的文件量继续增大,优先增加 executor 个数,再考虑加大单 executor 内存。
1unset SPARK_HOME
2unset SPARK_CONF_DIR
3
4/opt/bmr/spark3-iceberg-governor/bin/spark-shell \
5 --master yarn \
6 --deploy-mode client \
7 --name iceberg-migration-copy \
8 --conf spark.dynamicAllocation.enabled=false \
9 --conf spark.executor.instances=20 \
10 --conf spark.executor.cores=2 \
11 --conf spark.executor.memory=8g \
12 --conf spark.executor.memoryOverhead=2g \
13 --conf spark.driver.memory=4g
执行以下的 Scala 代码完成文件迁移。copyPlanLocation 是 rewrite_table_path 输出的 file-list 路径。statusOutput是存储 copy 状态的临时目录。
1import java.net.URI
2
3import org.apache.hadoop.fs.{FileSystem, Path}
4import org.apache.hadoop.io.IOUtils
5import org.apache.spark.util.SerializableConfiguration
6
7val copyPlanLocation = "hdfs://master-b0f35cf:8020/tmp/iceberg-migration/demo.orders/staging/file-list"
8val statusOutput = "hdfs://master-b0f35cf:8020/tmp/iceberg-migration/demo.orders/copy-status"
9val tempSuffix = ".iceberg-copying"
10val copyParallelism = math.max(spark.sparkContext.defaultParallelism, 40)
11
12val hadoopConf = new SerializableConfiguration(spark.sparkContext.hadoopConfiguration)
13
14{
15 val copyPlanPath = new Path(copyPlanLocation)
16 val copyPlanFs = FileSystem.get(new URI(copyPlanLocation), hadoopConf.value)
17 require(copyPlanFs.exists(copyPlanPath), s"copy plan does not exist: $copyPlanLocation")
18 require(copyPlanFs.isFile(copyPlanPath), s"copy plan must be a single file: $copyPlanLocation")
19}
20
21def parseCopyTask(line: String): Option[(String, String)] = {
22 val trimmed = line.trim
23 if (trimmed.isEmpty || trimmed.startsWith("#")) {
24 None
25 } else {
26 val parts = trimmed.split(",", 2)
27 require(parts.length == 2, s"invalid copy plan line: $trimmed")
28 Some(parts(0).trim -> parts(1).trim)
29 }
30}
31
32def copyOne(source: String, target: String, conf: org.apache.hadoop.conf.Configuration): String = {
33 val srcPath = new Path(source)
34 val dstPath = new Path(target)
35 val tmpPath = new Path(target + tempSuffix)
36 val srcFs = FileSystem.get(new URI(source), conf)
37 val dstFs = FileSystem.get(new URI(target), conf)
38
39 try {
40 val sourceSize = srcFs.getFileStatus(srcPath).getLen
41 if (dstFs.exists(dstPath)) {
42 val targetSize = dstFs.getFileStatus(dstPath).getLen
43 if (targetSize == sourceSize) {
44 s"$source $target SKIPPED $sourceSize $targetSize "
45 } else {
46 throw new IllegalStateException(
47 s"target exists with different size: source=$sourceSize target=$targetSize target=$dstPath"
48 )
49 }
50 } else {
51 val parent = dstPath.getParent
52 if (parent != null && !dstFs.exists(parent)) {
53 dstFs.mkdirs(parent)
54 }
55 if (dstFs.exists(tmpPath)) {
56 dstFs.delete(tmpPath, false)
57 }
58
59 val in = srcFs.open(srcPath)
60 val out = dstFs.create(tmpPath, true)
61 try {
62 IOUtils.copyBytes(in, out, conf, false)
63 } finally {
64 try out.close() finally in.close()
65 }
66
67 val targetSize = dstFs.getFileStatus(tmpPath).getLen
68 if (targetSize != sourceSize) {
69 throw new IllegalStateException(
70 s"copy size mismatch: source=$sourceSize target=$targetSize source=$srcPath target=$tmpPath"
71 )
72 }
73 if (!dstFs.rename(tmpPath, dstPath)) {
74 throw new IllegalStateException(s"rename failed: $tmpPath -> $dstPath")
75 }
76 s"$source $target SUCCESS $sourceSize $targetSize "
77 }
78 } catch {
79 case e: Throwable =>
80 try {
81 if (dstFs.exists(tmpPath)) {
82 dstFs.delete(tmpPath, false)
83 }
84 } catch {
85 case _: Throwable =>
86 }
87 s"$source $target FAILED -1 -1 ${Option(e.getMessage).getOrElse(e.getClass.getName)}"
88 }
89}
90
91val copiedResults = spark.sparkContext
92 .textFile(copyPlanLocation, copyParallelism)
93 .flatMap(parseCopyTask)
94 .repartition(copyParallelism)
95 .mapPartitions { iter =>
96 val conf = hadoopConf.value
97 iter.map { case (source, target) => copyOne(source, target, conf) }
98 }
99 .persist()
100
101copiedResults.saveAsTextFile(statusOutput)
102
103val statusCounts = copiedResults.map(_.split(" ", 6)(2)).countByValue()
104statusCounts.foreach(println)
105require(statusCounts.getOrElse("FAILED", 0L) == 0L, s"copy failed, see $statusOutput")
106
107copiedResults.unpersist()
执行完成后,再检查 statusOutput 中的结果。FAILED 必须为 0,SUCCESS + SKIPPED 必须等于 copy plan 的总行数。
生产环境不建议退回到单机 shell 循环复制;如果要继续抽象,可以把这段逻辑封装成独立的 Spark 作业,但仍然保持按 copy plan 逐行执行的语义。
在环境里输入:paste然后粘贴整块代码,然后输入 Ctrl + D结束输入。


目标集群注册目标 Catalog
1unset SPARK_HOME
2unset SPARK_CONF_DIR
3export TARGET_METASTORE_URI=thrift://master-b0f35cf:9083
4
5/opt/bmr/spark3-iceberg-governor/bin/spark-sql \
6 --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
7 --conf spark.sql.catalog.target_catalog=org.apache.iceberg.spark.SparkCatalog \
8 --conf spark.sql.catalog.target_catalog.type=hive \
9 --conf spark.sql.catalog.target_catalog.uri=${TARGET_METASTORE_URI}
metadata_file的值是 rewrite_table_path 的 target_prefix 参数值 + /metadata/+ rewrite_table_path 输出的文件名
1drop namespace if EXISTS target_catalog.demo_target CASCADE;
2CREATE NAMESPACE IF NOT EXISTS target_catalog.demo_target;
3
4
5CALL target_catalog.system.register_table(
6 table => 'demo_target.orders',
7 metadata_file => 'hdfs://master-b0f35cf:8020/warehouse/target-demo.db/orders/metadata/00001-961e7b9a-a9af-4276-9b24-adde82a7e7c0.metadata.json'
8);
返回示例:
1current_snapshot_id total_records_count total_data_files_count
2123456789 3 2

目标集群验证目标表
注册完成后,验证目标表应使用目标集群实际生产读写环境,也就是 Iceberg 1.6.1.2-baidu 对应的原 Spark 环境,而不是迁移专用的 Iceberg 1.11.0 环境。这样可以确认 1.11.0 生成的目标 metadata 能被目标环境当前引擎正常读取。
1spark-sql --master yarn --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
2 --conf spark.sql.catalog.spark_catalog=org.apache.iceberg.spark.SparkSessionCatalog \
3 --conf spark.sql.catalog.spark_catalog.type=hive
1SELECT dt, count(*) AS record_count
2FROM spark_catalog.demo_target.orders
3GROUP BY dt
4ORDER BY dt;
5
6SELECT snapshot_id, committed_at, operation, summary
7FROM spark_catalog.demo_target.orders.snapshots
8ORDER BY committed_at DESC
9LIMIT 5;
10
11## 可以看到文件路径为新地址
12SELECT file_path FROM spark_catalog.demo_target.orders.files;
目标表的记录数、分区分布、snapshot、data file 数量应与迁移源 snapshot 的预期一致。




评价此篇文章
