SQL 存储过程
使用 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
显示存储过程的参数和输出类型。
输入
| 参数名称 | 类型 | 是否必需 | 默认值 | 描述 |
|---|---|---|---|---|
| cmd | String | 否 | 无 | 存储过程的名称 |
输出
| 输出名称 | 类型 |
|---|---|
| result | String |
示例
call help(cmd => 'show_commits');参数:
| 参数 | 类型名称 | 默认值 | 必填 |
|---|---|---|---|
| table | string | None | true |
| limit | integer | 10 | false |
输出类型:
| 名称 | 类型名称 | 可为空 | 元数据 |
|---|---|---|---|
| commit_time | string | true | |
| action | string | true | |
| total_bytes_written | long | true | |
| total_files_added | long | true | |
| total_files_updated | long | true | |
| total_partitions_written | long | true | |
| total_records_written | long | true | |
| total_update_records_written | long | true | |
| total_errors | long | true |
提交管理
show_commits
显示提交信息。
输入
| 参数名称 | 类型 | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | None | Hudi 表名 |
| limit | Int | N | 10 | 返回的最大记录数 |
输出
| 输出名称 | 类型 |
|---|---|
| commit_time | String |
| state_transition_time | String |
| action | String |
| total_bytes_written | Long |
| total_files_added | Long |
| total_files_updated | Long |
| total_partitions_written | Long |
| total_records_written | Long |
| total_update_records_written | Long |
| total_errors | Long |
示例
call show_commits(table => 'test_hudi_table', limit => 10);| 提交时间 | 写入总字节数 | 新增文件总数 | 更新文件总数 | 写入分区总数 | 写入记录总数 | 更新写入记录总数 | 错误总数 |
|---|---|---|---|---|---|---|---|
| 20220216171049652 | 432653 | 0 | 1 | 1 | 0 | 0 | 0 |
| 20220216171027021 | 435346 | 1 | 0 | 1 | 1 | 0 | 0 |
| 20220216171019361 | 435349 | 1 | 0 | 1 | 1 | 0 | 0 |
show_commits_metadata
显示提交元数据。
输入
| 参数名称 | 类型 | 是否必需 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| limit | Int | 否 | 10 | 返回的最大记录数 |
输出
| 输出名称 | 类型 |
|---|---|
| commit_time | String |
| state_transition_time | String |
| action | String |
| partition | String |
| file_id | String |
| previous_commit | String |
| num_writes | Long |
| num_inserts | Long |
| num_deletes | Long |
| num_update_writes | Long |
| total_errors | Long |
| total_log_blocks | Long |
| total_corrupt_log_blocks | Long |
| total_rollback_blocks | Long |
| total_log_records | Long |
| total_updated_records_compacted | Long |
| total_bytes_written | Long |
示例
call show_commits_metadata(table => 'test_hudi_table');| commit_time | action | partition | file_id | previous_commit | num_writes | num_inserts | num_deletes | num_update_writes | total_errors | total_log_blocks | total_corrupt_logblocks | total_rollback_blocks | total_log_records | total_updated_records_compacted | total_bytes_written |
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| 20220109225319449 | commit | dt=2021-05-03 | d0073a12-085d-4f49-83e9-402947e7e90a-0 | null | 1 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 435349 |
| 20220109225311742 | commit | dt=2021-05-02 | b3b32bac-8a44-4c4d-b433-0cb1bf620f23-0 | 20220109214830592 | 1 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 435340 |
| 20220109225301429 | commit | dt=2021-05-01 | 0d7298b3-6b55-4cff-8d7d-b0772358b78a-0 | 20220109214830592 | 1 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 435340 |
| 20220109214830592 | commit | dt=2021-05-01 | 0d7298b3-6b55-4cff-8d7d-b0772358b78a-0 | 20220109191631015 | 0 | 0 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 432653 |
| 20220109214830592 | commit | dt=2021-05-02 | b3b32bac-8a44-4c4d-b433-0cb1bf620f23-0 | 20220109191648181 | 0 | 0 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 432653 |
| 20220109191648181 | commit | dt=2021-05-02 | b3b32bac-8a44-4c4d-b433-0cb1bf620f23-0 | null | 1 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 435341 |
| 20220109191631015 | commit | dt=2021-05-01 | 0d7298b3-6b55-4cff-8d7d-b0772358b78a-0 | null | 1 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 435341 |
show_commit_extra_metadata
显示提交的额外元数据。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | None | Hudi 表名 |
| limit | Int | N | 100 | 返回记录的最大数量 |
| instant_time | String | N | None | Instant 时间 |
| metadata_key | String | N | None | 元数据的键 |
输出
| 输出名称 | 类型 |
|---|---|
| instant_time | String |
| action | String |
| metadata_key | String |
| metadata_value | String |
示例
call show_commit_extra_metadata(table => 'test_hudi_table');| instant_time | action | metadata_key | metadata_value |
|---|---|---|---|
| 20230206174349556 | deltacommit | schema | {"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"}]} |
| 20230206174349556 | deltacommit | latest_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
显示已归档的提交。
输入
| 参数名 | 类型 | 是否必需 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | None | Hudi 表名 |
| limit | Int | N | 10 | 返回的最大记录数 |
| start_ts | String | N | "" | 提交的起始时间,默认:now - 10 天 |
| end_ts | String | N | "" | 提交的结束时间,默认:now - 1 天 |
输出
| 输出名称 | 类型 |
|---|---|
| commit_time | String |
| state_transition_time | String |
| total_bytes_written | Long |
| total_files_added | Long |
| total_files_updated | Long |
| total_partitions_written | Long |
| total_records_written | Long |
| total_update_records_written | Long |
| total_errors | Long |
示例
call show_archived_commits(table => 'test_hudi_table');| commit_time | total_bytes_written | total_files_added | total_files_updated | total_partitions_written | total_records_written | total_update_records_written | total_errors |
|---|---|---|---|---|---|---|---|
| 20220216171049652 | 432653 | 0 | 1 | 1 | 0 | 0 | 0 |
| 20220216171027021 | 435346 | 1 | 0 | 1 | 1 | 0 | 0 |
| 20220216171019361 | 435349 | 1 | 0 | 1 | 1 | 0 | 0 |
show_archived_commits_metadata
显示已归档提交的元数据。
输入
| 参数名 | 类型 | 是否必需 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| limit | Int | N | 10 | 返回记录的最大数量 |
| start_ts | String | N | "" | 提交的起始时间,默认:now - 10 天 |
| end_ts | String | N | "" | 提交的结束时间,默认:now - 1 天 |
输出
| 输出名 | 类型 |
|---|---|
| commit_time | String |
| state_transition_time | String |
| action | String |
| partition | String |
| file_id | String |
| previous_commit | String |
| num_writes | Long |
| num_inserts | Long |
| num_deletes | Long |
| num_update_writes | Long |
| total_errors | Long |
| total_log_blocks | Long |
| total_corrupt_log_blocks | Long |
| total_rollback_blocks | Long |
| total_log_records | Long |
| total_updated_records_compacted | Long |
| total_bytes_written | Long |
示例
call show_archived_commits_metadata(table => 'test_hudi_table');| commit_time | action | partition | file_id | previous_commit | num_writes | num_inserts | num_deletes | num_update_writes | total_errors | total_log_blocks | total_corrupt_logblocks | total_rollback_blocks | total_log_records | total_updated_records_compacted | total_bytes_written |
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| 20220109225319449 | commit | dt=2021-05-03 | d0073a12-085d-4f49-83e9-402947e7e90a-0 | null | 1 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 435349 |
| 20220109225311742 | commit | dt=2021-05-02 | b3b32bac-8a44-4c4d-b433-0cb1bf620f23-0 | 20220109214830592 | 1 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 435340 |
| 20220109225301429 | commit | dt=2021-05-01 | 0d7298b3-6b55-4cff-8d7d-b0772358b78a-0 | 20220109214830592 | 1 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 435340 |
| 20220109214830592 | commit | dt=2021-05-01 | 0d7298b3-6b55-4cff-8d7d-b0772358b78a-0 | 20220109191631015 | 0 | 0 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 432653 |
| 20220109214830592 | commit | dt=2021-05-02 | b3b32bac-8a44-4c4d-b433-0cb1bf620f23-0 | 20220109191648181 | 0 | 0 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 432653 |
| 20220109191648181 | commit | dt=2021-05-02 | b3b32bac-8a44-4c4d-b433-0cb1bf620f23-0 | null | 1 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 435341 |
| 20220109191631015 | commit | dt=2021-05-01 | 0d7298b3-6b55-4cff-8d7d-b0772358b78a-0 | null | 1 | 1 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 435341 |
call show_archived_commits(table => 'test_hudi_table');| 提交时间 | 写入总字节数 | 新增文件总数 | 更新文件总数 | 写入分区总数 | 写入记录总数 | 更新写入记录总数 | 错误总数 |
|---|---|---|---|---|---|---|---|
| 20220216171049652 | 432653 | 0 | 1 | 1 | 0 | 0 | 0 |
| 20220216171027021 | 435346 | 1 | 0 | 1 | 1 | 0 | 0 |
| 20220216171019361 | 435349 | 1 | 0 | 1 | 1 | 0 | 0 |
show_timeline
显示 Hudi 表的时间线条目。返回活动时间线以及可选的已归档时间线中所有时间线操作(提交、压缩、聚类、清理、回滚等)的实例级信息。结果按时间戳降序排列。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | N* | 无 | Hudi 表名(与 path 互斥) |
| path | String | N* | 无 | Hudi 表的根路径(与 table 互斥) |
| limit | Int | N | 20 | 返回时间线条目的最大数量(当同时设置了 startTime 和 endTime 时忽略该参数) |
| showArchived | Boolean | N | false | 是否包含已归档的时间线条目 |
| filter | String | N | "" | 用于对任意输出列进行结果过滤的 SQL 表达式 |
| startTime | String | N | "" | 用于过滤的起始时间戳(格式:yyyyMMddHHmmss,包含该时刻) |
| endTime | String | N | "" | 用于过滤的结束时间戳(格式:yyyyMMddHHmmss,包含该时刻) |
* table 与 path 必须至少提供其中一个。
输出
| 输出名称 | 类型 | 描述 |
|---|---|---|
| instant_time | String | 该 instant 的请求时间戳 |
| action | String | 操作类型:commit、deltacommit、compaction、clustering、clean、rollback 等 |
| state | String | instant 的状态:REQUESTED、INFLIGHT 或 COMPLETED |
| requested_time | String | 请求该 instant 时的挂钟时间(格式:MM-dd HH:mm:ss) |
| inflight_time | String | 该 instant 进入 inflight 状态时的挂钟时间(格式:MM-dd HH:mm:ss) |
| completed_time | String | 该 instant 完成时的挂钟时间(格式:MM-dd HH:mm:ss),或 null |
| timeline_type | String | ACTIVE 或 ARCHIVED |
| rollback_info | String | 对于回滚 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_time | action | state | requested_time | inflight_time | completed_time | timeline_type | rollback_info |
|---|---|---|---|---|---|---|---|
| 20251205143022001 | commit | COMPLETED | 12-05 14:30:20 | 12-05 14:30:21 | 12-05 14:30:22 | ACTIVE | null |
| 20251205141510003 | clean | COMPLETED | 12-05 14:15:09 | 12-05 14:15:10 | 12-05 14:15:10 | ACTIVE | null |
| 20251205140030002 | commit | COMPLETED | 12-05 14:00:28 | 12-05 14:00:29 | 12-05 14:00:30 | ACTIVE | null |
show_commit_files
显示一次提交(commit)所涉及的文件。
输入
| 参数名 | 类型 | 是否必需 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| limit | Int | 否 | 10 | 返回记录的最大数量 |
| instant_time | String | 是 | 无 | Instant 时间 |
输出
| 输出名 | 类型 |
|---|---|
| action | String |
| partition_path | String |
| file_id | String |
| previous_commit | String |
| total_records_updated | Long |
| total_records_written | Long |
| total_bytes_written | Long |
| total_errors | Long |
| file_size | Long |
示例
call show_commit_files(table => 'test_hudi_table', instant_time => '20230206174349556');| 操作 | 分区路径 | 文件 ID | 前一次提交 | 更新记录总数 | 写入记录总数 | 写入字节总数 | 错误总数 | 文件大小 |
|---|---|---|---|---|---|---|---|---|
| deltacommit | dt=2021-05-03 | 7fb52523-c7f6-41aa-84a6-629041477aeb-0 | null | 0 | 1 | 434768 | 0 | 434768 |
show_commit_partitions
显示某次提交的分区信息。
输入
| 参数名 | 类型 | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| limit | Int | N | 10 | 返回的最大记录数 |
| instant_time | String | Y | 无 | 实例时间 |
输出
| 输出名 | 类型 |
|---|---|
| action | String |
| partition_path | String |
| total_files_added | Long |
| total_files_updated | Long |
| total_records_inserted | Long |
| total_records_updated | Long |
| total_bytes_written | Long |
| total_errors | Long |
示例
call show_commit_partitions(table => 'test_hudi_table', instant_time => '20230206174349556');| action | partition_path | total_files_added | total_files_updated | total_records_inserted | total_records_updated | total_bytes_written | total_errors |
|---|---|---|---|---|---|---|---|
| deltacommit | dt=2021-05-03 | 7fb52523-c7f6-41aa-84a6-629041477aeb-0 | 0 | 1 | 434768 | 0 | 0 |
show_commit_write_stats
显示某次提交的写入统计信息。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| limit | Int | 否 | 10 | 返回的最大记录数 |
| instant_time | String | 是 | 无 | 实例时间(instant time) |
输出
| 输出名 | 类型 |
|---|---|
| action | String |
| total_bytes_written | Long |
| total_records_written | Long |
| avg_record_size | Long |
示例
call show_commit_write_stats(table => 'test_hudi_table', instant_time => '20230206174349556');| 操作 | total_bytes_written | total_records_written | avg_record_size |
|---|---|---|---|
| deltacommit | 434768 | 1 | 434768 |
show_rollbacks
显示回滚提交。
输入
| 参数名 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| limit | Int | 否 | 10 | 返回的最大记录数 |
输出
| 输出名 | 类型 |
|---|---|
| instant | String |
| rollback_instant | String |
| total_files_deleted | Int |
| time_taken_in_millis | Long |
| total_partitions | Int |
示例
call show_rollbacks(table => 'test_hudi_table');| instant | rollback_instant | total_files_deleted | time_taken_in_millis | total_partitions |
|---|---|---|---|---|
| deltacommit | 434768 | 1 | 434768 | 2 |
show_rollback_detail
显示回滚提交的详细信息。
输入
| 参数名 | 类型 | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名称 |
| limit | Int | 否 | 10 | 返回的最大记录数 |
| instant_time | String | 是 | 无 | Instant 时间 |
输出
| 输出名 | 类型 |
|---|---|
| instant | String |
| rollback_instant | String |
| partition | String |
| deleted_file | String |
| succeeded | Int |
示例
call show_rollback_detail(table => 'test_hudi_table', instant_time => '20230206174349556');| instant | rollback_instant | partition | deleted_file | succeeded |
|---|---|---|---|---|
| deltacommit | 434768 | 1 | 434768 | 2 |
commits_compare
将提交与另一个路径进行比较。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名称 |
| path | String | 是 | 无 | 表的路径 |
输出
| 输出名称 | 类型 |
|---|---|
| compare_detail | String |
示例
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
归档提交。
输入
| 参数名称 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | 否 | 无 | Hudi 表名 |
| path | String | 否 | 无 | 表路径 |
| min_commits | Int | 否 | 20 | 与 hoodie.keep.max.commits 类似,但用于控制活动时间线(active timeline)中保留的最少 instant 数量。 |
| max_commits | Int | 否 | 30 | 每次写入后,归档服务会将时间线中的旧条目移入归档日志,以使元数据开销保持恒定,即使表规模不断增长。此配置用于控制活动时间线中保留的最大 instant 数量。 |
| retain_commits | Int | 否 | 10 | instant 的归档以尽力而为的方式分批进行,以便将更多 instant 打包到单个归档日志中。此配置用于控制该归档批次的大小。 |
| enable_metadata | Boolean | 否 | true | 启用内部元数据表 |
输出
| 输出名称 | 类型 |
|---|---|
| result | Int |
示例
call archive_commits(table => 'test_hudi_table');| result |
|---|
| 0 |
export_instants
将 instants 导出到本地文件夹。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| local_folder | String | 是 | 无 | 本地文件夹 |
| limit | Int | 否 | -1 | 要导出的 instants 数量 |
| actions | String | 否 | clean,commit,deltacommit,rollback,savepoint,restore | 提交操作类型 |
| desc | Boolean | 否 | false | 是否降序排列 |
输出
| 输出名 | 类型 |
|---|---|
| export_detail | String |
示例
call export_instants(table => 'test_hudi_table', local_folder => '/tmp/folder');| export_detail |
|---|
| 已将 6 个 Instant 导出到 /tmp/folder |
rollback_to_instant
将表回滚到某个时间点时处于当前状态的那次提交。
输入
| 参数名 | 类型 | 是否必需 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| instant_time | String | 是 | 无 | Instant 时间 |
输出
| 输出名 | 类型 |
|---|---|
| rollback_result | Boolean |
示例
将 test_hudi_table 回滚到某一个 instant
call rollback_to_instant(table => 'test_hudi_table', instant_time => '20220109225319449');| rollback_result |
|---|
| true |
create_savepoint
为 Hudi 表创建一个保存点(savepoint)。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| commit_time | String | Y | 无 | 提交时间 |
| user | String | N | "" | 用户名 |
| comments | String | N | "" | 备注 |
输出
| 输出名 | 类型 |
|---|---|
| create_savepoint_result | Boolean |
示例
call create_savepoint(table => 'test_hudi_table', commit_time => '20220109225319449');| create_savepoint_result |
|---|
| true |
show_savepoints
显示保存点(savepoint)。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名称 |
输出
| 输出名 | 类型 |
|---|---|
| savepoint_time | String |
示例
call show_savepoints(table => 'test_hudi_table');| savepoint_time |
|---|
| 20220109225319449 |
| 20220109225311742 |
| 20220109225301429 |
delete_savepoint
删除 Hudi 表中的保存点(savepoint)。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| instant_time | String | 是 | 无 | Instant 时间 |
输出
| 输出名称 | 类型 |
|---|---|
| delete_savepoint_result | Boolean |
示例
从 test_hudi_table 中删除一个保存点
call delete_savepoint(table => 'test_hudi_table', instant_time => '20220109225319449');| delete_savepoint_result |
|---|
| true |
rollback_to_savepoint
将表回滚到某个时刻的当前提交。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| instant_time | String | Y | 无 | Instant 时间 |
输出
| 输出名称 | 类型 |
|---|---|
| rollback_savepoint_result | Boolean |
示例
将 test_hudi_table 回滚到某个保存点
call rollback_to_savepoint(table => 'test_hudi_table', instant_time => '20220109225319449');| rollback_savepoint_result |
|---|
| true |
copy_to_temp_view
将表复制到临时视图。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | Y | None | Hudi 表名 |
| query_type | String | N | "snapshot" | 需要读取数据的方式:incremental 模式(读取某个 instantTime 以来的新数据)、read_optimized 模式(基于基础文件获取最新视图)、snapshot 模式(通过合并基础文件和(如有)日志文件获取最新视图) |
| view_name | String | Y | None | 视图名称 |
| begin_instance_time | String | N | "" | 起始实例时间 |
| end_instance_time | String | N | "" | 结束实例时间 |
| as_of_instant | String | N | "" | 截至 instant 时间 |
| replace | Boolean | N | false | 是否替换已存在的视图 |
| global | Boolean | N | false | 是否为全局视图 |
输出
| 输出名称 | 类型 |
|---|---|
| status | Int |
示例
call copy_to_temp_view(table => 'test_hudi_table', view_name => 'copy_view_test_hudi_table');| 状态 |
|---|
| 0 |
copy_to_table
将表复制到新表。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | None | Hudi 表名 |
| query_type | String | N | "snapshot" | 数据读取方式:incremental 模式(读取某个 instantTime 之后的新数据)、read_optimized 模式(基于基础文件获取最新视图),或 snapshot 模式(通过合并基础文件与(若存在的)日志文件获取最新视图) |
| new_table | String | Y | None | 新表的名称 |
| begin_instance_time | String | N | "" | 起始实例时间 |
| end_instance_time | String | N | "" | 结束实例时间 |
| as_of_instant | String | N | "" | 截至某个即时时间 |
| save_mode | String | N | "overwrite" | 保存模式 |
| columns | String | N | "" | 需要从源表复制到新表的列 |
输出
| 输出名称 | 类型 |
|---|---|
| status | Int |
示例
call copy_to_table(table => 'test_hudi_table', new_table => 'copy_table_test_hudi_table');| 状态 |
|---|
| 0 |
元数据表管理
create_metadata_table
创建 Hudi 表的元数据表。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名称 |
输出
| 输出名称 | 类型 |
|---|---|
| result | String |
示例
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 表的元数据表。
输入
| 参数名称 | 类型 | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| read_only | Boolean | N | false | 是否只读 |
输出
| 输出名称 | 类型 |
|---|---|
| result | String |
示例
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 表的元数据表。
输入
| 参数名 | 类型 | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | None | Hudi 表名称 |
输出
| 输出名 | 类型 |
|---|---|
| result | String |
示例
call delete_metadata_table(table => 'test_hudi_table');| 结果 |
|---|
| 已从 hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table/.hoodie/metadata 中移除元数据表 |
show_metadata_table_partitions
显示 Hudi 表的分区。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名称 |
输出
| 输出名 | 类型 |
|---|---|
| partition | String |
示例
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 表的文件。
输入
| 参数名 | 类型 | 是否必需 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 否 | 无 | Hudi 表名 |
| path | String | 否 | 无 | 表路径 |
| partition | String | 否 | "" | 分区名称 |
| limit | Int | 否 | 100 | 限制返回数量 |
| filter | String | 否 | "" | 用于过滤结果的高级谓词表达式(例如 file_path LIKE '%.parquet') |
note
调用此存储过程时,参数 table 和 path 至少必须指定其中一个。如果同时指定两个参数,则 table 生效。
输出
| 输出名 | 类型 |
|---|---|
| file_path | String |
示例
显示 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 表的元数据表统计信息。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名称 |
输出
| 输出名称 | 类型 |
|---|---|
| stat_key | String |
| stat_value | String |
示例
call show_metadata_table_stats(table => 'test_hudi_table');| stat_key | stat_value |
|---|---|
| dt=2021-05-03.totalBaseFileSizeInBytes | 23142 |
validate_metadata_table_files
校验 Hudi 表的元数据表文件。
输入
| 参数名 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| verbose | Boolean | N | False | 是否打印所有文件 |
输出
| 输出名 | 类型 |
|---|---|
| partition | String |
| file_name | String |
| is_present_in_fs | Boolean |
| is_present_in_metadata | Boolean |
| fs_size | Long |
| metadata_size | Long |
示例
call validate_metadata_table_files(table => 'test_hudi_table');| partition | file_name | is_present_in_fs | is_present_in_metadata | fs_size | metadata_size |
|---|---|---|---|---|---|
| dt=2021-05-03 | ad1e5a3f-532f-4a13-9f60-223676798bf3-0_0-4-4_00000000000002.parquet | true | true | 43523 | 43523 |
表信息
show_table_properties
显示表的 Hudi 属性。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| path | String | 否 | 无 | 表路径 |
| limit | Int | 否 | 10 | 返回的最大记录数 |
输出
| 输出名称 | 类型 |
|---|---|
| key | String |
| value | String |
示例
call show_table_properties(table => 'test_hudi_table', limit => 10);| 键 | 值 |
|---|---|
| hoodie.table.ordering.fields | ts |
| hoodie.table.partition.fields | dt |
show_fs_path_detail
显示路径的详细信息。
输入
| 参数名称 | 类型 | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| path | String | Y | 无 | Hudi 表名 |
| is_sub | Boolean | N | false | 是否列出文件 |
| sort | Boolean | N | true | 按 storage_size 排序 |
| limit | Int | N | 100 | 限制数量 |
输出
| 输出名称 | 类型 |
|---|---|
| path_num | Long |
| file_num | Long |
| storage_size | Long |
| storage_size(unit) | String |
| storage_path | String |
| space_consumed | Long |
| quota | Long |
| space_quota | Long |
示例
call show_fs_path_detail(path => 'hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table');| path_num | file_num | storage_size | storage_size(单位) | storage_path | space_consumed | quota | space_quota |
|---|---|---|---|---|---|---|---|
| 22 | 58 | 2065612 | 1.97MB | hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table | -1 | 6196836 | -1 |
stats_file_sizes
展示表的文件大小。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| partition_path | String | 否 | "" | 分区路径 |
| limit | Int | 否 | 10 | 返回记录的最大数量 |
输出
| 输出名称 | 类型 |
|---|---|
| commit_time | String |
| min | Long |
| 10th | Double |
| 50th | Double |
| avg | Double |
| 95th | Double |
| max | Long |
| num_files | Int |
| std_dev | Double |
示例
call stats_file_sizes(table => 'test_hudi_table');| commit_time | min | 10th | 50th | avg | 95th | max | num_files | std_dev |
|---|---|---|---|---|---|---|---|---|
| 20230205134149455 | 435000 | 435000.0 | 435000.0 | 435000.0 | 435000.0 | 435000 | 1 | 0.0 |
stats_wa
显示表的写入统计信息与写放大情况。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| limit | Int | 否 | 10 | 返回记录的最大条数 |
输出
| 输出名 | 类型 |
|---|---|
| commit_time | String |
| total_upserted | Long |
| total_written | Long |
| write_amplification_factor | String |
示例
call stats_wa(table => 'test_hudi_table');| commit_time | total_upserted | total_written | write_amplification_factor |
|---|---|---|---|
| 总计 | 0 | 0 | 0 |
show_logfile_records
显示表中日志文件的记录。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | N | 无 | Hudi 表名 |
| path | String | N | 无 | 表的路径 |
| log_file_path_pattern | String | Y | 无 | 日志文件路径匹配模式 |
| merge | Boolean | N | false | 是否合并结果 |
| limit | Int | N | 10 | 返回的最大记录数 |
| filter | String | N | "" | 过滤表达式 |
输出
| 输出名 | 类型 |
|---|---|
| records | String |
示例
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 的元数据。
输入
| 参数名 | 类型 | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| log_file_path_pattern | String | Y | 10 | logfile 的路径模式 |
| merge | Boolean | N | false | 是否合并结果 |
| limit | Int | N | 10 | 返回记录的最大数量 |
输出
| 输出名 | 类型 |
|---|---|
| instant_time | String |
| record_count | Int |
| block_type | String |
| header_metadata | String |
| footer_metadata | String |
示例
call show_logfile_metadata(table => 'hudi_mor_tbl', log_file_path_pattern => 'hdfs://ns1/hive/warehouse/hudi.db/hudi_mor_tbl/*.log*');| instant_time | record_count | block_type | header_metadata | footer_metadata |
|---|---|---|---|---|
| 20230205133427059 | 1 | AVRO_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 文件。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
| path | String | Y | None | Hudi 表路径 |
| parallelism | Int | N | 100 | 并行度 |
| limit | Int | N | 100 | 返回条数限制 |
| needDelete | Boolean | N | false | 是否需要删除 |
| partitions | String | N | "" | 分区 |
| instants | String | N | "" | Instants |
| filter | String | N | "" | 过滤表达式 |
输出
| 输出名称 | 类型 |
|---|---|
| Path | String |
示例
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
显示表的文件系统视图。
输入
| 参数名称 | 类型 | 是否必需 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名称 |
| max_instant | String | 否 | "" | 最大 instant 时间 |
| include_max | Boolean | 否 | false | 是否包含最大 instant |
| include_in_flight | Boolean | 否 | false | 是否包含进行中的 instant |
| exclude_compaction | Boolean | 否 | false | 是否排除 compaction |
| limit | Int | 否 | 10 | 返回记录的最大数量 |
| path_regex | String | 否 | "ALL_PARTITIONS" | 路径的匹配模式 |
输出
| 输出名称 | 类型 |
|---|---|
| partition | String |
| file_id | String |
| base_instant | String |
| data_file | String |
| data_file_size | Long |
| num_delta_files | Long |
| total_delta_file_size | Long |
| delta_files | String |
示例
call show_fsview_all(table => 'test_hudi_table');| 分区 | file_id | base_instant | data_file | data_file_size | num_delta_files | total_delta_file_size | delta_files |
|---|---|---|---|---|---|---|---|
| dt=2021-05-03 | d0073a12-085d-4f49-83e9-402947e7e90a-0 | 20220109225319449 | 7fb52523-c7f6-41aa-84a6-629041477aeb-0_0-92-99_20220109225319449.parquet | 5319449 | 1 | 213193 | .7fb52523-c7f6-41aa-84a6-629041477aeb-0_20230205133217210.log.1_0-60-63 |
show_fsview_latest
显示表的最新文件系统视图。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| max_instant | String | N | "" | 最大 instant 时间 |
| include_max | Boolean | N | false | 是否包含最大 instant |
| include_in_flight | Boolean | N | false | 是否包含进行中的 |
| exclude_compaction | Boolean | N | false | 是否排除 compaction |
| path_regex | String | N | "ALL_PARTITIONS" | 路径匹配模式 |
| partition_path | String | N | "ALL_PARTITIONS" | 分区路径 |
| merge | Boolean | N | false | 是否合并结果 |
输出
| 输出名 | 类型 |
|---|---|
| partition | String |
| file_id | String |
| base_instant | String |
| data_file | String |
| data_file_size | Long |
| num_delta_files | Long |
| total_delta_file_size | Long |
| delta_size_compaction_scheduled | Long |
| delta_size_compaction_unscheduled | Long |
| delta_to_base_radio_compaction_scheduled | Double |
| delta_to_base_radio_compaction_unscheduled | Double |
| delta_files_compaction_scheduled | String |
| delta_files_compaction_unscheduled | String |
示例
call show_fsview_latest(table => 'test_hudi_table', partition => 'dt=2021-05-03');| partition | file_id | base_instant | data_file | data_file_size | num_delta_files | total_delta_file_size | delta_files |
|---|---|---|---|---|---|---|---|
| dt=2021-05-03 | d0073a12-085d-4f49-83e9-402947e7e90a-0 | 20220109225319449 | 7fb52523-c7f6-41aa-84a6-629041477aeb-0_0-92-99_20220109225319449.parquet | 5319449 | 1 | 213193 | .7fb52523-c7f6-41aa-84a6-629041477aeb-0_20230205133217210.log.1_0-60-63 |
表服务
run_clustering
在 hoodie 表上触发聚类。通过使用分区谓词,可以在指定分区上运行聚类任务,也可以指定排序列来对数据进行排序。
note
每次调用都会生成新的聚类 instant,或者执行一些处于 pending 状态的聚类 instant。调用此存储过程时,参数 table 和 path 至少必须指定其中一个;如果同时给出两个参数,则以 table 为准。
输入
| 参数名称 | 类型 | 是否必需 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 否 | 无 | 要进行聚类的表名 |
| path | String | 否 | 无 | 要进行聚类的表路径 |
| predicate | String | 否 | 无 | 用于过滤分区的谓词 |
| order | String | 否 | 无 | 排序列,以 , 分隔 |
| show_involved_partition | Boolean | 否 | false | 在输出中显示涉及的分区 |
| op | String | 否 | 无 | 操作类型,EXECUTE 或 SCHEDULE |
| order_strategy | String | 否 | 无 | 记录布局优化方式,linear/z-order/hilbert |
| options | String | 否 | 无 | 自定义 Hudi 配置,格式为 "key1=value1,key2=value2` |
| instants | String | 否 | 无 | 指定的 instant,以 , 分隔 |
| selected_partitions | String | 否 | 无 | 要执行聚类的分区,以 , 分隔 |
| partition_regex_pattern | String | 否 | 无 | 用于过滤分区的正则表达式(例如 2025.*) |
| limit | Int | 否 | 无 | 要执行的最大计划数 |
输出
| 输出名称 | 类型 |
|---|---|
| timestamp | String |
| input_group_size | Int |
| state | String |
| involved_partitions | String |
示例
使用表名对 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 为准。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 否 | 无 | 要进行聚类的表名 |
| path | String | 否 | 无 | 要进行聚类的表路径 |
| limit | Int | 否 | 20 | 返回的最大记录数 |
| show_involved_partition | Boolean | 否 | false | 是否在输出中显示涉及的分区 |
输出
| 输出名 | 类型 |
|---|---|
| timestamp | String |
| input_group_size | Int |
| state | String |
| involved_partitions | String |
示例
根据表名显示待处理的聚类任务
call show_clustering(table => 'test_hudi_table');| timestamp | groups |
|---|---|
| 20220408153707928 | 2 |
| 20220408153636963 | 3 |
按表路径显示待执行的 clustering 任务
call show_clustering(path => '/tmp/hoodie/test_hudi_table');| timestamp | groups |
|---|---|
| 20220408153707928 | 2 |
| 20220408153636963 | 3 |
显示待处理的聚类操作,包含表名和限制数量
call show_clustering(table => 'test_hudi_table', limit => 1);| timestamp | groups |
|---|---|
| 20220408153707928 | 2 |
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 生效。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| op | String | 是 | 无 | 操作类型,RUN 或 SCHEDULE |
| table | String | 否 | 无 | 要进行 compaction 的表名 |
| path | String | 否 | 无 | 要进行 compaction 的表路径 |
| timestamp | Long | 否 | 无 | Instant time |
| options | String | 否 | 无 | 以逗号分隔的 Hudi compaction 配置列表,格式为 "config1=value1,config2=value2" |
| instants | String | 否 | 无 | 以 , 分隔的指定 instants |
| limit | Int | 否 | 无 | 最多执行的计划数 |
输出
| 输出名 | 类型 |
|---|---|
| timestamp | String |
| operation_size | Int |
| state | String |
示例
使用表名执行 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 为准。
输入
| 参数名称 | 类型 | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 否 | 无 | 要显示 compaction 的表名 |
| path | String | 否 | 无 | 要显示 compaction 的表路径 |
| limit | Int | 否 | 20 | 返回的最大记录数 |
输出
| 参数名称 | 类型 | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| timestamp | String | 否 | 无 | Instant 时间 |
| operation_size | Int | 否 | 无 | 待压缩的 file slice 数量 |
| state | String | 否 | 无 | compaction 的状态 |
示例
按表名显示 compaction
call show_compaction(table => 'test_hudi_table');| timestamp | action | size |
|---|---|---|
| 20220408153707928 | compaction | 10 |
| 20220408153636963 | compaction | 10 |
显示带有表路径的 compaction
call show_compaction(path => '/tmp/hoodie/test_hudi_table');| 时间戳 | 操作 | 大小 |
|---|---|---|
| 20220408153707928 | compaction | 10 |
| 20220408153636963 | compaction | 10 |
显示 compaction,并附带表名和数量限制
call show_compaction(table => 'test_hudi_table', limit => 1);| timestamp | action | size |
|---|---|---|
| 20220408153707928 | compaction | 10 |
run_clean
在 hoodie 表上运行清理(cleaner)。
输入
| 参数名称 | 类型 | 是否必需 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | 是 | 无 | 要执行清理(clean)的表名称 |
| schedule_in_line | Boolean | 否 | true | 如果希望调度并立即运行一次 clean,请设置为 "true";如果已经调度了 clean 并希望运行该调度,则设置为 false。 |
| clean_policy | String | 否 | 无 | 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_commits | Int | 否 | 无 | 使用 KEEP_LATEST_COMMITS 清理策略时,需要保留而不被清理的提交次数。数据将保留 num_of_commits * time_between_commits(已调度)这么长的时间。这也直接决定了该表在增量查询方面所能支持的数据保留时长。 |
| hours_retained | Int | 否 | 无 | 使用 KEEP_LATEST_BY_HOURS 清理策略时,需要保留提交的小时数。与按提交次数保留相比,该配置提供了更灵活的选项。设置此属性后,提交时间早于所配置保留小时数的提交所对应的文件组中,除最新文件外的所有文件都将被清理。 |
| file_versions_retained | Int | 否 | 无 | 使用 KEEP_LATEST_FILE_VERSIONS 清理策略时,清理过程中每个文件组内需保留的最少文件切片数量。 |
| trigger_strategy | String | 否 | 无 | org.apache.hudi.table.action.clean.CleaningTriggerStrategy:控制何时调度清理。NUM_COMMITS(默认):每 N 次提交触发一次清理服务,N 由 hoodie.clean.max.commits 决定。 |
| trigger_max_commits | Int | 否 | 无 | 上一次清理操作之后,需要经过多少次提交才会尝试调度新的清理。 |
| options | String | 否 | 无 | 以逗号分隔的 Hudi 清理相关配置列表,格式为 "config1=value1,config2=value2"。 |
输出
| 参数名称 | 类型 |
|---|---|
| start_clean_time | String |
| time_taken_in_millis | Long |
| total_files_deleted | Int |
| earliest_commit_to_retain | String |
| bootstrap_part_metadata | String |
| version | Int |
示例
使用表名运行 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
显示已完成的清理操作,包含耗时、删除文件数以及保留策略等元数据。
输入
| 参数名称 | 类型 | 是否必需 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| limit | Int | N | 10 | 返回的最大记录数 |
| showArchived | Boolean | N | false | 是否包含已归档的时间线数据 |
| filter | String | N | "" | 用于过滤结果的高级谓词表达式(例如 total_files_deleted > 10) |
输出
| 输出名称 | 类型 |
|---|---|
| clean_time | String |
| state_transition_time | String |
| action | String |
| start_clean_time | String |
| time_taken_in_millis | Long |
| total_files_deleted | Int |
| earliest_commit_to_retain | String |
| last_completed_commit_timestamp | String |
| version | Int |
示例
显示指定表名的已完成清理操作
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)下的清理操作及其状态信息。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| limit | Int | N | 10 | 返回的最大记录数 |
| showArchived | Boolean | N | false | 是否包含已归档的时间线数据 |
| filter | String | N | "" | 用于过滤结果的高级谓词表达式(例如 state = 'COMPLETED') |
输出
| 输出名 | 类型 |
|---|---|
| plan_time | String |
| state | String |
| action | String |
| earliest_instant_to_retain | String |
| last_completed_commit_timestamp | String |
| policy | String |
| version | Int |
| total_partitions_to_clean | Int |
| total_partitions_to_delete | Int |
| extra_metadata | String |
示例
按表名显示清理计划
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
显示分区级别的清理详情,用于调试目的。
输入
| 参数名 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | Y | None | Hudi 表名 |
| limit | Int | N | 10 | 返回记录的最大数量 |
| showArchived | Boolean | N | false | 是否包含已归档的时间线数据 |
| filter | String | N | "" | 用于过滤结果的高级谓词表达式(例如 partition_path LIKE '2025%') |
输出
| 输出名 | 类型 |
|---|---|
| clean_time | String |
| state_transition_time | String |
| action | String |
| start_clean_time | String |
| partition_path | String |
| policy | String |
| delete_path_patterns | Int |
| success_delete_files | Int |
| failed_delete_files | Int |
| is_partition_deleted | Boolean |
| time_taken_in_millis | Long |
| total_files_deleted | Int |
示例
显示指定表名的清理元数据
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
查看指定文件的状态,例如是否已被删除、由哪个操作删除等。
输入
| 参数名 | 类型 | 是否必需 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 否 | 无 | Hudi 表名 |
| path | String | 否 | 无 | 表路径 |
| partition | String | 否 | 无 | 分区路径(分区表必须指定) |
| file | String | 是 | 无 | 要检查的文件名 |
| filter | String | 否 | "" | 用于过滤结果的高级谓词表达式 |
:::note
调用此存储过程时,table 和 path 参数至少需要指定其中一个。如果两个参数都指定,则以 table 为准。对于分区表,必须指定 partition 参数。
:::
输出
| 输出名 | 类型 |
|---|---|
| status | String |
| action | String |
| instant | String |
| timeline | String |
| full_path | String |
示例
通过表名查看文件状态
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,生存时间)操作,根据保留策略删除分区。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| ttl_policy | String | 否 | 无 | TTL 策略类型 |
| retain_days | Int | 否 | 无 | 分区保留天数 |
| options | String | 否 | 无 | TTL 相关的 Hudi 配置,以逗号分隔,格式为 "config1=value1,config2=value2" |
输出
| 输出名称 | 类型 |
|---|---|
| deleted_partitions | String |
示例
使用表名运行 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 文件。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | None | Hudi 表名 |
| instant_time | String | Y | None | Instant 名称 |
输出
| 输出名称 | 类型 |
|---|---|
| delete_marker_result | Boolean |
示例
call delete_marker(table => 'test_hudi_table', instant_time => '20230206174349556');| delete_marker_result |
|---|
| true |
sync_validate
验证同步过程。
输入
| 参数名称 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| src_table | String | 是 | 无 | 源表名称 |
| dst_table | String | 是 | 无 | 目标表名称 |
| mode | String | 是 | "complete" | 模式 |
| hive_server_url | String | 是 | 无 | Hive server 地址 |
| hive_pass | String | 是 | 无 | Hive 密码 |
| src_db | String | 否 | "rawdata" | 源数据库 |
| target_db | String | 否 | "dwh_hoodie" | 目标数据库 |
| partition_cnt | Int | 否 | 5 | 分区数量 |
| hive_user | String | 否 | "" | Hive 用户名 |
输出
| 输出名称 | 类型 |
|---|---|
| result | String |
示例
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。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| metastore_uri | String | 否 | "" | Metastore_uri |
| username | String | 否 | "" | 用户名 |
| password | String | 否 | "" | 密码 |
| use_jdbc | String | 否 | "" | 启用 Hive 同步时是否使用 JDBC |
| mode | String | 否 | "" | Hive 操作的模式。可选值为 hms、jdbc 和 hiveql。 |
| partition_fields | String | 否 | "" | 表中用于确定 Hive 分区列的字段。 |
| partition_extractor_class | String | 否 | "" | 实现 PartitionValueExtractor 以提取分区值的类,默认为 'org.apache.hudi.hive.MultiPartKeysValueExtractor'。 |
| strategy | String | 否 | "" | Hive 表同步策略。可选值:RO、RT、ALL。 |
| sync_incremental | String | 否 | "" | 是否将分区增量同步到 Metastore,即根据提交元数据仅同步新增、变更和删除的分区。若设置为 false,当分区丢失时,元数据同步将执行全量分区同步操作。 |
输出
| 输出名称 | 类型 |
|---|---|
| result | String |
示例
call hive_sync(table => 'test_hudi_table');| result |
|---|
| true |
hdfs_parquet_import
向 Hudi 表中添加 Parquet 文件。
输入
| 参数名 | 类型 | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| table_type | String | 是 | "" | 表类型,MERGE_ON_READ 或 COPY_ON_WRITE |
| src_path | String | 是 | "" | 源路径 |
| target_path | String | 是 | "" | 目标路径 |
| row_key | String | 是 | "" | 主键 |
| partition_key | String | 是 | "" | 分区键 |
| schema_file_path | String | 是 | "" | Schema 文件路径 |
| format | String | 否 | "parquet" | 文件格式 |
| command | String | 否 | "insert" | 导入命令 |
| retry | Int | 否 | 0 | 重试次数 |
| parallelism | Int | 否 | 无 | 并行度 |
| props_file_path | String | 否 | "" | 属性文件路径 |
输出
| 输出名 | 类型 |
|---|---|
| import_result | Int |
示例
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 表的分区添加操作。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | None | Hudi 表名称 |
| dry_run | Boolean | N | true | 试运行 |
输出
| 输出名称 | 类型 |
|---|---|
| partition_path | String |
| metadata_is_present | String |
| action | String |
示例
call repair_add_partition_meta(table => 'test_hudi_table');| partition_path | metadata_is_present | action |
|---|---|---|
| dt=2021-05-03 | 是 | 无 |
repair_corrupted_clean_files
修复 Hudi 表中损坏的 clean 文件。
输入
| 参数名称 | 类型 | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名称 |
输出
| 输出名称 | 类型 |
|---|---|
| result | Boolean |
示例
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)。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| duplicated_partition_path | String | Y | 无 | 重复的分区路径 |
| repaired_output_path | String | Y | 无 | 修复后的输出路径 |
| dry_run | Boolean | N | true | 试运行 |
| dedupe_type | String | N | "insert_type" | 去重类型 |
输出
| 输出名称 | 类型 |
|---|---|
| result | String |
示例
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 表。
输入参数
| 参数名 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | 是 | 无 | Hudi 表名 |
| dry_run | Boolean | 否 | true | 试运行 |
输出
| 输出名 | 类型 |
|---|---|
| partition_path | String |
| text_metafile_present | String |
| base_metafile_present | String |
| action | String |
示例
call repair_migrate_partition_meta(table => 'test_hudi_table');repair_overwrite_hoodie_props
覆盖 Hudi 表属性。
输入
| 参数名称 | 类型 | 是否必需 | 默认值 | 说明 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| new_props_file_path | String | Y | 无 | 新属性文件的路径 |
输出
| 输出名称 | 类型 |
|---|---|
| property | String |
| old_value | String |
| new_value | String |
示例
call repair_overwrite_hoodie_props(table => 'test_hudi_table', new_props_file_path = > '/tmp/props');| 属性 | old_value | new_value |
|---|---|---|
| hoodie.file.index.enable | true | false |
引导(Bootstrap)
run_bootstrap
将现有表转换为 Hudi 表。
输入
| 参数名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | None | 要进行聚类的表名 |
| table_type | String | Y | None | 表类型,MERGE_ON_READ 或 COPY_ON_WRITE |
| bootstrap_path | String | Y | None | 需要作为 Hudi 表进行引导(bootstrap)的数据集基础路径 |
| base_path | String | Y | None | 基础路径 |
| rowKey_field | String | Y | None | 主键字段 |
| base_file_format | String | N | "PARQUET" | 基础文件格式 |
| partition_path_field | String | N | "" | 分区列字段 |
| bootstrap_index_class | String | N | "org.apache.hudi.common.bootstrap.index.HFileBootstrapIndex" | 用于将骨架基础文件映射到引导基础文件的实现类。 |
| selector_class | String | N | "org.apache.hudi.client.bootstrap.selector.MetadataOnlyBootstrapModeSelector" | 选择引导数据集中每个文件/分区的引导模式 |
| key_generator_class | String | N | "org.apache.hudi.keygen.SimpleKeyGenerator" | 键生成器类 |
| full_bootstrap_input_provider | String | N | "org.apache.hudi.bootstrap.SparkParquetBootstrapDataProvider" | 完整引导输入提供者类 |
| schema_provider_class | String | N | "" | Schema 提供者类 |
| payload_class | String | N | "org.apache.hudi.common.model.OverwriteWithLatestAvroPayload" | Payload 类 |
| parallelism | Int | N | 1500 | 对于仅元数据的引导,Hudi 会并行化该操作,使每个表分区由一个 Spark 任务处理。此配置限制并行度。如果表分区数大于该配置值,则使用配置的并行度;如果表分区数小于该配置值,则将并行度设为表分区数。对于全记录引导,即记录的 BULK_INSERT 操作,该配置值将作为 BULK_INSERT shuffle 并行度(hoodie.bulkinsert.shuffle.parallelism)传入,决定 BULK_INSERT 的写入行为。如果由于并行度受限导致引导较慢,可以调大此值。 |
| enable_hive_sync | Boolean | N | false | 是否启用 Hive 同步 |
| props_file_path | String | N | "" | 属性文件路径 |
| bootstrap_overwrite | Boolean | N | false | 覆盖引导路径 |
输出
| 输出名称 | 类型 |
|---|---|
| status | Int |
示例
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 表的映射文件。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 否 | 无 | Hudi 表名 |
| path | String | 否 | 无 | 表的路径 |
| partition_path | String | 否 | "" | 分区路径 |
| file_ids | String | 否 | "" | 文件 id |
| limit | Int | 否 | 10 | 返回记录的最大数量 |
| sort_by | String | 否 | "partition" | 排序列 |
| desc | Boolean | 否 | false | 是否降序排列 |
| filter | String | 否 | "" | 用于过滤结果的高级谓词表达式(例如 partition LIKE '2025%') |
note
调用此过程时,参数 table 与 path 至少需要指定其中一个。如果两个参数都指定,则 table 生效。
输出
| 输出名 | 类型 |
|---|---|
| partition | String |
| file_id | String |
| source_base_path | String |
| source_partition | String |
| source_file | String |
示例
根据表名显示 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'");| partition | file_id | source_base_path | source_partition | source_file |
|---|---|---|---|---|
| dt=2021-05-03 | d0073a12-085d-4f49-83e9-402947e7e90a-0 | hdfs://ns1/hive/warehouse/hudi.db/test_hudi_table/dt=2021-05-03/d0073a12-085d-4f49-83e9-402947e7e90a-0_0-2-2_00000000000002.parquet | dt=2021-05-03 | hdfs://ns1/tmp/dt=2021-05-03/00001.parquet |
show_bootstrap_partitions
显示 bootstrap 表的分区。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | 是 | 无 | 要进行聚簇的表的名称 |
输出
| 参数名 | 类型 |
|---|---|
| indexed_partitions | String |
示例
call show_bootstrap_partitions(table => 'test_hudi_table');| indexed_partitions |
|---|
| dt=2021-05-03 |
版本管理
upgrade_table
将 Hudi 表升级到指定版本。
输入
| 参数名 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | 无 | Hudi 表名 |
| to_version | String | Y | 无 | hoodie 表的版本号 |
输出
| 输出名 | 类型 |
|---|---|
| result | Boolean |
示例
call upgrade_table(table => 'test_hudi_table', to_version => 'FIVE');| result |
|---|
| true |
downgrade_table
将 Hudi 表降级到指定版本。
输入
| 参数名称 | 类型 | 是否必需 | 默认值 | 描述 |
|---|---|---|---|---|
| table | String | Y | None | Hudi 表名称 |
| to_version | String | Y | None | hoodie 表的目标版本 |
输出
| 输出名称 | 类型 |
|---|---|
| result | Boolean |
示例
call downgrade_table(table => 'test_hudi_table', to_version => 'FOUR');| 结果 |
|---|
| true |
评论
登录后参与评论
KnowForge