Delta Lake 迁移
Delta Lake 表迁移
Delta Lake 是一种表格式,支持 Parquet 文件格式,并提供时间旅行和版本管理功能。在将数据从 Delta Lake 迁移到 Iceberg 时,通常会迁移所有快照以保留数据的历史记录。
目前,Iceberg 支持通过 Snapshot Table 操作从 Delta Lake 迁移到 Iceberg 表。由于 Delta Lake 表会维护事务,因此所有可用事务都会按顺序作为事务提交到新的 Iceberg 表中。对于 Delta Lake 表,初始迁移之后新增的任何额外数据文件都会包含在相应的事务中,并随后通过 Add Transaction 操作添加到新的 Iceberg 表中。Add Transaction 操作是 Add File 操作的一个变体,目前仍在开发中。
启用从 Delta Lake 到 Iceberg 的迁移
iceberg-delta-lake 模块未随 Spark 和 Flink 引擎运行时一起打包。要启用 Delta Lake 迁移功能,所需的最低依赖项如下:
兼容性
该模块使用 Delta Standalone:0.6.0 进行构建和测试,并支持以下协议版本的 Delta Lake 表:
minReaderVersion:1minWriterVersion:2
有关 Delta Lake 协议版本的更多详情,请参阅 Delta Lake 表协议版本管理。
API
iceberg-delta-lake 模块提供了一个名为 DeltaLakeToIcebergMigrationActionsProvider 的接口,其中包含有助于从 Delta Lake 转换为 Iceberg 的操作。支持的操作包括:
snapshotDeltaLakeTable:将现有的 Delta Lake 表快照为 Iceberg 表
默认实现
iceberg-delta-lake 模块还提供了该接口的默认实现,可以通过
DeltaLakeToIcebergMigrationActionsProvider defaultActions = DeltaLakeToIcebergMigrationActionsProvider.defaultActions()将 Delta Lake 表快照到 Iceberg
snapshotDeltaLakeTable 操作会读取 Delta Lake 表的事务,并在一次 Iceberg 事务中将其转换为一个具有相同 schema 和分区方式的新 Iceberg 表。原始 Delta Lake 表保持不变。
新创建的表可以进行修改或写入,而不会影响源表,但该快照使用的是原始表的数据文件。现有数据文件会被添加到 Iceberg 表的元数据中,并可以通过由原始表 schema 创建的名称到 ID 的映射进行读取。
当对快照执行插入或覆盖操作时,新文件会被放置在快照表的位置。该位置默认与源 Delta Lake 表的位置相同。用户也可以为快照表指定不同的位置。
Info
由于通过 snapshotDeltaLakeTable 创建的表并不是其数据文件的唯一所有者,因此禁止它们执行 expire_snapshots 等会物理删除数据文件的操作。仅影响元数据的 Iceberg 删除操作仍然允许。此外,任何影响原始数据文件的操作都会破坏快照的完整性。针对原始 Delta Lake 表执行的 DELETE 语句会删除原始数据文件,届时 snapshotDeltaLakeTable 表将无法再访问这些文件。
用法
| 必需输入 | 配置方式 | 说明 |
|---|---|---|
| 源表位置 | 参数 sourceTableLocation | 源 Delta Lake 表的位置 |
| 新 Iceberg 表标识符 | 配置 API as | 该标识符指定了新 Iceberg 表的命名空间和表名 |
| Iceberg Catalog | 配置 API icebergCatalog | 用于创建新 Iceberg 表的 catalog |
| Hadoop 配置 | 配置 API deltaLakeConfiguration | 用于读取源 Delta Lake 表的 Hadoop 配置 |
有关详细用法和其他可选配置,请参阅 SnapshotDeltaLakeTable API
输出
| 输出名称 | 类型 | 说明 |
|---|---|---|
imported_files_count | long | 添加到新表中的文件数量 |
新增的表属性
默认情况下,以下表属性会被添加到将要创建的 Iceberg 表中:
| 属性名称 | 值 | 描述 |
|---|---|---|
snapshot_source | delta | 表明该表是从 Delta Lake 表快照而来 |
original_location | Delta Lake 表的位置 | 原始 Delta Lake 表位置的绝对路径 |
schema.name-mapping.default | 根据 schema 推导出的 JSON 名称映射 | 用于读取 Delta Lake 表数据文件的名称映射字符串 |
示例
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 = ...; // Iceberg Catalog fetched from engines like Spark or created via CatalogUtil.loadCatalog
Configuration hadoopConf = ...; // Hadoop Configuration fetched from engines like Spark and have proper file system configuration to access the Delta Lake table.
DeltaLakeToIcebergMigrationActionsProvider.defaultActions()
.snapshotDeltaLakeTable(sourceDeltaLakeTableLocation)
.as(destTableIdentifier)
.icebergCatalog(icebergCatalog)
.tableLocation(destTableLocation)
.deltaLakeConfiguration(hadoopConf)
.tableProperty("my_property", "my_value")
.execute();评论
登录后参与评论
KnowForge