运维 Hudi

SQL 存储过程

师成师成· 更新于 2026-09-29· 阅读 151 分钟· 0 次阅读

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

使用 Hudi SparkSQL 扩展时,在所有 Spark 版本中均可使用存储过程。

用法

CALL 支持按名称传递参数(推荐),也支持按位置传递参数。同时支持混合使用位置参数和命名参数。

命名参数

所有过程参数都有名称。按名称传递参数时,参数可以按任意顺序排列,并且可以省略任意可选参数。

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

位置参数

按位置传递参数时,如果参数是可选的,则可以省略。

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

注意

此处的 system 没有实际意义,存储过程的完整名称为 system.procedure_name。

help

显示存储过程的参数和输出类型。

输入

参数名称类型是否必需默认值描述
cmdString否无存储过程的名称

输出

输出名称类型
resultString

示例

call help(cmd => 'show_commits');

参数:

参数类型名称默认值必填
tablestringNonetrue
limitinteger10false

输出类型:

名称类型名称可为空元数据
commit_timestringtrue
actionstringtrue
total_bytes_writtenlongtrue
total_files_addedlongtrue
total_files_updatedlongtrue
total_partitions_writtenlongtrue
total_records_writtenlongtrue
total_update_records_writtenlongtrue
total_errorslongtrue

提交管理

show_commits

显示提交信息。

输入

参数名称类型必填默认值描述
tableStringYNoneHudi 表名
limitIntN10返回的最大记录数

输出

输出名称类型
commit_timeString
state_transition_timeString
actionString
total_bytes_writtenLong
total_files_addedLong
total_files_updatedLong
total_partitions_writtenLong
total_records_writtenLong
total_update_records_writtenLong
total_errorsLong

示例

call show_commits(table => 'test_hudi_table', limit => 10);
提交时间写入总字节数新增文件总数更新文件总数写入分区总数写入记录总数更新写入记录总数错误总数
20220216171049652432653011000
20220216171027021435346101100
20220216171019361435349101100

show_commits_metadata

显示提交元数据。

输入

参数名称类型是否必需默认值描述
tableString是无Hudi 表名
limitInt否10返回的最大记录数

输出

输出名称类型
commit_timeString
state_transition_timeString
actionString
partitionString
file_idString
previous_commitString
num_writesLong
num_insertsLong
num_deletesLong
num_update_writesLong
total_errorsLong
total_log_blocksLong
total_corrupt_log_blocksLong
total_rollback_blocksLong
total_log_recordsLong
total_updated_records_compactedLong
total_bytes_writtenLong

示例

call show_commits_metadata(table => 'test_hudi_table');
commit_timeactionpartitionfile_idprevious_commitnum_writesnum_insertsnum_deletesnum_update_writestotal_errorstotal_log_blockstotal_corrupt_logblockstotal_rollback_blockstotal_log_recordstotal_updated_records_compactedtotal_bytes_written
20220109225319449commitdt=2021-05-03d0073a12-085d-4f49-83e9-402947e7e90a-0null1100000000435349
20220109225311742commitdt=2021-05-02b3b32bac-8a44-4c4d-b433-0cb1bf620f23-0202201092148305921100000000435340
20220109225301429commitdt=2021-05-010d7298b3-6b55-4cff-8d7d-b0772358b78a-0202201092148305921100000000435340
20220109214830592commitdt=2021-05-010d7298b3-6b55-4cff-8d7d-b0772358b78a-0202201091916310150010000000432653
20220109214830592commitdt=2021-05-02b3b32bac-8a44-4c4d-b433-0cb1bf620f23-0202201091916481810010000000432653
20220109191648181commitdt=2021-05-02b3b32bac-8a44-4c4d-b433-0cb1bf620f23-0null1100000000435341
20220109191631015commitdt=2021-05-010d7298b3-6b55-4cff-8d7d-b0772358b78a-0null1100000000435341

show_commit_extra_metadata

显示提交的额外元数据。

输入

参数名称类型是否必填默认值描述
tableStringYNoneHudi 表名
limitIntN100返回记录的最大数量
instant_timeStringNNoneInstant 时间
metadata_keyStringNNone元数据的键

输出

输出名称类型
instant_timeString
actionString
metadata_keyString
metadata_valueString

示例

call show_commit_extra_metadata(table => 'test_hudi_table');
instant_timeactionmetadata_keymetadata_value
20230206174349556deltacommitschema{"type":"record","name":"hudi_mor_tbl","fields":[{"name":"_hoodie_commit_time","type":["null","string"],"doc":"","default":null},{"name":"_hoodie_commit_seqno","type":["null","string"],"doc":"","default":null},{"name":"_hoodie_record_key","type":["null","string"],"doc":"","default":null},{"name":"_hoodie_partition_path","type":["null","string"],"doc":"","default":null},{"name":"_hoodie_file_name","type":["null","string"],"doc":"","default":null},{"name":"id","type":"int"},{"name":"ts","type":"long"}]}
20230206174349556deltacommitlatest_schema{"max_column_id":8,"version_id":20230206174349556,"type":"record","fields":[{"id":0,"name":"_hoodie_commit_time","optional":true,"type":"string","doc":""},{"id":1,"name":"_hoodie_commit_seqno","optional":true,"type":"string","doc":""},{"id":2,"name":"_hoodie_record_key","optional":true,"type":"string","doc":""},{"id":3,"name":"_hoodie_partition_path","optional":true,"type":"string","doc":""},{"id":4,"name":"_hoodie_file_name","optional":true,"type":"string","doc":""},{"id":5,"name":"id","optional":false,"type":"int"},{"id":8,"name":"ts","optional":false,"type":"long"}]}

show_archived_commits

显示已归档的提交。

输入

参数名类型是否必需默认值描述
tableStringYNoneHudi 表名
limitIntN10返回的最大记录数
start_tsStringN""提交的起始时间,默认:now - 10 天
end_tsStringN""提交的结束时间,默认:now - 1 天

输出

输出名称类型
commit_timeString
state_transition_timeString
total_bytes_writtenLong
total_files_addedLong
total_files_updatedLong
total_partitions_writtenLong
total_records_writtenLong
total_update_records_writtenLong
total_errorsLong

示例

call show_archived_commits(table => 'test_hudi_table');
commit_timetotal_bytes_writtentotal_files_addedtotal_files_updatedtotal_partitions_writtentotal_records_writtentotal_update_records_writtentotal_errors
20220216171049652432653011000
20220216171027021435346101100
20220216171019361435349101100

show_archived_commits_metadata

显示已归档提交的元数据。

输入

参数名类型是否必需默认值说明
tableStringY无Hudi 表名
limitIntN10返回记录的最大数量
start_tsStringN""提交的起始时间,默认:now - 10 天
end_tsStringN""提交的结束时间,默认:now - 1 天

输出

输出名类型
commit_timeString
state_transition_timeString
actionString
partitionString
file_idString
previous_commitString
num_writesLong
num_insertsLong
num_deletesLong
num_update_writesLong
total_errorsLong
total_log_blocksLong
total_corrupt_log_blocksLong
total_rollback_blocksLong
total_log_recordsLong
total_updated_records_compactedLong
total_bytes_writtenLong

示例

call show_archived_commits_metadata(table => 'test_hudi_table');
commit_timeactionpartitionfile_idprevious_commitnum_writesnum_insertsnum_deletesnum_update_writestotal_errorstotal_log_blockstotal_corrupt_logblockstotal_rollback_blockstotal_log_recordstotal_updated_records_compactedtotal_bytes_written
20220109225319449commitdt=2021-05-03d0073a12-085d-4f49-83e9-402947e7e90a-0null1100000000435349
20220109225311742commitdt=2021-05-02b3b32bac-8a44-4c4d-b433-0cb1bf620f23-0202201092148305921100000000435340
20220109225301429commitdt=2021-05-010d7298b3-6b55-4cff-8d7d-b0772358b78a-0202201092148305921100000000435340
20220109214830592commitdt=2021-05-010d7298b3-6b55-4cff-8d7d-b0772358b78a-0202201091916310150010000000432653
20220109214830592commitdt=2021-05-02b3b32bac-8a44-4c4d-b433-0cb1bf620f23-0202201091916481810010000000432653
20220109191648181commitdt=2021-05-02b3b32bac-8a44-4c4d-b433-0cb1bf620f23-0null1100000000435341
20220109191631015commitdt=2021-05-010d7298b3-6b55-4cff-8d7d-b0772358b78a-0null1100000000435341
call show_archived_commits(table => 'test_hudi_table');
提交时间写入总字节数新增文件总数更新文件总数写入分区总数写入记录总数更新写入记录总数错误总数
20220216171049652432653011000
20220216171027021435346101100
20220216171019361435349101100

show_timeline

显示 Hudi 表的时间线条目。返回活动时间线以及可选的已归档时间线中所有时间线操作(提交、压缩、聚类、清理、回滚等)的实例级信息。结果按时间戳降序排列。

输入

参数名类型是否必填默认值描述
tableStringN*无Hudi 表名(与 path 互斥)
pathStringN*无Hudi 表的根路径(与 table 互斥)
limitIntN20返回时间线条目的最大数量(当同时设置了 startTime 和 endTime 时忽略该参数)
showArchivedBooleanNfalse是否包含已归档的时间线条目
filterStringN""用于对任意输出列进行结果过滤的 SQL 表达式
startTimeStringN""用于过滤的起始时间戳(格式:yyyyMMddHHmmss,包含该时刻)
endTimeStringN""用于过滤的结束时间戳(格式:yyyyMMddHHmmss,包含该时刻)

* table 与 path 必须至少提供其中一个。

输出

输出名称类型描述
instant_timeString该 instant 的请求时间戳
actionString操作类型:commit、deltacommit、compaction、clustering、clean、rollback 等
stateStringinstant 的状态:REQUESTED、INFLIGHT 或 COMPLETED
requested_timeString请求该 instant 时的挂钟时间(格式:MM-dd HH:mm:ss)
inflight_timeString该 instant 进入 inflight 状态时的挂钟时间(格式:MM-dd HH:mm:ss)
completed_timeString该 instant 完成时的挂钟时间(格式:MM-dd HH:mm:ss),或 null
timeline_typeStringACTIVE 或 ARCHIVED
rollback_infoString对于回滚 instant:回滚了哪些内容;对于被回滚的 instant:是哪个回滚 instant 将其回滚;其他情况为 null

示例

-- Show the 20 most recent timeline entries
call show_timeline(table => 'test_hudi_table');

-- Show up to 50 entries including archived timeline
call show_timeline(table => 'test_hudi_table', limit => 50, showArchived => true);

-- Filter to completed commits in a time range
call show_timeline(
  table => 'test_hudi_table',
  startTime => '20251201000000',
  endTime => '20251231235959',
  filter => "action = 'commit' AND state = 'COMPLETED'"
);

-- Look up by base path instead of table name
call show_timeline(path => 'hdfs:///user/hive/warehouse/test_hudi_table');
instant_timeactionstaterequested_timeinflight_timecompleted_timetimeline_typerollback_info
20251205143022001commitCOMPLETED12-05 14:30:2012-05 14:30:2112-05 14:30:22ACTIVEnull
20251205141510003cleanCOMPLETED12-05 14:15:0912-05 14:15:1012-05 14:15:10ACTIVEnull
20251205140030002commitCOMPLETED12-05 14:00:2812-05 14:00:2912-05 14:00:30ACTIVEnull

show_commit_files

显示一次提交(commit)所涉及的文件。

输入

参数名类型是否必需默认值说明
tableString是无Hudi 表名
limitInt否10返回记录的最大数量
instant_timeString是无Instant 时间

输出

输出名类型
actionString
partition_pathString
file_idString
previous_commitString
total_records_updatedLong
total_records_writtenLong
total_bytes_writtenLong
total_errorsLong
file_sizeLong

示例

call show_commit_files(table => 'test_hudi_table', instant_time => '20230206174349556');
操作分区路径文件 ID前一次提交更新记录总数写入记录总数写入字节总数错误总数文件大小
deltacommitdt=2021-05-037fb52523-c7f6-41aa-84a6-629041477aeb-0null014347680434768

show_commit_partitions

显示某次提交的分区信息。

输入

参数名类型必填默认值描述
tableStringY无Hudi 表名
limitIntN10返回的最大记录数
instant_timeStringY无实例时间

输出

输出名类型
actionString
partition_pathString
total_files_addedLong
total_files_updatedLong
total_records_insertedLong
total_records_updatedLong
total_bytes_writtenLong
total_errorsLong

示例

call show_commit_partitions(table => 'test_hudi_table', instant_time => '20230206174349556');
actionpartition_pathtotal_files_addedtotal_files_updatedtotal_records_insertedtotal_records_updatedtotal_bytes_writtentotal_errors
deltacommitdt=2021-05-037fb52523-c7f6-41aa-84a6-629041477aeb-00143476800

show_commit_write_stats

显示某次提交的写入统计信息。

输入

参数名类型是否必填默认值描述
tableString是无Hudi 表名
limitInt否10返回的最大记录数
instant_timeString是无实例时间(instant time)

输出

输出名类型
actionString
total_bytes_writtenLong
total_records_writtenLong
avg_record_sizeLong

示例

call show_commit_write_stats(table => 'test_hudi_table', instant_time => '20230206174349556');
操作total_bytes_writtentotal_records_writtenavg_record_size
deltacommit4347681434768

show_rollbacks

显示回滚提交。

输入

参数名类型必填默认值说明
tableString是无Hudi 表名
limitInt否10返回的最大记录数

输出

输出名类型
instantString
rollback_instantString
total_files_deletedInt
time_taken_in_millisLong
total_partitionsInt

示例

call show_rollbacks(table => 'test_hudi_table');
instantrollback_instanttotal_files_deletedtime_taken_in_millistotal_partitions
deltacommit43476814347682

show_rollback_detail

显示回滚提交的详细信息。

输入

参数名类型必填默认值描述
tableString是无Hudi 表名称
limitInt否10返回的最大记录数
instant_timeString是无Instant 时间

输出

输出名类型
instantString
rollback_instantString
partitionString
deleted_fileString
succeededInt

示例

call show_rollback_detail(table => 'test_hudi_table', instant_time => '20230206174349556');
instantrollback_instantpartitiondeleted_filesucceeded
deltacommit43476814347682

commits_compare

将提交与另一个路径进行比较。

输入

参数名称类型是否必填默认值描述
tableString是无Hudi 表名称
pathString是无表的路径

输出

输出名称类型
compare_detailString

示例

call commits_compare(table => 'test_hudi_table', path => 'hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table');
compare_detail
源表 test_hudi_table 领先 0 个提交。需要追赶的提交数 - []

archive_commits

归档提交。

输入

参数名称类型必填默认值说明
tableString否无Hudi 表名
pathString否无表路径
min_commitsInt否20与 hoodie.keep.max.commits 类似,但用于控制活动时间线(active timeline)中保留的最少 instant 数量。
max_commitsInt否30每次写入后,归档服务会将时间线中的旧条目移入归档日志,以使元数据开销保持恒定,即使表规模不断增长。此配置用于控制活动时间线中保留的最大 instant 数量。
retain_commitsInt否10instant 的归档以尽力而为的方式分批进行,以便将更多 instant 打包到单个归档日志中。此配置用于控制该归档批次的大小。
enable_metadataBoolean否true启用内部元数据表

输出

输出名称类型
resultInt

示例

call archive_commits(table => 'test_hudi_table');
result
0

export_instants

将 instants 导出到本地文件夹。

输入

参数名类型是否必填默认值描述
tableString是无Hudi 表名
local_folderString是无本地文件夹
limitInt否-1要导出的 instants 数量
actionsString否clean,commit,deltacommit,rollback,savepoint,restore提交操作类型
descBoolean否false是否降序排列

输出

输出名类型
export_detailString

示例

call export_instants(table => 'test_hudi_table', local_folder => '/tmp/folder');
export_detail
已将 6 个 Instant 导出到 /tmp/folder

rollback_to_instant

将表回滚到某个时间点时处于当前状态的那次提交。

输入

参数名类型是否必需默认值说明
tableString是无Hudi 表名
instant_timeString是无Instant 时间

输出

输出名类型
rollback_resultBoolean

示例

将 test_hudi_table 回滚到某一个 instant

call rollback_to_instant(table => 'test_hudi_table', instant_time => '20220109225319449');
rollback_result
true

create_savepoint

为 Hudi 表创建一个保存点(savepoint)。

输入

参数名类型是否必填默认值说明
tableStringY无Hudi 表名
commit_timeStringY无提交时间
userStringN""用户名
commentsStringN""备注

输出

输出名类型
create_savepoint_resultBoolean

示例

call create_savepoint(table => 'test_hudi_table', commit_time => '20220109225319449');
create_savepoint_result
true

show_savepoints

显示保存点(savepoint)。

输入

参数名类型是否必填默认值描述
tableStringY无Hudi 表名称

输出

输出名类型
savepoint_timeString

示例

call show_savepoints(table => 'test_hudi_table');
savepoint_time
20220109225319449
20220109225311742
20220109225301429

delete_savepoint

删除 Hudi 表中的保存点(savepoint)。

输入

参数名称类型是否必填默认值说明
tableString是无Hudi 表名
instant_timeString是无Instant 时间

输出

输出名称类型
delete_savepoint_resultBoolean

示例

从 test_hudi_table 中删除一个保存点

call delete_savepoint(table => 'test_hudi_table', instant_time => '20220109225319449');
delete_savepoint_result
true

rollback_to_savepoint

将表回滚到某个时刻的当前提交。

输入

参数名称类型是否必填默认值说明
tableStringY无Hudi 表名
instant_timeStringY无Instant 时间

输出

输出名称类型
rollback_savepoint_resultBoolean

示例

将 test_hudi_table 回滚到某个保存点

call rollback_to_savepoint(table => 'test_hudi_table', instant_time => '20220109225319449');
rollback_savepoint_result
true

copy_to_temp_view

将表复制到临时视图。

输入

参数名称类型是否必填默认值说明
tableStringYNoneHudi 表名
query_typeStringN"snapshot"需要读取数据的方式:incremental 模式(读取某个 instantTime 以来的新数据)、read_optimized 模式(基于基础文件获取最新视图)、snapshot 模式(通过合并基础文件和(如有)日志文件获取最新视图)
view_nameStringYNone视图名称
begin_instance_timeStringN""起始实例时间
end_instance_timeStringN""结束实例时间
as_of_instantStringN""截至 instant 时间
replaceBooleanNfalse是否替换已存在的视图
globalBooleanNfalse是否为全局视图

输出

输出名称类型
statusInt

示例

call copy_to_temp_view(table => 'test_hudi_table', view_name => 'copy_view_test_hudi_table');
状态
0

copy_to_table

将表复制到新表。

输入

参数名称类型是否必填默认值描述
tableStringYNoneHudi 表名
query_typeStringN"snapshot"数据读取方式:incremental 模式(读取某个 instantTime 之后的新数据)、read_optimized 模式(基于基础文件获取最新视图),或 snapshot 模式(通过合并基础文件与(若存在的)日志文件获取最新视图)
new_tableStringYNone新表的名称
begin_instance_timeStringN""起始实例时间
end_instance_timeStringN""结束实例时间
as_of_instantStringN""截至某个即时时间
save_modeStringN"overwrite"保存模式
columnsStringN""需要从源表复制到新表的列

输出

输出名称类型
statusInt

示例

call copy_to_table(table => 'test_hudi_table', new_table => 'copy_table_test_hudi_table');
状态
0

元数据表管理

create_metadata_table

创建 Hudi 表的元数据表。

输入

参数名称类型是否必填默认值说明
tableString是无Hudi 表名称

输出

输出名称类型
resultString

示例

call create_metadata_table(table => 'test_hudi_table');
result
Created Metadata Table in hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table/.hoodie/metadata (duration=2.777secs)

init_metadata_table

初始化 Hudi 表的元数据表。

输入

参数名称类型必填默认值描述
tableStringY无Hudi 表名
read_onlyBooleanNfalse是否只读

输出

输出名称类型
resultString

示例

call init_metadata_table(table => 'test_hudi_table');
result
Initialized Metadata Table in hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table/.hoodie/metadata (duration=0.023sec)

delete_metadata_table

删除 Hudi 表的元数据表。

输入

参数名类型必填默认值描述
tableStringYNoneHudi 表名称

输出

输出名类型
resultString

示例

call delete_metadata_table(table => 'test_hudi_table');
结果
已从 hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table/.hoodie/metadata 中移除元数据表

show_metadata_table_partitions

显示 Hudi 表的分区。

输入

参数名类型是否必填默认值描述
tableStringY无Hudi 表名称

输出

输出名类型
partitionString

示例

call show_metadata_table_partitions(table => 'test_hudi_table');
partition
dt=2021-05-01
dt=2021-05-02
dt=2021-05-03

show_metadata_table_files

显示 Hudi 表的文件。

输入

参数名类型是否必需默认值描述
tableString否无Hudi 表名
pathString否无表路径
partitionString否""分区名称
limitInt否100限制返回数量
filterString否""用于过滤结果的高级谓词表达式(例如 file_path LIKE '%.parquet')

note

调用此存储过程时,参数 table 和 path 至少必须指定其中一个。如果同时指定两个参数,则 table 生效。

输出

输出名类型
file_pathString

示例

显示 Hudi 表在某个分区下的文件。

call show_metadata_table_files(table => 'test_hudi_table', partition => 'dt=20230220');
文件路径
.d3cdf6ff-250a-4cee-9af4-ab179fdb9bfb-0_20230220190948086.log.1_0-111-123
d3cdf6ff-250a-4cee-9af4-ab179fdb9bfb-0_0-78-81_20230220190948086.parquet

show_metadata_table_stats

显示 Hudi 表的元数据表统计信息。

输入

参数名称类型是否必填默认值说明
tableString是无Hudi 表名称

输出

输出名称类型
stat_keyString
stat_valueString

示例

call show_metadata_table_stats(table => 'test_hudi_table');
stat_keystat_value
dt=2021-05-03.totalBaseFileSizeInBytes23142

validate_metadata_table_files

校验 Hudi 表的元数据表文件。

输入

参数名类型必填默认值说明
tableStringY无Hudi 表名
verboseBooleanNFalse是否打印所有文件

输出

输出名类型
partitionString
file_nameString
is_present_in_fsBoolean
is_present_in_metadataBoolean
fs_sizeLong
metadata_sizeLong

示例

call validate_metadata_table_files(table => 'test_hudi_table');
partitionfile_nameis_present_in_fsis_present_in_metadatafs_sizemetadata_size
dt=2021-05-03ad1e5a3f-532f-4a13-9f60-223676798bf3-0_0-4-4_00000000000002.parquettruetrue4352343523

表信息

show_table_properties

显示表的 Hudi 属性。

输入

参数名称类型是否必填默认值描述
tableString是无Hudi 表名
pathString否无表路径
limitInt否10返回的最大记录数

输出

输出名称类型
keyString
valueString

示例

call show_table_properties(table => 'test_hudi_table', limit => 10);
键值
hoodie.table.ordering.fieldsts
hoodie.table.partition.fieldsdt

show_fs_path_detail

显示路径的详细信息。

输入

参数名称类型必填默认值描述
pathStringY无Hudi 表名
is_subBooleanNfalse是否列出文件
sortBooleanNtrue按 storage_size 排序
limitIntN100限制数量

输出

输出名称类型
path_numLong
file_numLong
storage_sizeLong
storage_size(unit)String
storage_pathString
space_consumedLong
quotaLong
space_quotaLong

示例

call show_fs_path_detail(path => 'hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table');
path_numfile_numstorage_sizestorage_size(单位)storage_pathspace_consumedquotaspace_quota
225820656121.97MBhdfs://ns1/hive/warehouse/hudi.db/test_hudi_table-16196836-1

stats_file_sizes

展示表的文件大小。

输入

参数名称类型是否必填默认值描述
tableString是无Hudi 表名
partition_pathString否""分区路径
limitInt否10返回记录的最大数量

输出

输出名称类型
commit_timeString
minLong
10thDouble
50thDouble
avgDouble
95thDouble
maxLong
num_filesInt
std_devDouble

示例

call stats_file_sizes(table => 'test_hudi_table');
commit_timemin10th50thavg95thmaxnum_filesstd_dev
20230205134149455435000435000.0435000.0435000.0435000.043500010.0

stats_wa

显示表的写入统计信息与写放大情况。

输入

参数名类型是否必填默认值描述
tableString是无Hudi 表名
limitInt否10返回记录的最大条数

输出

输出名类型
commit_timeString
total_upsertedLong
total_writtenLong
write_amplification_factorString

示例

call stats_wa(table => 'test_hudi_table');
commit_timetotal_upsertedtotal_writtenwrite_amplification_factor
总计000

show_logfile_records

显示表中日志文件的记录。

输入

参数名类型是否必填默认值描述
tableStringN无Hudi 表名
pathStringN无表的路径
log_file_path_patternStringY无日志文件路径匹配模式
mergeBooleanNfalse是否合并结果
limitIntN10返回的最大记录数
filterStringN""过滤表达式

输出

输出名类型
recordsString

示例

call show_logfile_records(table => 'test_hudi_table', log_file_path_pattern => 'hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table/*.log*');
记录
{"_hoodie_commit_time": "20230205133427059", "_hoodie_commit_seqno": "20230205133427059_0_10", "_hoodie_record_key": "1", "_hoodie_partition_path": "", "_hoodie_file_name": "3438e233-7b50-4eff-adbb-70b1cd76f518-0", "id": 1, "name": "a1", "price": 40.0, "ts": 1111}

show_logfile_metadata

展示表中 logfile 的元数据。

输入

参数名类型必填默认值描述
tableStringY无Hudi 表名
log_file_path_patternStringY10logfile 的路径模式
mergeBooleanNfalse是否合并结果
limitIntN10返回记录的最大数量

输出

输出名类型
instant_timeString
record_countInt
block_typeString
header_metadataString
footer_metadataString

示例

call show_logfile_metadata(table => 'hudi_mor_tbl', log_file_path_pattern => 'hdfs://ns1/hive/warehouse/hudi.db/hudi_mor_tbl/*.log*');
instant_timerecord_countblock_typeheader_metadatafooter_metadata
202302051334270591AVRO_DATA_BLOCK{"INSTANT_TIME":"20230205133427059","SCHEMA":"{"type":"record","name":"hudi_mor_tbl_record","namespace":"hoodie.hudi_mor_tbl","fields":[{"name":"_hoodie_commit_time","type":["null","string"],"doc":"","default":null},{"name":"_hoodie_commit_seqno","type":["null","string"],"doc":"","default":null},{"name":"_hoodie_record_key","type":["null","string"],"doc":"","default":null},{"name":"_hoodie_partition_path","type":["null","string"],"doc":"","default":null},{"name":"_hoodie_file_name","type":["null","string"],"doc":"","default":null},{"name":"id","type":"int"},{"name":"name","type":"string"},{"name":"price","type":"double"},{"name":"ts","type":"long"}]}"}

show_invalid_parquet

显示表中无效的 parquet 文件。

输入

参数名称类型是否必填默认值说明
pathStringYNoneHudi 表路径
parallelismIntN100并行度
limitIntN100返回条数限制
needDeleteBooleanNfalse是否需要删除
partitionsStringN""分区
instantsStringN""Instants
filterStringN""过滤表达式

输出

输出名称类型
PathString

示例

call show_invalid_parquet(path => 'hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table');
路径
hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table/7fb52523-c7f6-41aa-84a6-629041477aeb-0_0-92-99_20230205133532199.parquet

show_fsview_all

显示表的文件系统视图。

输入

参数名称类型是否必需默认值说明
tableString是无Hudi 表名称
max_instantString否""最大 instant 时间
include_maxBoolean否false是否包含最大 instant
include_in_flightBoolean否false是否包含进行中的 instant
exclude_compactionBoolean否false是否排除 compaction
limitInt否10返回记录的最大数量
path_regexString否"ALL_PARTITIONS"路径的匹配模式

输出

输出名称类型
partitionString
file_idString
base_instantString
data_fileString
data_file_sizeLong
num_delta_filesLong
total_delta_file_sizeLong
delta_filesString

示例

call show_fsview_all(table => 'test_hudi_table');
分区file_idbase_instantdata_filedata_file_sizenum_delta_filestotal_delta_file_sizedelta_files
dt=2021-05-03d0073a12-085d-4f49-83e9-402947e7e90a-0202201092253194497fb52523-c7f6-41aa-84a6-629041477aeb-0_0-92-99_20220109225319449.parquet53194491213193.7fb52523-c7f6-41aa-84a6-629041477aeb-0_20230205133217210.log.1_0-60-63

show_fsview_latest

显示表的最新文件系统视图。

输入

参数名类型是否必填默认值说明
tableStringY无Hudi 表名
max_instantStringN""最大 instant 时间
include_maxBooleanNfalse是否包含最大 instant
include_in_flightBooleanNfalse是否包含进行中的
exclude_compactionBooleanNfalse是否排除 compaction
path_regexStringN"ALL_PARTITIONS"路径匹配模式
partition_pathStringN"ALL_PARTITIONS"分区路径
mergeBooleanNfalse是否合并结果

输出

输出名类型
partitionString
file_idString
base_instantString
data_fileString
data_file_sizeLong
num_delta_filesLong
total_delta_file_sizeLong
delta_size_compaction_scheduledLong
delta_size_compaction_unscheduledLong
delta_to_base_radio_compaction_scheduledDouble
delta_to_base_radio_compaction_unscheduledDouble
delta_files_compaction_scheduledString
delta_files_compaction_unscheduledString

示例

call show_fsview_latest(table => 'test_hudi_table', partition => 'dt=2021-05-03');
partitionfile_idbase_instantdata_filedata_file_sizenum_delta_filestotal_delta_file_sizedelta_files
dt=2021-05-03d0073a12-085d-4f49-83e9-402947e7e90a-0202201092253194497fb52523-c7f6-41aa-84a6-629041477aeb-0_0-92-99_20220109225319449.parquet53194491213193.7fb52523-c7f6-41aa-84a6-629041477aeb-0_20230205133217210.log.1_0-60-63

表服务

run_clustering

在 hoodie 表上触发聚类。通过使用分区谓词,可以在指定分区上运行聚类任务,也可以指定排序列来对数据进行排序。

note

每次调用都会生成新的聚类 instant,或者执行一些处于 pending 状态的聚类 instant。调用此存储过程时,参数 table 和 path 至少必须指定其中一个;如果同时给出两个参数,则以 table 为准。

输入

参数名称类型是否必需默认值描述
tableString否无要进行聚类的表名
pathString否无要进行聚类的表路径
predicateString否无用于过滤分区的谓词
orderString否无排序列,以 , 分隔
show_involved_partitionBoolean否false在输出中显示涉及的分区
opString否无操作类型,EXECUTE 或 SCHEDULE
order_strategyString否无记录布局优化方式,linear/z-order/hilbert
optionsString否无自定义 Hudi 配置,格式为 "key1=value1,key2=value2`
instantsString否无指定的 instant,以 , 分隔
selected_partitionsString否无要执行聚类的分区,以 , 分隔
partition_regex_patternString否无用于过滤分区的正则表达式(例如 2025.*)
limitInt否无要执行的最大计划数

输出

输出名称类型
timestampString
input_group_sizeInt
stateString
involved_partitionsString

示例

使用表名对 test_hudi_table 进行聚类

call run_clustering(table => 'test_hudi_table');

使用表路径对 test_hudi_table 表执行聚类操作。

call run_clustering(path => '/tmp/hoodie/test_hudi_table');

使用表名、谓词和排序列对 test_hudi_table 进行聚类

call run_clustering(table => 'test_hudi_table', predicate => 'ts <= 20220408L', order => 'ts');

对 test_hudi_table 进行聚类,参数为表名与 show_involved_partition

call run_clustering(table => 'test_hudi_table', show_involved_partition => true);

使用表名和 op 对 test_hudi_table 进行聚类

call run_clustering(table => 'test_hudi_table', op => 'schedule');

Clustering test_hudi_table,表名为 order_strategy

call run_clustering(table => 'test_hudi_table', order_strategy => 'z-order');

使用 table name、op、options 对 test_hudi_table 执行聚类

call run_clustering(table => 'test_hudi_table', op => 'schedule', options => '
hoodie.clustering.plan.strategy.target.file.max.bytes=1024*1024*1024,
hoodie.clustering.plan.strategy.max.bytes.per.group=2*1024*1024*1024');

对 test_hudi_table 表执行聚类,参数为 table name、op、instants。

call run_clustering(table => 'test_hudi_table', op => 'execute', instants => 'ts1,ts2');

使用表名、op、selected_partitions 对 test_hudi_table 进行聚类

call run_clustering(table => 'test_hudi_table', op => 'execute', selected_partitions => 'par1,par2');

使用表名、op 与 limit 对 test_hudi_table 进行聚类

call run_clustering(table => 'test_hudi_table', op => 'execute', limit => 10);

使用表名和分区正则模式对 test_hudi_table 进行聚类

call run_clustering(table => 'test_hudi_table', partition_regex_pattern => '2025.*');

注意

limit 参数仅在 op 为 execute 时有效。

show_clustering

显示 hoodie 表上待处理的聚类任务。

注意

调用此过程时,至少需要指定 table 和 path 两个参数中的一个。如果两个参数都提供了,将以 table 为准。

输入

参数名类型是否必填默认值描述
tableString否无要进行聚类的表名
pathString否无要进行聚类的表路径
limitInt否20返回的最大记录数
show_involved_partitionBoolean否false是否在输出中显示涉及的分区

输出

输出名类型
timestampString
input_group_sizeInt
stateString
involved_partitionsString

示例

根据表名显示待处理的聚类任务

call show_clustering(table => 'test_hudi_table');
timestampgroups
202204081537079282
202204081536369633

按表路径显示待执行的 clustering 任务

call show_clustering(path => '/tmp/hoodie/test_hudi_table');
timestampgroups
202204081537079282
202204081536369633

显示待处理的聚类操作,包含表名和限制数量

call show_clustering(table => 'test_hudi_table', limit => 1);
timestampgroups
202204081537079282

run_compaction

在 Hudi 表上调度或执行 compaction。

note

关于调度 compaction:如果指定了 timestamp,新调度的 compaction 将使用给定的时间戳作为 instant time;否则,compaction 将使用当前系统时间进行调度。

关于执行 compaction:给定的 timestamp 必须是已存在的待处理(pending)compaction instant time,若不存在则会抛出异常。同时,如果指定了 timestamp 且存在待处理的 compaction,所有待处理的 compaction 都将被执行,而不会生成新的 compaction instant。

调用此存储过程时,参数 table 和 path 至少需要指定其中一个。如果两个参数都提供了,则 table 生效。

输入

参数名类型是否必填默认值描述
opString是无操作类型,RUN 或 SCHEDULE
tableString否无要进行 compaction 的表名
pathString否无要进行 compaction 的表路径
timestampLong否无Instant time
optionsString否无以逗号分隔的 Hudi compaction 配置列表,格式为 "config1=value1,config2=value2"
instantsString否无以 , 分隔的指定 instants
limitInt否无最多执行的计划数

输出

输出名类型
timestampString
operation_sizeInt
stateString

示例

使用表名执行 compaction

call run_compaction(op => 'run', table => 'test_hudi_table');

通过表路径运行 compaction

call run_compaction(op => 'run', path => '/tmp/hoodie/test_hudi_table');

使用表路径和时间戳运行压缩(compaction)

call run_compaction(op => 'run', path => '/tmp/hoodie/test_hudi_table', timestamp => '20220408153658568');

运行 compaction 并指定选项

call run_compaction(op => 'run', table => 'test_hudi_table', options => hoodie.compaction.strategy=org.apache.hudi.table.action.compact.strategy.LogFileNumBasedCompactionStrategy,hoodie.compaction.logfile.num.threshold=3);

按表名调度 Compaction

call run_compaction(op => 'schedule', table => 'test_hudi_table');
instant
20220408153650834

根据表路径调度 compaction

call run_compaction(op => 'schedule', path => '/tmp/hoodie/test_hudi_table');
instant
20220408153650834

使用表路径和时间戳调度 compaction

call run_compaction(op => 'schedule', path => '/tmp/hoodie/test_hudi_table', timestamp => '20220408153658568');
instant
20220408153658568

show_compaction

显示 hoodie 表上的所有 compaction,包括进行中和已完成的 compaction,结果按触发时间倒序排列。

note

调用此 procedure 时,参数 table 和 path 至少需要指定其中一个。如果两个参数都给了,将以 table 为准。

输入

参数名称类型必填默认值描述
tableString否无要显示 compaction 的表名
pathString否无要显示 compaction 的表路径
limitInt否20返回的最大记录数

输出

参数名称类型必填默认值描述
timestampString否无Instant 时间
operation_sizeInt否无待压缩的 file slice 数量
stateString否无compaction 的状态

示例

按表名显示 compaction

call show_compaction(table => 'test_hudi_table');
timestampactionsize
20220408153707928compaction10
20220408153636963compaction10

显示带有表路径的 compaction

call show_compaction(path => '/tmp/hoodie/test_hudi_table');
时间戳操作大小
20220408153707928compaction10
20220408153636963compaction10

显示 compaction,并附带表名和数量限制

call show_compaction(table => 'test_hudi_table', limit => 1);
timestampactionsize
20220408153707928compaction10

run_clean

在 hoodie 表上运行清理(cleaner)。

输入

参数名称类型是否必需默认值说明
tableString是无要执行清理(clean)的表名称
schedule_in_lineBoolean否true如果希望调度并立即运行一次 clean,请设置为 "true";如果已经调度了 clean 并希望运行该调度,则设置为 false。
clean_policyString否无org.apache.hudi.common.model.HoodieCleaningPolicy:要使用的清理策略。清理服务会删除较旧的文件切片文件以回收空间。长时间运行的查询计划通常会引用较旧的文件切片,如果这些文件切片在查询有机会执行之前就被清理掉,查询将会失败。因此,最好确保数据的保留时间超过最长的查询执行时间。默认情况下,清理策略根据用户显式设置的以下配置之一来确定(这些配置最多只能设置一个;否则将使用 KEEP_LATEST_COMMITS 清理策略)。KEEP_LATEST_FILE_VERSIONS:保留最近写入的 N 个文件切片版本;仅在显式设置 "hoodie.cleaner.fileversions.retained" 时使用。KEEP_LATEST_COMMITS(默认):保留最近 N 次提交写入的文件切片;仅在显式设置 "hoodie.cleaner.commits.retained" 时使用。KEEP_LATEST_BY_HOURS:根据提交时间保留最近 N 小时内写入的文件切片;仅在显式设置 "hoodie.cleaner.hours.retained" 时使用。
retain_commitsInt否无使用 KEEP_LATEST_COMMITS 清理策略时,需要保留而不被清理的提交次数。数据将保留 num_of_commits * time_between_commits(已调度)这么长的时间。这也直接决定了该表在增量查询方面所能支持的数据保留时长。
hours_retainedInt否无使用 KEEP_LATEST_BY_HOURS 清理策略时,需要保留提交的小时数。与按提交次数保留相比,该配置提供了更灵活的选项。设置此属性后,提交时间早于所配置保留小时数的提交所对应的文件组中,除最新文件外的所有文件都将被清理。
file_versions_retainedInt否无使用 KEEP_LATEST_FILE_VERSIONS 清理策略时,清理过程中每个文件组内需保留的最少文件切片数量。
trigger_strategyString否无org.apache.hudi.table.action.clean.CleaningTriggerStrategy:控制何时调度清理。NUM_COMMITS(默认):每 N 次提交触发一次清理服务,N 由 hoodie.clean.max.commits 决定。
trigger_max_commitsInt否无上一次清理操作之后,需要经过多少次提交才会尝试调度新的清理。
optionsString否无以逗号分隔的 Hudi 清理相关配置列表,格式为 "config1=value1,config2=value2"。

输出

参数名称类型
start_clean_timeString
time_taken_in_millisLong
total_files_deletedInt
earliest_commit_to_retainString
bootstrap_part_metadataString
versionInt

示例

使用表名运行 clean 操作

call run_clean(table => 'test_hudi_table');

运行 clean,使用保留最新文件版本策略

call run_clean(table => 'test_hudi_table', trigger_max_commits => 2, clean_policy => 'KEEP_LATEST_FILE_VERSIONS', file_versions_retained => 1)

show_cleans

显示已完成的清理操作,包含耗时、删除文件数以及保留策略等元数据。

输入

参数名称类型是否必需默认值描述
tableStringY无Hudi 表名
limitIntN10返回的最大记录数
showArchivedBooleanNfalse是否包含已归档的时间线数据
filterStringN""用于过滤结果的高级谓词表达式(例如 total_files_deleted > 10)

输出

输出名称类型
clean_timeString
state_transition_timeString
actionString
start_clean_timeString
time_taken_in_millisLong
total_files_deletedInt
earliest_commit_to_retainString
last_completed_commit_timestampString
versionInt

示例

显示指定表名的已完成清理操作

call show_cleans(table => 'test_hudi_table');

显示包含已归档数据的已完成清理操作

call show_cleans(table => 'test_hudi_table', showArchived => true);

显示已完成的清理操作(带过滤条件)

call show_cleans(table => 'test_hudi_table', filter => "total_files_deleted > 10");

show_clean_plans

显示所有状态(REQUESTED、INFLIGHT、COMPLETED)下的清理操作及其状态信息。

输入

参数名类型是否必填默认值描述
tableStringY无Hudi 表名
limitIntN10返回的最大记录数
showArchivedBooleanNfalse是否包含已归档的时间线数据
filterStringN""用于过滤结果的高级谓词表达式(例如 state = 'COMPLETED')

输出

输出名类型
plan_timeString
stateString
actionString
earliest_instant_to_retainString
last_completed_commit_timestampString
policyString
versionInt
total_partitions_to_cleanInt
total_partitions_to_deleteInt
extra_metadataString

示例

按表名显示清理计划

call show_clean_plans(table => 'test_hudi_table');

显示包含已归档数据的清理计划

call show_clean_plans(table => 'test_hudi_table', showArchived => true);

显示带过滤条件的清理计划

call show_clean_plans(table => 'test_hudi_table', filter => "state = 'COMPLETED'");

show_cleans_metadata

显示分区级别的清理详情,用于调试目的。

输入

参数名类型必填默认值说明
tableStringYNoneHudi 表名
limitIntN10返回记录的最大数量
showArchivedBooleanNfalse是否包含已归档的时间线数据
filterStringN""用于过滤结果的高级谓词表达式(例如 partition_path LIKE '2025%')

输出

输出名类型
clean_timeString
state_transition_timeString
actionString
start_clean_timeString
partition_pathString
policyString
delete_path_patternsInt
success_delete_filesInt
failed_delete_filesInt
is_partition_deletedBoolean
time_taken_in_millisLong
total_files_deletedInt

示例

显示指定表名的清理元数据

call show_cleans_metadata(table => 'test_hudi_table');

显示已归档数据的清理元数据

call show_cleans_metadata(table => 'test_hudi_table', showArchived => true);

显示带过滤条件的清理元数据

call show_cleans_metadata(table => 'test_hudi_table', filter => "partition_path LIKE '2025%' AND success_delete_files > 0");

show_file_status

查看指定文件的状态,例如是否已被删除、由哪个操作删除等。

输入

参数名类型是否必需默认值描述
tableString否无Hudi 表名
pathString否无表路径
partitionString否无分区路径(分区表必须指定)
fileString是无要检查的文件名
filterString否""用于过滤结果的高级谓词表达式

:::note
调用此存储过程时,table 和 path 参数至少需要指定其中一个。如果两个参数都指定,则以 table 为准。对于分区表,必须指定 partition 参数。
:::

输出

输出名类型
statusString
actionString
instantString
timelineString
full_pathString

示例

通过表名查看文件状态

call show_file_status(table => 'test_hudi_table', partition => 'dt=2021-05-03', file => 'd0073a12-085d-4f49-83e9-402947e7e90a-0_0-2-2_00000000000002.parquet');

显示指定表路径的文件状态

call show_file_status(path => '/tmp/hoodie/test_hudi_table', partition => 'dt=2021-05-03', file => 'd0073a12-085d-4f49-83e9-402947e7e90a-0_0-2-2_00000000000002.parquet');
状态操作时间戳时间线完整路径
已删除清理20220109225319449活跃

run_ttl

执行 TTL(Time To Live,生存时间)操作,根据保留策略删除分区。

输入

参数名称类型是否必填默认值描述
tableString是无Hudi 表名
ttl_policyString否无TTL 策略类型
retain_daysInt否无分区保留天数
optionsString否无TTL 相关的 Hudi 配置,以逗号分隔,格式为 "config1=value1,config2=value2"

输出

输出名称类型
deleted_partitionsString

示例

使用表名运行 TTL

call run_ttl(table => 'test_hudi_table');

以保留天数运行 TTL

call run_ttl(table => 'test_hudi_table', retain_days => 30);

使用策略和保留天数运行 TTL

call run_ttl(table => 'test_hudi_table', ttl_policy => 'PARTITION_LEVEL', retain_days => 30);
deleted_partitions
dt=2021-05-01
dt=2021-05-02

delete_marker

删除 Hudi 表的 marker 文件。

输入

参数名称类型是否必填默认值描述
tableStringYNoneHudi 表名
instant_timeStringYNoneInstant 名称

输出

输出名称类型
delete_marker_resultBoolean

示例

call delete_marker(table => 'test_hudi_table', instant_time => '20230206174349556');
delete_marker_result
true

sync_validate

验证同步过程。

输入

参数名称类型必填默认值说明
src_tableString是无源表名称
dst_tableString是无目标表名称
modeString是"complete"模式
hive_server_urlString是无Hive server 地址
hive_passString是无Hive 密码
src_dbString否"rawdata"源数据库
target_dbString否"dwh_hoodie"目标数据库
partition_cntInt否5分区数量
hive_userString否""Hive 用户名

输出

输出名称类型
resultString

示例

 call sync_validate(hive_server_url=>'jdbc:hive2://localhost:10000/default', src_table => 'test_hudi_table_src', dst_table=> 'test_hudi_table_dst', mode=>'complete', hive_pass=>'', src_db=> 'default', target_db=>'default');

hive_sync

将表的最新结构同步到 Hive Metastore。

输入

参数名称类型是否必填默认值描述
tableString是无Hudi 表名
metastore_uriString否""Metastore_uri
usernameString否""用户名
passwordString否""密码
use_jdbcString否""启用 Hive 同步时是否使用 JDBC
modeString否""Hive 操作的模式。可选值为 hms、jdbc 和 hiveql。
partition_fieldsString否""表中用于确定 Hive 分区列的字段。
partition_extractor_classString否""实现 PartitionValueExtractor 以提取分区值的类,默认为 'org.apache.hudi.hive.MultiPartKeysValueExtractor'。
strategyString否""Hive 表同步策略。可选值:RO、RT、ALL。
sync_incrementalString否""是否将分区增量同步到 Metastore,即根据提交元数据仅同步新增、变更和删除的分区。若设置为 false,当分区丢失时,元数据同步将执行全量分区同步操作。

输出

输出名称类型
resultString

示例

call hive_sync(table => 'test_hudi_table');
result
true

hdfs_parquet_import

向 Hudi 表中添加 Parquet 文件。

输入

参数名类型必填默认值描述
tableString是无Hudi 表名
table_typeString是""表类型,MERGE_ON_READ 或 COPY_ON_WRITE
src_pathString是""源路径
target_pathString是""目标路径
row_keyString是""主键
partition_keyString是""分区键
schema_file_pathString是""Schema 文件路径
formatString否"parquet"文件格式
commandString否"insert"导入命令
retryInt否0重试次数
parallelismInt否无并行度
props_file_pathString否""属性文件路径

输出

输出名类型
import_resultInt

示例

call hdfs_parquet_import(table => 'test_hudi_table', table_type => 'COPY_ON_WRITE', src_path => '', target_path => '', row_key => 'id', partition_key => 'dt', schema_file_path => '');
import_result
0

repair_add_partition_meta

修复 Hudi 表的分区添加操作。

输入

参数名称类型是否必填默认值描述
tableStringYNoneHudi 表名称
dry_runBooleanNtrue试运行

输出

输出名称类型
partition_pathString
metadata_is_presentString
actionString

示例

call repair_add_partition_meta(table => 'test_hudi_table');
partition_pathmetadata_is_presentaction
dt=2021-05-03是无

repair_corrupted_clean_files

修复 Hudi 表中损坏的 clean 文件。

输入

参数名称类型必填默认值描述
tableString是无Hudi 表名称

输出

输出名称类型
resultBoolean

示例

call repair_corrupted_clean_files(table => 'test_hudi_table');
result
true

repair_deduplicate

修复 Hudi 表的重复记录。该作业对 duplicated_partition_path 中的数据进行去重,并将结果写入 repaired_output_path。作业结束时,repaired_output_path 中的数据会被复制回原始路径(duplicated_partition_path)。

输入

参数名称类型是否必填默认值描述
tableStringY无Hudi 表名
duplicated_partition_pathStringY无重复的分区路径
repaired_output_pathStringY无修复后的输出路径
dry_runBooleanNtrue试运行
dedupe_typeStringN"insert_type"去重类型

输出

输出名称类型
resultString

示例

call repair_deduplicate(table => 'test_hudi_table', duplicated_partition_path => 'dt=2021-05-03', repaired_output_path => '/tmp/repair_path/');
result
重复文件已放置于:/tmp/repair_path/。

repair_migrate_partition_meta

降级 Hudi 表。

输入参数

参数名类型必填默认值说明
tableString是无Hudi 表名
dry_runBoolean否true试运行

输出

输出名类型
partition_pathString
text_metafile_presentString
base_metafile_presentString
actionString

示例

call repair_migrate_partition_meta(table => 'test_hudi_table');

repair_overwrite_hoodie_props

覆盖 Hudi 表属性。

输入

参数名称类型是否必需默认值说明
tableStringY无Hudi 表名
new_props_file_pathStringY无新属性文件的路径

输出

输出名称类型
propertyString
old_valueString
new_valueString

示例

call repair_overwrite_hoodie_props(table => 'test_hudi_table', new_props_file_path = > '/tmp/props');
属性old_valuenew_value
hoodie.file.index.enabletruefalse

引导(Bootstrap)

run_bootstrap

将现有表转换为 Hudi 表。

输入

参数名称类型是否必填默认值描述
tableStringYNone要进行聚类的表名
table_typeStringYNone表类型,MERGE_ON_READ 或 COPY_ON_WRITE
bootstrap_pathStringYNone需要作为 Hudi 表进行引导(bootstrap)的数据集基础路径
base_pathStringYNone基础路径
rowKey_fieldStringYNone主键字段
base_file_formatStringN"PARQUET"基础文件格式
partition_path_fieldStringN""分区列字段
bootstrap_index_classStringN"org.apache.hudi.common.bootstrap.index.HFileBootstrapIndex"用于将骨架基础文件映射到引导基础文件的实现类。
selector_classStringN"org.apache.hudi.client.bootstrap.selector.MetadataOnlyBootstrapModeSelector"选择引导数据集中每个文件/分区的引导模式
key_generator_classStringN"org.apache.hudi.keygen.SimpleKeyGenerator"键生成器类
full_bootstrap_input_providerStringN"org.apache.hudi.bootstrap.SparkParquetBootstrapDataProvider"完整引导输入提供者类
schema_provider_classStringN""Schema 提供者类
payload_classStringN"org.apache.hudi.common.model.OverwriteWithLatestAvroPayload"Payload 类
parallelismIntN1500对于仅元数据的引导,Hudi 会并行化该操作,使每个表分区由一个 Spark 任务处理。此配置限制并行度。如果表分区数大于该配置值,则使用配置的并行度;如果表分区数小于该配置值,则将并行度设为表分区数。对于全记录引导,即记录的 BULK_INSERT 操作,该配置值将作为 BULK_INSERT shuffle 并行度(hoodie.bulkinsert.shuffle.parallelism)传入,决定 BULK_INSERT 的写入行为。如果由于并行度受限导致引导较慢,可以调大此值。
enable_hive_syncBooleanNfalse是否启用 Hive 同步
props_file_pathStringN""属性文件路径
bootstrap_overwriteBooleanNfalse覆盖引导路径

输出

输出名称类型
statusInt

示例

call run_bootstrap(table => 'test_hudi_table', table_type => 'COPY_ON_WRITE', bootstrap_path => 'hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table', base_path => 'hdfs://ns1//tmp/hoodie/test_hudi_table', rowKey_field => 'id', partition_path_field => 'dt',bootstrap_overwrite => true);
status
0

show_bootstrap_mapping

显示 bootstrap 表的映射文件。

输入

参数名类型是否必填默认值描述
tableString否无Hudi 表名
pathString否无表的路径
partition_pathString否""分区路径
file_idsString否""文件 id
limitInt否10返回记录的最大数量
sort_byString否"partition"排序列
descBoolean否false是否降序排列
filterString否""用于过滤结果的高级谓词表达式(例如 partition LIKE '2025%')

note

调用此过程时,参数 table 与 path 至少需要指定其中一个。如果两个参数都指定,则 table 生效。

输出

输出名类型
partitionString
file_idString
source_base_pathString
source_partitionString
source_fileString

示例

根据表名显示 bootstrap 映射

call show_bootstrap_mapping(table => 'test_hudi_table');

显示带有表路径的 bootstrap 映射

call show_bootstrap_mapping(path => '/tmp/hoodie/test_hudi_table');

显示带过滤条件的 bootstrap 映射

call show_bootstrap_mapping(table => 'test_hudi_table', filter => "partition LIKE '2025%' AND file_id > '20251006'");
partitionfile_idsource_base_pathsource_partitionsource_file
dt=2021-05-03d0073a12-085d-4f49-83e9-402947e7e90a-0hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table/dt=2021-05-03/d0073a12-085d-4f49-83e9-402947e7e90a-0_0-2-2_00000000000002.parquetdt=2021-05-03hdfs://ns1/tmp/dt=2021-05-03/00001.parquet

show_bootstrap_partitions

显示 bootstrap 表的分区。

输入

参数名类型是否必填默认值描述
tableString是无要进行聚簇的表的名称

输出

参数名类型
indexed_partitionsString

示例

call show_bootstrap_partitions(table => 'test_hudi_table');
indexed_partitions
dt=2021-05-03

版本管理

upgrade_table

将 Hudi 表升级到指定版本。

输入

参数名类型是否必填默认值描述
tableStringY无Hudi 表名
to_versionStringY无hoodie 表的版本号

输出

输出名类型
resultBoolean

示例

call upgrade_table(table => 'test_hudi_table', to_version => 'FIVE');
result
true

downgrade_table

将 Hudi 表降级到指定版本。

输入

参数名称类型是否必需默认值描述
tableStringYNoneHudi 表名称
to_versionStringYNonehoodie 表的目标版本

输出

输出名称类型
resultBoolean

示例

call downgrade_table(table => 'test_hudi_table', to_version => 'FOUR');
结果
true

评论

登录后参与评论

正在加载评论…