Apache Spark

存储过程

师成师成· 更新于 2026-09-28· 阅读 76 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

Spark 存储过程

要在 Spark 中使用 Iceberg,首先配置 Spark catalog。
对于 Spark 3.x,存储过程仅在 Spark 中使用 Iceberg SQL 扩展时可用。
对于 Spark 4.0,存储过程原生受支持,无需 Iceberg SQL 扩展。但请注意,Spark 4.0 中的存储过程是区分大小写的。

用法

存储过程可以从任何已配置的 Iceberg catalog 中通过 CALL 调用。所有存储过程都位于 system 命名空间下。

CALL 支持按名称(推荐)或按位置传递参数。不支持混合使用位置参数和命名参数。

命名参数

所有存储过程参数均为命名参数。按名称传递参数时,参数的顺序可以任意,并且可以省略任意可选参数。

CALL catalog_name.system.procedure_name(arg_name_2 => arg_2, arg_name_1 => arg_1);

位置参数

按位置传递参数时,只有可选的末尾参数可以省略。

CALL catalog_name.system.procedure_name(arg_1, arg_2, ... arg_n);

快照管理

rollback_to_snapshot

将表回滚到指定的快照 ID。

若要回滚到指定的时间点,请使用 rollback_to_timestamp。

Info

此过程会使所有引用了受影响表的已缓存 Spark 计划失效。

用法

参数名是否必需类型描述
table✔️string要更新的表的名称
snapshot_id✔️long要回滚到的快照 ID

输出

输出名类型描述
previous_snapshot_idlong回滚前的当前快照 ID
current_snapshot_idlong回滚后的新的当前快照 ID

示例

将表 db.sample 回滚到快照 ID 1:

CALL catalog_name.system.rollback_to_snapshot('db.sample', 1);

rollback_to_timestamp

将表回滚到某个时间点时所处的快照状态。

信息

此过程会使所有引用受影响表的已缓存 Spark 计划失效。

用法

参数名称是否必需类型描述
table✔️string要更新的表的名称
timestamp✔️timestamp要回滚到的时间戳

输出

输出名称类型描述
previous_snapshot_idlong回滚前的当前快照 ID
current_snapshot_idlong新的当前快照 ID

示例

将 db.sample 回滚到某一天的特定时刻。

CALL catalog_name.system.rollback_to_timestamp('db.sample', TIMESTAMP '2021-06-30 00:00:00.000');

set_current_snapshot

设置表的当前快照 ID。

与回滚不同,该快照不要求是当前表状态的祖先快照。

信息

此过程会使所有引用受影响表的已缓存 Spark 计划失效。

用法

参数名称是否必需类型描述
table✔️string要更新的表名称
snapshot_idlong要设置为当前值的快照 ID
refstring要设置为当前值的快照引用(分支或标签)

必须提供 snapshot_id 或 ref 其中之一,但不能同时提供两者。

输出

输出名称类型描述
previous_snapshot_idlong回滚前的当前快照 ID
current_snapshot_idlong新的当前快照 ID

示例

将 db.sample 的当前快照设置为 1:

CALL catalog_name.system.set_current_snapshot('db.sample', 1);

将 db.sample 的当前快照设置为标签 s1:

CALL catalog_name.system.set_current_snapshot(table => 'db.sample', ref => 's1');

cherrypick_snapshot

将某个快照中的更改拣选到当前表状态中。

拣选会基于现有快照创建一个新的快照,而不会修改或删除原始快照。

只有追加(append)和动态覆盖(dynamic overwrite)类型的快照可以被拣选。

Info

此过程会使所有引用受影响表的已缓存 Spark 计划失效。

用法

参数名称是否必需类型描述
table✔️string要更新的表的名称
snapshot_id✔️long要拣选的快照 ID

输出

输出名称类型描述
source_snapshot_idlong执行拣选前表的当前快照
current_snapshot_idlong应用拣选操作后创建的快照 ID

示例

拣选快照 1

CALL catalog_name.system.cherrypick_snapshot('my_table', 1);

使用命名参数 cherry-pick 快照 1

CALL catalog_name.system.cherrypick_snapshot(snapshot_id => 1, table => 'my_table' );

publish_changes

将暂存的 WAP ID 中的变更发布到当前表状态。

publish_changes 会基于已有快照创建一个新快照,而不会修改或删除原有快照。

只有 append 和 dynamic overwrite 类型的快照可以成功发布。

如果表中存在多个具有所提供 wap_id 的快照,publish_changes 存储过程将会失败。

Info

此存储过程会使所有引用受影响表的已缓存 Spark 计划失效。

用法

参数名称是否必需类型描述
table✔️string要更新的表的名称
wap_id✔️string要从 stage 发布到 prod 的 wap_id

输出

输出名称类型描述
source_snapshot_idlong发布变更之前表的当前快照 ID
current_snapshot_idlong应用变更后创建的快照 ID

示例

使用 WAP ID 'wap_id_1' 执行 publish_changes

CALL catalog_name.system.publish_changes('my_table', 'wap_id_1');

带命名参数的 publish_changes

CALL catalog_name.system.publish_changes(wap_id => 'wap_id_2', table => 'my_table');

fast_forward

将一个分支的当前快照快进(fast-forward)到另一个分支的最新快照。

用法

参数名称是否必填类型描述
table✔️string要更新的表名
branch✔️string要执行快进的分支名
to✔️string

输出

输出名称类型描述
branch_updatedstring已被快进的分支名称
previous_reflong执行快进前的快照 ID
updated_reflong执行快进后的当前快照 ID

示例

将 main 分支快进到 audit-branch 的最新位置。

CALL catalog_name.system.fast_forward('my_table', 'main', 'audit-branch');

元数据管理

许多维护操作都可以通过 Iceberg 存储过程来执行。

expire_snapshots

Iceberg 中的每次写入/更新/删除/更新插入(upsert)/合并(compaction)都会生成一个新快照,同时保留旧的数据和元数据,以实现快照隔离和时间旅行。expire_snapshots 存储过程可用于移除不再需要的旧快照及其文件。

该存储过程会移除旧快照以及仅由这些旧快照唯一依赖的数据文件。这意味着 expire_snapshots 存储过程绝不会删除仍被未过期快照所依赖的文件。

用法

参数名称是否必需类型说明
table✔️string要更新的表的名称
older_than️timestamp在此时间戳之前的快照将被删除(默认:5 天前)
retain_lastint无论 older_than 如何设置,都需要保留的祖先快照数量(默认为 1)
max_concurrent_deletesint用于删除文件操作的线程池大小(默认情况下不使用线程池)
stream_resultsboolean当为 true 时,待删除的文件将按 RDD 分区发送到 Spark driver(默认情况下,所有文件会一次性发送到 Spark driver)。建议将此选项设置为 true,以防止文件过大导致 Spark driver 出现 OOM
snapshot_idslong 数组要过期的快照 ID 数组
clean_expired_metadataboolean当为 true 时,会清理不再被任何快照引用的元数据,例如分区规范和 schema

如果省略了 older_than 和 retain_last,将使用表的过期属性。仍被分支或标签引用的快照不会被删除。默认情况下,分支和标签永不过期,但可以通过表属性 history.expire.max-ref-age-ms 更改其保留策略。main 分支永不过期。

输出

输出名称类型描述
deleted_data_files_countlong此操作删除的数据文件数量
deleted_position_delete_files_countlong此操作删除的位置删除文件数量
deleted_equality_delete_files_countlong此操作删除的等值删除文件数量
deleted_manifest_files_countlong此操作删除的 manifest 文件数量
deleted_manifest_lists_countlong此操作删除的 manifest 列表文件数量
deleted_statistics_files_countlong此操作删除的统计信息文件数量

示例

删除早于指定日期和时间的快照,但保留最近的 100 个快照:

CALL hive_prod.system.expire_snapshots('db.sample', TIMESTAMP '2021-06-30 00:00:00.000', 100);

删除快照 ID 为 123 的快照(注意,该快照 ID 不应是当前快照):

CALL hive_prod.system.expire_snapshots(table => 'db.sample', snapshot_ids => ARRAY(123));

remove_orphan_files

用于移除未被 Iceberg 表的任何元数据文件引用、因此可视为「孤立」的文件。

用法

参数名 是否必需 类型 说明

table ✔️ string 要清理的表名

older_than ️ timestamp 移除此时间戳之前创建的孤立文件(默认为 3 天前)

location string 用于查找文件的目录(默认为表的位置)

dry_run boolean 为 true 时不实际删除文件(默认为 false)

max_concurrent_deletes int 用于执行删除文件操作的线程池大小(默认不使用线程池)

stream_results boolean 为 true 时,孤立文件将按 RDD 分区发送到 Spark driver(默认会将所有文件一次性发送到 Spark driver)。为避免大量文件导致 Spark driver 内存溢出(OOM),建议将此选项设置为 true。启用后,输出中将包含最多 20,000 个文件路径的样本

file_list_view string 用于查找文件的数据集(跳过目录列举操作)

equal_schemes map 用于视为相等的文件系统 scheme 映射。键为以逗号分隔的 scheme 列表,值为一个 scheme(默认为 map('s3a,s3n','s3'))。

equal_authorities map 用于视为相等的文件系统 authority 映射。键为以逗号分隔的 authority 列表,值为一个 authority。

prefix_mismatch_mode string 当位置前缀(scheme/authority)不匹配时的行为:

  • ERROR - 抛出异常。(默认)
  • IGNORE - 不执行任何操作。
  • DELETE - 删除文件。

prefix_listing boolean 为 true 时,通过 SupportsPrefixOperations 接口使用基于前缀的文件列举。启用此标志时,表的 FileIO 实现必须支持 SupportsPrefixOperations(默认为 false)

输出

输出名称类型说明
orphan_file_locationString此命令判定为孤立文件的每个文件的路径

示例

对该表执行 remove_orphan_files 命令的试运行(dry run),在不实际删除的情况下列出所有可能被移除的文件:

CALL catalog_name.system.remove_orphan_files(table => 'db.sample', dry_run => true);

删除 tablelocation/data 文件夹中表 db.sample 不认识的任何文件。

CALL catalog_name.system.remove_orphan_files(table => 'db.sample', location => 'tablelocation/data');

删除 files_view 视图中所有不被表 db.sample 识别的文件。

Dataset<Row> compareToFileList =
    spark
        .createDataFrame(allFiles, FilePathLastModifiedRecord.class)
        .withColumnRenamed("filePath", "file_path")
        .withColumnRenamed("lastModified", "last_modified");
String fileListViewName = "files_view";
compareToFileList.createOrReplaceTempView(fileListViewName);
CALL catalog_name.system.remove_orphan_files(table => 'db.sample', file_list_view => 'files_view');

当某个文件与元数据文件中的引用除位置前缀(scheme/authority)之外均匹配时,默认情况下会抛出错误。通过将 prefix_mismatch_mode 设置为 IGNORE,可以忽略该错误并跳过该文件。

CALL catalog_name.system.remove_orphan_files(table => 'db.sample', prefix_mismatch_mode => 'IGNORE');

通过将 prefix_mismatch_mode 设置为 DELETE,该文件仍可被删除。

CALL catalog_name.system.remove_orphan_files(table => 'db.sample', prefix_mismatch_mode => 'DELETE');

也可以通过将不匹配的前缀视为相等来删除该文件。

CALL catalog_name.system.remove_orphan_files(table => 'db.sample', equal_schemes => map('file', 'file1'));
CALL catalog_name.system.remove_orphan_files(table => 'db.sample', equal_authorities => map('ns1', 'ns2'));

列出使用前缀列表(prefix listing)得到的所有可能被删除的文件。

CALL catalog_name.system.remove_orphan_files(table => 'db.sample', prefix_listing => true);

rewrite_data_files

Iceberg 会跟踪表中的每一个数据文件。数据文件越多,存储在 manifest 文件中的元数据就越多;而小数据文件则会导致元数据过多,并因文件打开开销而降低查询效率。

Iceberg 可以借助 Spark 通过 rewriteDataFiles action 并行压缩数据文件。该操作会将小文件合并为较大的文件,从而减少元数据开销和运行时的文件打开开销。

用法

参数名称 是否必填 类型 描述

table ✔️ string 要更新的表的名称

strategy string 策略名称——binpack 或 sort。默认为 binpack 策略

sort_order string 使用 Zorder 时,写为 zorder() 内以逗号分隔的列列表。例如:zorder(c1,c2,c3)。
否则,写为以逗号分隔的排序规则,格式为 (ColumnName SortDirection NullOrder)。
其中 SortDirection 可以是 ASC 或 DESC,NullOrder 可以是 NULLS FIRST 或 NULLS LAST。
默认使用表自身的排序规则

options ️ map 供该 action 使用的选项

where ️ string 以字符串形式表示的谓词,用于过滤文件。注意,所有可能包含匹配该过滤条件的数据的文件都会被选中进行重写

选项

通用选项

名称 默认值 描述

max-concurrent-file-group-rewrites 5 同时被重写的文件组的最大数量

partial-progress.enabled false 在整个重写完成之前,先提交部分文件组

partial-progress.max-commits 10 若启用部分进度,本次重写最多允许产生的提交次数

partial-progress.max-failed-commits partial-progress.max-commits 的值 若启用部分进度,在作业失败之前最多允许的失败提交次数

use-starting-sequence-number true 使用压缩开始时快照的序号,而非新产生快照的序号

rewrite-job-order none 根据该值强制指定重写作业的顺序。

  • 如果 rewrite-job-order=bytes-asc,则先重写最小的作业组。
  • 如果 rewrite-job-order=bytes-desc,则先重写最大的作业组。
  • 如果 rewrite-job-order=files-asc,则先重写文件数量最少的作业组。
  • 如果 rewrite-job-order=files-desc,则先重写文件数量最多的作业组。
  • 如果 rewrite-job-order=none,则按作业组被规划的顺序进行重写(无特定顺序)。

target-file-size-bytes 536870912(512 MB,即 表属性 中 write.target-file-size-bytes 的默认值) 目标输出文件大小

min-file-size-bytes 目标文件大小的 75% 小于该阈值的文件,无论其他条件如何,都会被考虑进行重写

max-file-size-bytes 目标文件大小的 180% 大于该阈值的文件,无论其他条件如何,都会被考虑进行重写

min-input-files 5 文件数达到或超过此数值的任何文件组,无论其他条件如何都会被重写(该文件组至少应包含两个文件)

rewrite-all false 强制重写所有提供的文件,覆盖其他选项

max-file-group-size-bytes 107374182400(100GB)单个文件组中应重写的最大数据量。整个重写操作会按分区拆分,并在分区内按大小拆分为文件组。这有助于拆分超大分区的重写工作,否则这些分区可能因集群资源限制而无法重写。

delete-file-threshold 2147483647 数据文件需要关联的删除记录的最小数量,达到该数量后才会被考虑重写

delete-ratio-threshold 0.3 数据文件需要关联的最小删除比例,达到该比例后才会被考虑重写

output-spec-id 当前分区规范 ID 输出分区规范的标识符。重写过程中数据将被重新组织,以与输出分区方式保持一致。

remove-dangling-deletes false 重写后移除悬空的位置删除和相等删除。如果一个删除文件不适用于任何有效的数据文件,则视为悬空。启用此选项将为移除操作生成一次额外的提交。

max-files-to-rewrite null 此选项设置将被重写的合格文件数量的上限。如果未指定此选项,所有合格文件都将被重写。

Info

悬空删除文件仅根据数据序列号来移除。此操作不适用于全局相等删除,也不适用于删除条件不匹配任何数据文件的无效相等删除,同样不适用于其中的位置删除已不再匹配任何有效数据文件的位置删除文件。

排序策略的选项
名称默认值描述
compression-factor1.0Spark 排序产生的 shuffle 分区数,进而产生的输出文件数量,是基于该文件重写器所使用的输入数据文件的大小来决定的。由于压缩的缘故,磁盘上的文件大小可能无法准确反映输出文件的实际大小。此参数允许用户调整用于估算实际输出数据大小时所使用的文件大小。大于 1.0 的系数会基于磁盘文件大小生成比预期更多的文件;小于 1.0 的值则会基于磁盘大小生成比预期更少的文件。
shuffle-partitions-per-file1每个输出文件使用的 shuffle 分区数量。Iceberg 会使用自定义的合并(coalesce)操作,将这些已排序的分区重新拼接成一个有序文件。
zorder 排序策略的选项(Options for sort strategy with zorder sort_order)
名称默认值描述
var-length-contribution8对于可变长度类型(String、Binary)的输入列,所考虑的字节数
max-output-size2147483647ZOrder 算法中交错处理的字节数量

输出

输出名称类型描述
rewritten_data_files_countint此命令重写的数据文件数量
added_data_files_countint此命令写入的新数据文件数量
rewritten_bytes_countlong此命令写入的字节数量
failed_data_files_countint当 partial-progress.enabled 为 true 时,重写失败的数据文件数量
removed_delete_files_countint此命令删除的删除文件数量

示例

使用默认的 bin-packing 重写算法重写表 db.sample 中的数据文件,以合并小文件,并根据表的默认写入大小拆分大文件。

CALL catalog_name.system.rewrite_data_files('db.sample');

重写表 db.sample 中的数据文件:按 id 和 name 对所有数据进行排序,并使用与 bin-pack 相同的默认值来确定需要重写的文件。

CALL catalog_name.system.rewrite_data_files(table => 'db.sample', strategy => 'sort', sort_order => 'id DESC NULLS LAST,name ASC NULLS FIRST');

对表 db.sample 中的数据文件按列 c1 和 c2 执行 zOrdering 重写。使用与 bin-pack 相同的默认设置来确定需要重写哪些文件。

CALL catalog_name.system.rewrite_data_files(table => 'db.sample', strategy => 'sort', sort_order => 'zorder(c1,c2)');

在表 db.sample 中,对至少需要重写两个文件的每个分区,使用 bin-pack(装箱)策略重写数据文件,然后移除所有游离的删除文件。

CALL catalog_name.system.rewrite_data_files(table => 'db.sample', options => map('min-input-files', '2', 'remove-dangling-deletes', 'true'));

重写表 db.sample 中的数据文件,并选择可能包含与过滤条件 (id = 3 and name = "foo") 匹配的数据的文件进行重写。

CALL catalog_name.system.rewrite_data_files(table => 'db.sample', where => 'id = 3 and name = "foo"');

rewrite_manifests

重写表的 manifest,以优化扫描计划的制定。

manifest 中的数据文件按照分区规范(partition spec)中的字段进行排序。该过程通过 Spark 作业并行执行。

Info

此过程会使引用受影响表的所有已缓存 Spark 计划失效。

用法

参数名称是否必需类型描述
table✔️string要更新的表的名称
use_caching️boolean操作期间是否使用 Spark 缓存(默认为 false)。启用缓存可能会增加执行器的内存占用。
spec_id️int要重写的 manifest 所属的 spec id(默认为当前 spec id)
sort_by️array用于对 manifest 进行分组聚类的分区转换名称列表。选择常被查询的分区转换,可以通过跳过不必要的 manifest 来缩短计划制定时间。如果未设置,manifest 将按 spec 中所有分区转换的顺序进行排序。

输出

输出名称类型描述
rewritten_manifests_countint此命令重写的 manifest 数量
added_manifests_countint此命令写入的新 manifest 文件数量

示例

重写表 db.sample 中的 manifest,并使 manifest 文件与表的分区方式保持一致。

CALL catalog_name.system.rewrite_manifests('db.sample');

重写表 db.sample 中分区规范 1 上的 manifest 文件。

CALL catalog_name.system.rewrite_manifests(table => 'db.sample', spec_id => 1);

重写表 db.sample 中的清单文件,并按分区字段 category 对清单条目进行聚簇。当查询经常对 category 进行过滤时,这可以提升扫描规划的性能。

CALL catalog_name.system.rewrite_manifests(table => 'db.sample', sort_by => array('category'));

rewrite_position_delete_files

Iceberg 可以重写位置删除文件,这有两个目的:

  • 小规模合并(Minor Compaction):将小的位置删除文件合并为较大的文件。这可以减小 manifest 文件中存储的元数据规模,并降低打开小删除文件的开销。
  • 清除悬空删除(Remove Dangling Deletes):过滤掉那些指向已不再存活的数据文件的位置删除记录。执行 rewrite_data_files 之后,指向已被重写数据文件的位置删除记录并不总会被标记为删除,仍可能被表的存活快照元数据所跟踪。这就是所谓的"悬空删除"问题。

用法

参数名是否必需类型说明
table✔️string要更新的表的名称
options️map该过程使用的选项
where️string用于过滤文件的谓词(以字符串形式表示)。

重写时始终会过滤掉悬空删除。

选项

名称 默认值 说明

max-concurrent-file-group-rewrites 5 同时重写的文件组的最大数量

partial-progress.enabled false 在整个重写完成之前,允许提交部分文件组

partial-progress.max-commits 10 如果启用了部分进度,本次重写允许产生的最大提交次数

rewrite-job-order none 根据该值强制指定重写作业的顺序。

  • 如果 rewrite-job-order=bytes-asc,则先重写最小的作业组。
  • 如果 rewrite-job-order=bytes-desc,则先重写最大的作业组。
  • 如果 rewrite-job-order=files-asc,则先重写文件最少的作业组。
  • 如果 rewrite-job-order=files-desc,则先重写文件最多的作业组。
  • 如果 rewrite-job-order=none,则按计划生成的顺序重写作业组(无特定排序)。

target-file-size-bytes 67108864(64MB,即表属性中 write.delete.target-file-size-bytes 的默认值) 目标输出文件大小

min-file-size-bytes 目标文件大小的 75% 小于该阈值的文件,无论其他条件如何,都会被考虑进行重写

max-file-size-bytes 目标文件大小的 180% 大于该阈值的文件,无论其他条件如何,都会被考虑进行重写

min-input-files 5 文件数量超过该数量的文件组,无论其他条件如何,都会被重写

rewrite-all false 强制重写所提供的全部文件,覆盖其他选项

max-file-group-size-bytes 107374182400(100GB)单个文件组中应当被重写的最大数据量。整个重写操作会先按分区拆分,再在分区内按大小拆分为文件组。这有助于拆分对超大分区的重写工作,否则这些分区可能因集群资源限制而无法重写。

max-files-to-rewrite null 该选项用于设置将要重写的符合条件文件数量的上限。如果未指定该选项,则所有符合条件的文件都会被重写。

输出

输出名称类型描述
rewritten_delete_files_countint由此命令删除的删除文件(delete file)数量
added_delete_files_countint由此命令新增的删除文件数量
rewritten_bytes_countlong由此命令删除的所有删除文件的总字节数
added_bytes_countlong由此命令新增的所有删除文件的总字节数

示例

重写表 db.sample 中的位置删除文件。该操作会选取符合默认重写条件的位置删除文件,并按目标大小 target-file-size-bytes 写入新文件。游离的删除记录会从重写后的删除文件中移除。

CALL catalog_name.system.rewrite_position_delete_files('db.sample');

重写表 db.sample 中的所有 position delete 文件,新文件大小为 target-file-size-bytes。重写后的删除文件中会移除悬空删除(dangling deletes)。

CALL catalog_name.system.rewrite_position_delete_files(table => 'db.sample', options => map('rewrite-all', 'true'));

重写表 db.sample 中的位置删除文件。该操作会基于大小条件,选出位于包含两个或更多位置删除文件的分区中的位置删除文件进行重写。重写后的删除文件中会移除悬空删除(dangling deletes)。

CALL catalog_name.system.rewrite_position_delete_files(table => 'db.sample', options => map('min-input-files','2'));

表迁移

snapshot 和 migrate 过程用于测试并将现有的 Hive 或 Spark 表迁移到 Iceberg。

snapshot

创建表的轻量级临时副本以便测试,且不修改源表。

新创建的表可以进行更改或写入操作,而不会影响源表,但快照使用的是原表的数据文件。

当在快照上执行插入或覆盖操作时,新文件会被写入快照表的位置,而不是原表的位置。

测试完成后,可通过执行 DROP TABLE 清理快照表。

Info

由于 snapshot 创建的表并非其数据文件的唯一所有者,因此禁止它们执行 expire_snapshots 等会物理删除数据文件的操作。仅影响元数据的 Iceberg 删除操作仍然允许。此外,任何影响原始数据文件的操作都会破坏快照的完整性。针对原始 Hive 表执行的 DELETE 语句会删除原始数据文件,snapshot 表将无法再访问这些文件。

参见 migrate,了解如何用 Iceberg 表替换现有表。

用法

参数名是否必需类型说明
source_table✔️string要创建快照的表名
table✔️string要创建的新 Iceberg 表名
locationstring新表的表位置(默认由 catalog 决定)
properties️map要添加到新创建表中的属性
parallelismint用于读取文件的线程数(默认为 1)

输出

输出名类型说明
imported_files_countlong添加到新表中的文件数

示例

创建一个独立的 Iceberg 表,引用表 db.sample,表名为 db.snap,位于 catalog 为 db.snap 默认的位置。

CALL catalog_name.system.snapshot('db.sample', 'db.snap');

将一个孤立的 Iceberg 表迁移到手动指定的位置 /tmp/temptable/,该表名为 db.snap,引用表为 db.sample。

CALL catalog_name.system.snapshot('db.sample', 'db.snap', '/tmp/temptable/');

migrate

将表替换为 Iceberg 表,并装载源表的数据文件。

表的架构、分区、属性和位置都将从源表复制而来。

如果表的任何分区使用了不受支持的格式,Migrate 将会失败。受支持的格式有 Avro、Parquet 和 ORC。如果表使用了分桶,Migrate 也会失败,因为分桶信息不会被保留。已有的数据文件会被添加到 Iceberg 表的元数据中,并可通过根据原始表架构创建的名称到 ID 的映射进行读取。

若想在测试时保留原始表不变,请使用 snapshot 创建一个共享源数据文件和架构的新的临时表。

默认情况下,原始表会以 table_BACKUP_ 为名称保留。

用法

参数名是否必需类型说明
table✔️string要迁移的表的名称
properties️map新 Iceberg 表的属性
drop_backupboolean为 true 时,原始表不会作为备份保留(默认为 false)
backup_table_namestring作为备份保留的表的名称(默认为 table_BACKUP_)
parallelismint用于读取文件的线程数(默认为 1)

输出

输出名类型说明
migrated_files_countlong追加到 Iceberg 表的文件数量

示例

将 Spark 默认 catalog 中的表 db.sample 迁移为 Iceberg 表,并添加一个值为 'bar' 的属性 'foo':

CALL catalog_name.system.migrate('spark_catalog.db.sample', map('foo', 'bar'));

将当前 catalog 中的 db.sample 迁移为 Iceberg 表,且不添加任何额外属性:

CALL catalog_name.system.migrate('db.sample');

add_files

尝试直接将 Hive 表或基于文件的表中的文件添加到指定的 Iceberg 表中。与 migrate 或 snapshot 不同,add_files 可以从特定的一个或多个分区导入文件,并且不会创建新的 Iceberg 表。此命令会为新文件创建元数据,而不会移动这些文件。此过程不会分析文件的 schema,以判断它们是否真的与 Iceberg 表的 schema 匹配。执行完成后,Iceberg 表将把这些文件视为自身所拥有文件集合的一部分。这意味着后续的 expire_snapshot 调用将能够物理删除这些已添加的文件。如果可以使用 migrate 或 snapshot,则不应使用此方法。

警告

请注意,add_files 过程只会获取每个待添加文件的 Parquet 元数据一次。如果你使用的是分层存储(例如 Amazon S3 Intelligent-Tiering 存储类型),底层文件将从归档中检索,并在一段设定的时间内保留在较高层级上。

用法

参数名称必需类型描述
table✔️string将要添加文件的表
source_table✔️string文件来源表,也可以使用 file_format.path 形式的路径
partition_filtermap要从中导入的源表分区的映射
check_duplicate_filesboolean是否阻止添加已存在于表中的文件(默认为 true)
parallelismint用于读取文件的线程数(默认为 1)

警告:不会校验 schema,向 Iceberg 表中添加 schema 不同的文件会导致问题。

警告:此方法添加的文件可能会被 Iceberg 操作物理删除。

输出

输出名称类型描述
added_files_countlong此命令添加的文件数量
changed_partition_countlong此命令变更的分区数量(如果已知)

警告

当表属性 compatibility.snapshot-id-inheritance.enabled 设置为 true,或者表格式版本大于 1 时,changed_partition_count 将为 NULL。

示例

将表 db.src_table 中的文件添加到 Iceberg 表 db.tbl,其中 db.src_table 是会话 Catalog 中注册的 Hive 表或 Spark 表。仅添加 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')
);

将位于 path/to/table 路径下的基于 parquet 文件的表中的文件添加到 Iceberg 表 db.tbl 中。无论文件属于哪个分区,都将其全部添加。

CALL spark_catalog.system.add_files(
  table => 'db.tbl',
  source_table => '`parquet`.`path/to/table`'
);

register_table

为已存在但没有对应目录标识符的 metadata.json 文件创建目录条目。

用法

参数名称是否必需类型说明
table✔️string要注册的表
metadata_file✔️string要作为新目录标识符注册的元数据文件

警告

将同一个 metadata.json 注册到多个目录中,可能会导致更新丢失、数据丢失以及表损坏。仅当该表已不再在现有目录中注册,或你正在将表从一个目录迁移到另一个目录时,才使用此过程。

输出

输出名称类型说明
current_snapshot_idlong新注册的 Iceberg 表的当前快照 ID
total_records_countlong新注册的 Iceberg 表的总记录数
total_data_files_countlong新注册的 Iceberg 表的数据文件总数

示例

将一个新表以 db.tbl 注册到 spark_catalog,并指向元数据文件 path/to/metadata/file.json。

CALL spark_catalog.system.register_table(
  table => 'db.tbl',
  metadata_file => 'path/to/metadata/file.json'
);

元数据信息

ancestors_of

报告指定快照的父级实时快照 ID

用法

参数名必填?类型描述
table✔️string要报告实时快照 ID 的表名
snapshot_id️long使用指定的快照来获取其父级的实时快照 ID

提示:使用 snapshot_id

给定如下快照历史,其中回滚到 B 并新增了 C' -> D'

A -> B - > C -> D
      \ -> C' -> (D')

不指定快照 ID 将返回 A -> B -> C' -> D',而提供 D 的快照 ID 作为参数则返回 A-> B -> C -> D

输出

输出名类型描述
snapshot_idlong祖先快照 ID
timestamplong快照创建时间

示例

获取当前快照(默认)的所有快照祖先

CALL spark_catalog.system.ancestors_of('db.tbl');

获取某个特定快照的所有快照祖先

CALL spark_catalog.system.ancestors_of('db.tbl', 1);
CALL spark_catalog.system.ancestors_of(snapshot_id => 1, table => 'db.tbl');

变更数据捕获

create_changelog_view

创建一个包含指定表变更内容的视图。

用法

参数名称是否必需类型描述
table✔️string变更日志的源表名称
changelog_viewstring要创建的视图名称
optionsmap要使用的 Spark 读取选项映射
net_changesboolean是否输出净变更(详见下文)。默认为 false。当 compute_updates 为 true 时,该参数必须为 false。
compute_updatesboolean是否计算更新前/更新后的镜像(详见下文)。如果提供了 identifer_columns,则默认为 true;否则默认为 false。
identifier_columnsarray用于计算更新的标识列列表。如果参数 compute_updates 被设置为 true 但未提供 identifier_columns,则会使用表当前的标识字段。

以下是常用的 Spark 读取选项:

  • start-snapshot-id:起始快照 ID(不含该快照)。若未提供,则从表的第一个快照开始读取(含该快照)。
  • end-snapshot-id:结束快照 ID(含该快照),默认为表的当前快照。
  • start-timestamp:起始时间戳(不含该时间点)。若未提供,则从表的第一个快照开始读取(含该快照)。
  • end-timestamp:结束时间戳(含该时间点),默认为表的当前快照。

输出

输出名称类型描述
changelog_viewstring所创建的变更日志视图名称

示例

基于快照 1(不含)与快照 2(含)之间发生的变更,创建一个变更日志视图 tbl_changes。

CALL spark_catalog.system.create_changelog_view(
  table => 'db.tbl',
  options => map('start-snapshot-id','1','end-snapshot-id', '2')
);

基于时间戳 1678335750489(不含)到 1678992105265(含)之间发生的变更,创建变更日志视图 my_changelog_view。

CALL spark_catalog.system.create_changelog_view(
  table => 'db.tbl',
  options => map('start-timestamp','1678335750489','end-timestamp', '1678992105265'),
  changelog_view => 'my_changelog_view'
);

创建一个变更日志视图,该视图基于标识符列 id 和 name 计算更新。

CALL spark_catalog.system.create_changelog_view(
  table => 'db.tbl',
  options => map('start-snapshot-id','1','end-snapshot-id', '2'),
  identifier_columns => array('id', 'name')
);

创建变更日志视图后,你可以查询该视图,查看两个快照之间发生的变更。

SELECT * FROM tbl_changes;
SELECT * FROM tbl_changes where _change_type = 'INSERT' AND id = 3 ORDER BY _change_ordinal;

请注意,changelog 视图包含变更数据捕获(CDC)元数据列,这些列提供了有关所跟踪变更的额外信息。这些列包括:

  • _change_type:变更类型。其取值为以下之一:INSERT、DELETE、UPDATE_BEFORE 或 UPDATE_AFTER。
  • _change_ordinal:变更的顺序编号
  • _commit_snapshot_id:发生该变更的快照 ID

以下是相应的结果示例。它显示第一个快照插入了 2 条记录,第二个快照删除了 1 条记录。

idname_change_type_change_ordinal_commit_snapshot_id
1AliceINSERT05390529835796506035
2BobINSERT05390529835796506035
1AliceDELETE18764748981452218370

净变更

该 procedure 可以移除跨多个快照的中间变更,仅输出净变更。以下示例展示了如何创建一个计算净变更的 changelog 视图。

CALL spark_catalog.system.create_changelog_view(
  table => 'db.tbl',
  options => map('end-snapshot-id', '87647489814522183702'),
  net_changes => true
);

经过这些净变更,上面的变更日志视图只剩下面这一行,因为 Alice 在第一个快照中被插入,在第二个快照中被删除。

idname_change_type_change_ordinal_commit_snapshot_id
2BobINSERT05390529835796506035

携带行(Carry-over Rows)

该过程默认会移除携带行。携带行是使用写时复制(copy-on-write)执行行级操作(MERGE、UPDATE 和 DELETE)时产生的结果。例如,给定一个包含 row1 (id=1, name='Alice') 和 row2 (id=2, name='Bob') 的文件。对 row2 执行写时复制删除时,需要擦除该文件,并将 row1 保存到一个新文件中。变更日志表会将其报告为下面这一对行,尽管这并不是对表的实际变更。

idname_change_type
1AliceDELETE
1AliceINSERT

要查看携带行,请按如下方式查询 SparkChangelogTable:

SELECT * FROM spark_catalog.db.tbl.changes;

更新前/更新后镜像

如果配置了相关选项,该过程会计算更新前/更新后镜像。更新前/更新后镜像由一对删除行和插入行转换而来。标识列用于判断插入记录和删除记录是否指向同一行。如果两条记录在标识列上具有相同的值,则认为它们是同一行的更新前状态和更新后状态。你可以在表结构中设置标识字段,也可以将其作为过程参数输入。

下面的示例展示了使用标识列(id)计算更新前/更新后镜像的过程,其中 id 相同的行删除和行插入被视为一次更新操作。具体来说,假设我们有如下一对行:

idname_change_type
3RobertDELETE
3DanINSERT

在这种情况下,该过程会将更新前的行标记为 UPDATE_BEFORE 镜像,将更新后的行标记为 UPDATE_AFTER 镜像,得到如下的更新前/更新后镜像:

idname_change_type
3RobertUPDATE_BEFORE
3DanUPDATE_AFTER

表统计信息

compute_table_stats

该过程用于计算特定表的不同值数量(NDV)统计信息。默认情况下,统计信息基于表的当前快照针对所有列计算。也可以选择配置该过程,针对特定快照和/或列子集计算统计信息。

参数名是否必需类型说明
table✔️string表的名称
snapshot_idstring要收集统计信息的快照 ID
columnsarray要收集统计信息的列

输出

输出名类型说明
statistics_filestring由该命令创建的统计信息文件路径

示例

收集表 my_table 最新快照的统计信息

CALL catalog_name.system.compute_table_stats('my_table');

收集表 my_table 中 ID 为 snap1 的快照的统计信息。

CALL catalog_name.system.compute_table_stats(table => 'my_table', snapshot_id => 'snap1' );

收集表 my_table 中 ID 为 snap1 的快照在列 col1 和 col2 上的统计信息

CALL catalog_name.system.compute_table_stats(table => 'my_table', snapshot_id => 'snap1', columns => array('col1', 'col2'));

分区统计

compute_partition_stats

此过程(procedure)会从最近一个包含 PartitionStatisticsFile 的快照开始,增量计算分区统计信息,一直计算到指定的快照(若未指定则使用当前快照),并将合并后的结果写入 PartitionStatisticsFile。如果之前的分区统计文件不存在,则会执行完整计算。该过程还会将 PartitionStatisticsFile 注册到表元数据中。

参数名是否必需类型说明
table✔️string表的名称
snapshot_idstring用于计算分区统计信息的快照 ID。默认为当前快照 ID

输出

输出名类型说明
partition_statistics_filestring由该命令创建的分区统计文件的路径

示例

收集表 my_table 最新快照的分区统计信息

CALL catalog_name.system.compute_partition_stats('my_table');

收集表 my_table 中 ID 为 snap1 的快照的分区统计信息

CALL catalog_name.system.compute_partition_stats(table => 'my_table', snapshot_id => 'snap1');

表复制

rewrite_table_path 存储过程用于为将 Iceberg 表复制到其他位置做准备。

rewrite_table_path

将 Iceberg 表的元数据文件暂存一份副本,并将其中所有的绝对路径源前缀替换为指定的目标前缀。
这可以作为将 Iceberg 表完整或增量复制到新位置的起点。

信息

该过程仅暂存改写后的元数据文件,并准备待复制文件的列表。实际的文件复制不包含在此过程中。

参数名称是否必需默认值类型描述
table✔️string表的名称
source_prefix✔️string要替换的现有前缀
target_prefix✔️string用于替换 source_prefix 的目标前缀
start_version表元数据日志中的第一个 metadata.jsonstring按时间顺序排在最前面的、要改写的 metadata.json 的名称或路径
end_version表元数据日志中的最新 metadata.jsonstring按时间顺序排在最后面的、要改写的 metadata.json 的名称或路径
staging_location表元数据目录下的新目录string新改写的元数据文件的输出位置
create_file_listtrueboolean是否生成包含已改写元数据路径的文件列表

运作模式

  • 完整改写:完整改写会改写所有可达的元数据文件(包括 metadata.json、manifest 列表、manifest 以及 position delete 文件),并将所有可达文件返回到 file_list_location 中。这是该过程的默认运作模式。
  • 增量改写:可以选择提供 start_version 和 end_version,将范围限定为一次增量改写。增量改写只会改写在 start_version 与 end_version 之间新增的元数据文件,并且仅将该区间内新增的文件返回到 file_list_location 中。

输出

输出名称类型描述
latest_versionstring本过程重写的最新元数据文件的名称
file_list_locationstring包含源路径到目标路径映射的 CSV 文件的路径
rewritten_manifest_file_paths_countint路径被重写的清单文件数量
rewritten_delete_file_paths_countint路径被重写的删除文件数量
文件列表

该文件包含在 start_version 与 end_version 之间添加到表中的所有文件的复制计划。

对于每个文件,它指定了:

  • 源路径:文件在表中的原始文件路径;如果文件已被重写,则为暂存位置
  • 目标路径:带有替换前缀的路径

下面的示例展示了三个文件的复制计划:

sourcepath/datafile1.parquet,targetpath/datafile1.parquet
sourcepath/datafile2.parquet,targetpath/datafile2.parquet
stagingpath/manifest.avro,targetpath/manifest.avro

示例

此示例将 my_table 的元数据路径从 HDFS 中的源位置完全重写为 S3 中的目标位置。它会在该表元数据目录下的默认暂存位置中生成一套新的元数据。

CALL catalog_name.system.rewrite_table_path(
    table => 'db.my_table',
    source_prefix => 'hdfs://nn:8020/path/to/source_table',
    target_prefix => 's3a://bucket/prefix/db.db/my_table'
);

此示例在元数据版本 v2.metadata.json 与 v20.metadata.json 之间增量重写 my_table 的元数据路径,并将新的元数据文件写入显式的暂存位置。

CALL catalog_name.system.rewrite_table_path(
    table => 'db.my_table',
    source_prefix => 's3a://bucketOne/prefix/db.db/my_table',
    target_prefix => 's3a://bucketTwo/prefix/db.db/my_table',
    start_version => 'v2.metadata.json',
    end_version => 'v20.metadata.json',
    staging_location => 's3a://bucketStaging/my_table'
);

重写完成后,第三方工具(例如 Distcp)可以将新创建的元数据文件和数据文件复制到目标位置。

最后,可以使用 register_table 过程将复制到目标位置的表注册到目录(catalog)中。

警告

带有分区统计文件的 Iceberg 表目前不支持路径重写。

评论

登录后参与评论

正在加载评论…