스파크 프로시저
스파크 프로시저 (Spark Procedures)
아이스버그는 스파크의 저장 프로시저(stored procedures)를 제공해서 테이블 관리, 마이그레이션, 메타데이터 정보, CDC, 테이블 통계, 테이블 복제 등의 작업을 수행할 수 있게 해요. 이 문서에서는 각 프로시저의 사용법, 인자, 출력, 예시를 상세히 알려드릴게요. Spark 3.x에서는 아이스버그 SQL 확장이 필요해요.
출처: 문서
본문
스파크에서 아이스버그를 사용하려면 먼저 스파크 카탈로그를 구성해요. Spark 3.x의 경우 저장 프로시저는 스파크에서 아이스버그 SQL 확장을 사용할 때만 사용할 수 있어요. Spark 4.0의 경우 저장 프로시저는 아이스버그 SQL 확장 없이 네이티브로 지원돼요. 다만 Spark 4.0에서는 대소문자를 구분한다는 점에 유의해주세요.
사용법 (Usage)
프로시저는 CALL로 구성된 어떤 아이스버그 카탈로그에서든 사용할 수 있어요. 모든 프로시저는 system 네임스페이스에 있어요.
CALL은 이름(권장) 또는 위치로 인자를 전달하는 것을 지원해요. 위치와 이름 인자를 섞는 것은 지원되지 않아요.
이름 있는 인자 (Named arguments)
모든 프로시저 인자는 이름이 있어요. 이름으로 인자를 전달할 때 인자는 어떤 순서로든 올 수 있고 어떤 선택 인자도 생략할 수 있어요.
CALL catalog_name.system.procedure_name(arg_name_2 => arg_2, arg_name_1 => arg_1);
위치 인자 (Positional arguments)
위치로 인자를 전달할 때는 끝 인자만 선택이면 생략할 수 있어요.
CALL catalog_name.system.procedure_name(arg_1, arg_2, ... arg_n);
스냅샷 관리 (Snapshot management)
rollback_to_snapshot
테이블을 특정 스냅샷 ID로 롤백해요.
특정 시간으로 롤백하려면 rollback_to_timestamp를 사용해요.
정보 (Info)
이 프로시저는 영향받는 테이블을 참조하는 모든 캐시된 스파크 계획을 무효화해요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 업데이트할 테이블 이름 |
| snapshot_id | ✔️ | long | 롤백할 스냅샷 ID |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| previous_snapshot_id | long | 롤백 전의 현재 스냅샷 ID |
| current_snapshot_id | long | 새 현재 스냅샷 ID |
예시 (Example)
테이블 db.sample을 스냅샷 ID 1로 롤백:
CALL catalog_name.system.rollback_to_snapshot('db.sample', 1);
rollback_to_timestamp
테이블을 어떤 시점의 현재 스냅샷으로 롤백해요.
정보 (Info)
이 프로시저는 영향받는 테이블을 참조하는 모든 캐시된 스파크 계획을 무효화해요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 업데이트할 테이블 이름 |
| timestamp | ✔️ | timestamp | 롤백할 타임스탬프 |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| previous_snapshot_id | long | 롤백 전의 현재 스냅샷 ID |
| current_snapshot_id | long | 새 현재 스냅샷 ID |
예시 (Example)
db.sample을 특정 날짜와 시간으로 롤백해요.
CALL catalog_name.system.rollback_to_timestamp('db.sample', TIMESTAMP '2021-06-30 00:00:00.000');
set_current_snapshot
테이블의 현재 스냅샷 ID를 설정해요.
롤백과 달리, 스냅샷이 현재 테이블 상태의 조상일 필요는 없어요.
정보 (Info)
이 프로시저는 영향받는 테이블을 참조하는 모든 캐시된 스파크 계획을 무효화해요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 업데이트할 테이블 이름 |
| snapshot_id | long | 현재로 설정할 스냅샷 ID | |
| ref | string | 현재로 설정할 스냅샷 참조(브랜치 또는 태그) |
snapshot_id 또는 ref 중 하나는 제공해야 하지만 둘 다는 안 돼요.
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| previous_snapshot_id | long | 롤백 전의 현재 스냅샷 ID |
| current_snapshot_id | long | 새 현재 스냅샷 ID |
예시 (Example)
db.sample의 현재 스냅샷을 1로 설정:
CALL catalog_name.system.set_current_snapshot('db.sample', 1);
db.sample의 현재 스냅샷을 태그 s1로 설정:
CALL catalog_name.system.set_current_snapshot(table => 'db.sample', ref => 's1');
cherrypick_snapshot
스냅샷의 변경을 현재 테이블 상태로 체리픽해요.
체리픽은 원본을 변경하거나 제거하지 않고 기존 스냅샷에서 새 스냅샷을 만들어요.
append와 dynamic overwrite 스냅샷만 체리픽할 수 있어요.
정보 (Info)
이 프로시저는 영향받는 테이블을 참조하는 모든 캐시된 스파크 계획을 무효화해요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 업데이트할 테이블 이름 |
| snapshot_id | ✔️ | long | 체리픽할 스냅샷 ID |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| source_snapshot_id | long | 체리픽 전의 테이블 현재 스냅샷 |
| current_snapshot_id | long | 체리픽 적용으로 만들어진 스냅샷 ID |
예시 (Examples)
스냅샷 1 체리픽:
CALL catalog_name.system.cherrypick_snapshot('my_table', 1);
이름 있는 인자로 스냅샷 1 체리픽:
CALL catalog_name.system.cherrypick_snapshot(snapshot_id => 1, table => 'my_table' );
publish_changes
스테이징된 WAP ID의 변경을 현재 테이블 상태로 게시해요.
publish_changes는 원본을 변경하거나 제거하지 않고 기존 스냅샷에서 새 스냅샷을 만들어요.
append와 dynamic overwrite 스냅샷만 성공적으로 게시할 수 있어요.
publish_changes 프로시저는 제공된 wap_id를 가진 스냅샷이 테이블에 여러 개 있으면 실패해요.
정보 (Info)
이 프로시저는 영향받는 테이블을 참조하는 모든 캐시된 스파크 계획을 무효화해요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 업데이트할 테이블 이름 |
| wap_id | ✔️ | string | 스테이지에서 프로덕션으로 게시할 wap_id |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| source_snapshot_id | long | 변경을 게시하기 전의 테이블 현재 스냅샷 |
| current_snapshot_id | long | 변경 적용으로 만들어진 스냅샷 ID |
예시 (Examples)
WAP ID 'wap_id_1'로 publish_changes:
CALL catalog_name.system.publish_changes('my_table', 'wap_id_1');
이름 있는 인자로 publish_changes:
CALL catalog_name.system.publish_changes(wap_id => 'wap_id_2', table => 'my_table');
fast_forward
한 브랜치의 현재 스냅샷을 다른 브랜치의 최신 스냅샷으로 fast-forward해요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 업데이트할 테이블 이름 |
| branch | ✔️ | string | fast-forward할 브랜치 이름 |
| to | ✔️ | string |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| branch_updated | string | fast-forward된 브랜치 이름 |
| previous_ref | long | fast-forward 적용 전의 스냅샷 ID |
| updated_ref | long | fast-forward 적용 후의 현재 스냅샷 ID |
예시 (Examples)
main 브랜치를 audit-branch의 헤드로 fast-forward:
CALL catalog_name.system.fast_forward('my_table', 'main', 'audit-branch');
메타데이터 관리 (Metadata management)
많은 유지보수 작업을 아이스버그 저장 프로시저로 수행할 수 있어요.
expire_snapshots
아이스버그의 각 쓰기/업데이트/삭제/업서트/컴팩션은 스냅샷 격리와 타임 트래블을 위해 기존 데이터와 메타데이터를 보관하면서 새 스냅샷을 만들어요. expire_snapshots 프로시저는 더 이상 필요 없는 이전 스냅샷과 그 파일들을 제거하는 데 사용할 수 있어요.
이 프로시저는 구 스냅샷과 그 구 스냅샷이 고유하게 필요로 하는 데이터 파일을 제거해요. 즉 expire_snapshots 프로시저는 만료되지 않은 스냅샷이 여전히 필요로 하는 파일을 절대 제거하지 않아요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 업데이트할 테이블 이름 |
| older_than | timestamp | 이 타임스탬프 이전의 스냅샷이 제거됨 (기본값: 5일 전) | |
| retain_last | int | older_than과 무관하게 보존할 조상 스냅샷 수 (기본값 1) | |
| max_concurrent_deletes | int | 삭제 파일 작업에 사용되는 스레드 풀 크기 (기본적으로 스레드 풀 미사용) | |
| stream_results | boolean | true이면 삭제 파일이 RDD 파티션으로 스파크 드라이버에 전송됨 (기본적으로 모든 파일이 스파크 드라이버로 전송). 큰 파일 크기로 인한 스파크 드라이버 OOM을 방지하려면 이 옵션을 true로 설정하는 것을 권장해요 | |
| snapshot_ids | long 배열 | 만료시킬 스냅샷 ID 배열 | |
| clean_expired_metadata | boolean | true이면 스냅샷이 더 이상 참조하지 않는 파티션 스펙과 스키마 같은 만료된 메타데이터를 정리해요 |
older_than과 retain_last가 생략되면 테이블의 만료 속성이 사용돼요. 브랜치나 태그가 여전히 참조하는 스냅샷은 제거되지 않아요. 기본적으로 브랜치와 태그는 만료되지 않지만, 테이블 속성 history.expire.max-ref-age-ms로 보존 정책을 변경할 수 있어요. main 브랜치는 만료되지 않아요.
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| deleted_data_files_count | long | 이 연산으로 삭제된 데이터 파일 수 |
| deleted_position_delete_files_count | long | 이 연산으로 삭제된 position delete 파일 수 |
| deleted_equality_delete_files_count | long | 이 연산으로 삭제된 equality delete 파일 수 |
| deleted_manifest_files_count | long | 이 연산으로 삭제된 매니페스트 파일 수 |
| deleted_manifest_lists_count | long | 이 연산으로 삭제된 매니페스트 리스트 파일 수 |
| deleted_statistics_files_count | long | 이 연산으로 삭제된 통계 파일 수 |
예시 (Examples)
특정 날짜와 시간보다 오래된 스냅샷을 제거하되, 마지막 100개 스냅샷은 보존:
CALL hive_prod.system.expire_snapshots('db.sample', TIMESTAMP '2021-06-30 00:00:00.000', 100);
스냅샷 ID 123으로 스냅샷 제거 (이 스냅샷 ID는 현재 스냅샷이 아니어야 한다는 점에 유의):
CALL hive_prod.system.expire_snapshots(table => 'db.sample', snapshot_ids => ARRAY(123));
remove_orphan_files
아이스버그 테이블의 어떤 메타데이터 파일에서도 참조되지 않아 "고아(orphaned)"로 간주할 수 있는 파일을 제거하는 데 사용돼요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 정리할 테이블 이름 |
| older_than | timestamp | 이 타임스탬프 이전에 생성된 고아 파일 제거 (기본값: 3일 전) | |
| location | string | 파일을 찾을 디렉터리 (기본값: 테이블 위치) | |
| dry_run | boolean | true이면 실제로 파일을 제거하지 않음 (기본값: false) | |
| max_concurrent_deletes | int | 삭제 파일 작업에 사용되는 스레드 풀 크기 (기본적으로 스레드 풀 미사용) | |
| stream_results | boolean | true이면 고아 파일이 RDD 파티션으로 스파크 드라이버에 전송됨 (기본적으로 모든 파일이 스파크 드라이버로 전송). 이 옵션은 큰 파일 크기로 인한 스파크 드라이버 OOM을 방지하기 위해 true로 설정하는 것을 권장해요. 활성화되면 출력은 최대 20,000개의 파일 경로 샘플을 포함해요 | |
| file_list_view | string | 파일을 찾을 데이터셋 (디렉터리 나열을 건너뜀) | |
| equal_schemes | map | 동일하게 간주할 파일 시스템 스킴 매핑. 키는 쉼표 구분 스킴 목록, 값은 스킴 (기본값: map('s3a,s3n','s3')). | |
| equal_authorities | map | 동일하게 간주할 파일 시스템 authority 매핑. 키는 쉼표 구분 authority 목록, 값은 authority | |
| prefix_mismatch_mode | string | 위치 프리픽스(스킴/authority)가 불일치할 때의 동작: ERROR - 예외 발생 (기본값). IGNORE - 아무 조치 없음. DELETE - 파일 삭제 | |
| prefix_listing | boolean | true이면 SupportsPrefixOperations 인터페이스를 통해 프리픽스 기반 파일 나열 사용. 이 플래그를 활성화하면 Table FileIO 구현이 SupportsPrefixOperations를 지원해야 해요 (기본값: false) |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| orphan_file_location | String | 이 명령이 고아로 판단한 각 파일의 경로 |
예시 (Examples)
이 테이블에서 remove_orphan_files 명령의 dry run을 수행해서, 실제로 제거하지 않고 제거 후보인 모든 파일을 나열:
CALL catalog_name.system.remove_orphan_files(table => 'db.sample', dry_run => true);
테이블 db.sample이 알지 못하는 tablelocation/data 폴더의 모든 파일 제거:
CALL catalog_name.system.remove_orphan_files(table => 'db.sample', location => 'tablelocation/data');
테이블 db.sample이 알지 못하는 files_view 뷰의 모든 파일 제거:
Dataset<Row> compareToFileList =
spark
.createDataFrame(allFiles, FilePathLastModifiedRecord.class)
.withColumnRenamed("filePath", "file_path")
.withColumnRenamed("lastModified", "last_modified");
String fileListViewName = "files_view";
compareToFileList.createOrReplaceTempView(fileListViewName);
CALL catalog_name.system.remove_orphan_files(table => 'db.sample', file_list_view => 'files_view');
파일이 메타데이터 파일의 참조와 위치 프리픽스(스킴/authority)를 제외하고 일치하면 기본적으로 오류가 발생해요. prefix_mismatch_mode를 IGNORE로 설정하면 오류를 무시하고 파일을 건너뛰어요.
CALL catalog_name.system.remove_orphan_files(table => 'db.sample', prefix_mismatch_mode => 'IGNORE');
prefix_mismatch_mode를 DELETE로 설정하면 파일을 계속 삭제할 수 있어요.
CALL catalog_name.system.remove_orphan_files(table => 'db.sample', prefix_mismatch_mode => 'DELETE');
불일치한 프리픽스를 동등하게 간주해서 파일을 삭제할 수도 있어요.
CALL catalog_name.system.remove_orphan_files(table => 'db.sample', equal_schemes => map('file', 'file1'));
CALL catalog_name.system.remove_orphan_files(table => 'db.sample', equal_authorities => map('ns1', 'ns2'));
프리픽스 나열을 사용해서 제거 후보인 모든 파일을 나열:
CALL catalog_name.system.remove_orphan_files(table => 'db.sample', prefix_listing => true);
rewrite_data_files
아이스버그는 테이블의 각 데이터 파일을 추적해요. 데이터 파일이 많을수록 매니페스트 파일에 저장되는 메타데이터가 많아지고, 작은 데이터 파일은 파일을 여는 비용 때문에 불필요한 메타데이터와 덜 효율적인 쿼리를 만들 수 있어요.
아이스버그는 rewriteDataFiles 액션으로 스파크를 사용해 데이터 파일을 병렬로 컴팩션할 수 있어요. 작은 파일을 더 큰 파일로 결합해서 메타데이터 오버헤드와 런타임 파일 오픈 비용을 줄여요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 업데이트할 테이블 이름 |
| strategy | string | 전략 이름 - binpack 또는 sort. 기본값은 binpack 전략 | |
| sort_order | string | Zorder의 경우 zorder() 안에 쉼표 구분 컬럼 목록. 예: zorder(c1,c2,c3). 그 외에는 (ColumnName SortDirection NullOrder) 형식의 쉼표 구분 정렬 순서. SortDirection은 ASC 또는 DESC 가능. NullOrder는 NULLS FIRST 또는 NULLS LAST 가능. 기본값은 테이블의 정렬 순서 | |
| options | map | 작업에 사용할 옵션 | |
| where | string | 파일 필터링에 사용되는 문자열 조건. 필터와 일치하는 데이터를 포함할 수 있는 모든 파일이 재작성에 선택된다는 점에 유의해주세요. |
옵션 (Options)
일반 옵션 (General Options)
| 이름 | 기본값 | 설명 |
|---|---|---|
| max-concurrent-file-group-rewrites | 5 | 동시에 재작성할 최대 파일 그룹 수 |
| partial-progress.enabled | false | 전체 재작성이 완료되기 전에 파일 그룹 커밋 활성화 |
| partial-progress.max-commits | 10 | 부분 진행이 활성화되면 이 재작성이 만들 수 있는 최대 커밋 수 |
| partial-progress.max-failed-commits | partial-progress.max-commits 값 | 부분 진행이 활성화되면 작업 실패 전에 허용되는 최대 실패 커밋 수 |
| use-starting-sequence-number | true | 새로 생성된 스냅샷의 시퀀스 번호 대신 컴팩션 시작 시점의 스냅샷 시퀀스 번호 사용 |
| rewrite-job-order | none | 값에 따라 재작성 작업 순서를 강제. rewrite-job-order=bytes-asc이면 가장 작은 작업 그룹을 먼저 재작성. rewrite-job-order=bytes-desc이면 가장 큰 작업 그룹을 먼저 재작성. rewrite-job-order=files-asc이면 파일이 가장 적은 작업 그룹을 먼저 재작성. rewrite-job-order=files-desc이면 파일이 가장 많은 작업 그룹을 먼저 재작성. rewrite-job-order=none이면 계획된 순서대로 작업 그룹 재작성 (특정 순서 없음). |
| target-file-size-bytes | 536870912 (512 MB, 테이블 속성의 write.target-file-size-bytes 기본값) | 목표 출력 파일 크기 |
| min-file-size-bytes | 목표 파일 크기의 75% | 이 임계값 미만의 파일은 다른 기준과 무관하게 재작성 대상으로 고려 |
| max-file-size-bytes | 목표 파일 크기의 180% | 이 임계값보다 큰 파일은 다른 기준과 무관하게 재작성 대상으로 고려 |
| min-input-files | 5 | 이 수 이상의 파일을 가진 파일 그룹은 다른 기준과 무관하게 재작성됨 (파일 그룹은 최소 2개 파일이 있어야 함) |
| rewrite-all | false | 다른 옵션을 재정의해 제공된 모든 파일의 재작성 강제 |
| max-file-group-size-bytes | 107374182400 (100GB) | 단일 파일 그룹에서 재작성해야 하는 가장 큰 데이터 양. 전체 재작성 연산은 파티셔닝을 기준으로, 파티션 안에서는 크기를 기준으로 파일 그룹으로 나뉘어요. 이는 클러스터의 리소스 제약으로 재작성할 수 없는 매우 큰 파티션의 재작성을 나누는 데 도움이 돼요. |
| delete-file-threshold | 2147483647 | 데이터 파일이 재작성 대상으로 고려되기 위해 연관돼야 하는 최소 삭제 수 |
| delete-ratio-threshold | 0.3 | 데이터 파일이 재작성 대상으로 고려되기 위해 연관돼야 하는 최소 삭제 비율 |
| output-spec-id | 현재 파티션 스펙 id | 출력 파티션 스펙의 식별자. 재작성 중 데이터가 출력 파티셔닝에 맞춰 재구성됨. |
| remove-dangling-deletes | false | 재작성 후 dangling position·equality 삭제 제거. 삭제 파일이 어떤 라이브 데이터 파일에도 적용되지 않으면 dangling으로 간주됨. 활성화하면 제거를 위한 추가 커밋이 생성됨. |
| max-files-to-rewrite | null | 이 옵션은 재작성될 적격 파일 수에 대한 상한을 설정해요. 이 옵션이 지정되지 않으면 모든 적격 파일이 재작성됨 |
정보 (Info)
Dangling 삭제 파일은 데이터 시퀀스 번호만을 기반으로 제거돼요. 이 작업은 삭제 조건이 어떤 데이터 파일과도 일치하지 않는 전역 equality 삭제나 유효하지 않은 equality 삭제, 또는 더 이상 어떤 라이브 데이터 파일과도 일치하지 않는 position 삭제가 포함된 position delete 파일에는 적용되지 않아요.
sort 전략 옵션 (Options for sort strategy)
| 이름 | 기본값 | 설명 |
|---|---|---|
| compression-factor | 1.0 | 스파크 정렬이 만드는 셔플 파티션 수와 결과적으로 생성되는 출력 파일 수는 이 파일 재작성기에 사용된 입력 데이터 파일의 크기를 기반으로 해요. 압축 때문에 디스크 파일 크기가 출력 파일 크기를 정확히 나타내지 않을 수 있어요. 이 파라미터는 실제 출력 데이터 크기를 추정하는 데 사용되는 파일 크기를 조정할 수 있게 해줘요. 1.0보다 큰 인자는 디스크 파일 크기를 기반으로 예상되는 것보다 더 많은 파일을 생성해요. 1.0보다 작은 값은 디스크 크기를 기반으로 예상되는 것보다 더 적은 파일을 만들어요. |
| shuffle-partitions-per-file | 1 | 각 출력 파일에 사용할 셔플 파티션 수. 아이스버그는 커스텀 coalesce 연산을 사용해 이 정렬된 파티션들을 다시 단일 정렬 파일로 이어 붙여요. |
zorder sort_order를 가진 sort 전략 옵션 (Options for sort strategy with zorder sort_order)
| 이름 | 기본값 | 설명 |
|---|---|---|
| var-length-contribution | 8 | 가변 길이 타입(String, Binary)의 입력 컬럼에서 고려되는 바이트 수 |
| max-output-size | 2147483647 | ZOrder 알고리즘에서 인터리브되는 바이트 양 |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| rewritten_data_files_count | int | 이 명령으로 재작성된 데이터 수 |
| added_data_files_count | int | 이 명령으로 쓰여진 새 데이터 파일 수 |
| rewritten_bytes_count | long | 이 명령으로 쓰여진 바이트 수 |
| failed_data_files_count | int | partial-progress.enabled가 true일 때 재작성에 실패한 데이터 파일 수 |
| removed_delete_files_count | int | 이 명령으로 제거된 삭제 파일 수 |
예시 (Examples)
작은 파일을 결합하고 테이블의 기본 쓰기 크기에 따라 큰 파일을 분할하기 위해 기본 재작성 알고리즘인 bin-packing을 사용해 db.sample 테이블의 데이터 파일을 재작성:
CALL catalog_name.system.rewrite_data_files('db.sample');
어떤 파일을 재작성할지 결정하는 데 bin-pack과 같은 기본값을 사용해 모든 데이터를 id와 name으로 정렬해서 db.sample 테이블의 데이터 파일을 재작성:
CALL catalog_name.system.rewrite_data_files(table => 'db.sample', strategy => 'sort', sort_order => 'id DESC NULLS LAST,name ASC NULLS FIRST');
어떤 파일을 재작성할지 결정하는 데 bin-pack과 같은 기본값을 사용해 c1과 c2 컬럼에 zOrdering해서 db.sample 테이블의 데이터 파일을 재작성:
CALL catalog_name.system.rewrite_data_files(table => 'db.sample', strategy => 'sort', sort_order => 'zorder(c1,c2)');
최소 2개 파일의 재작성이 필요한 모든 파티션에서 bin-pack 전략을 사용해 db.sample 테이블의 데이터 파일을 재작성한 뒤, dangling 삭제 파일을 제거:
CALL catalog_name.system.rewrite_data_files(table => 'db.sample', options => map('min-input-files', '2', 'remove-dangling-deletes', 'true'));
필터(id = 3 and name = "foo")와 일치하는 데이터를 포함할 수 있는 파일을 선택해서 db.sample 테이블의 데이터 파일을 재작성:
CALL catalog_name.system.rewrite_data_files(table => 'db.sample', where => 'id = 3 and name = "foo"');
rewrite_manifests
스캔 계획을 최적화하기 위해 테이블의 매니페스트를 재작성해요.
매니페스트의 데이터 파일은 파티션 스펙의 필드로 정렬돼요. 이 프로시저는 스파크 작업을 사용해 병렬로 실행돼요.
정보 (Info)
이 프로시저는 영향받는 테이블을 참조하는 모든 캐시된 스파크 계획을 무효화해요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 업데이트할 테이블 이름 |
| use_caching | boolean | 작업 중 스파크 캐싱 사용 (기본값: false). 캐싱을 활성화하면 실행자의 메모리 사용량이 늘 수 있어요. | |
| spec_id | int | 재작성할 매니페스트의 스펙 id (기본값: 현재 스펙 id) | |
| sort_by | array | 매니페스트를 클러스터링할 파티션 변환 이름 목록. 자주 쿼리되는 파티션 변환을 선택하면 불필요한 매니페스트를 건너뛰어 계획 시간을 줄일 수 있어요. 설정하지 않으면 매니페스트가 스펙 순서의 모든 파티션 변환으로 정렬됨 |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| rewritten_manifests_count | int | 이 명령으로 재작성된 매니페스트 수 |
| added_manifests_count | int | 이 명령으로 쓰여진 새 매니페스트 파일 수 |
예시 (Examples)
db.sample 테이블의 매니페스트를 재작성하고 매니페스트 파일을 테이블 파티셔닝에 맞춰 정렬:
CALL catalog_name.system.rewrite_manifests('db.sample');
db.sample 테이블의 파티션 스펙 1에서 매니페스트를 재작성:
CALL catalog_name.system.rewrite_manifests(table => 'db.sample', spec_id => 1);
db.sample 테이블의 매니페스트를 재작성하고 파티션 필드 category로 매니페스트 엔트리를 클러스터링. 쿼리가 category를 자주 필터링할 때 스캔 계획 성능을 개선할 수 있어요:
CALL catalog_name.system.rewrite_manifests(table => 'db.sample', sort_by => array('category'));
rewrite_position_delete_files
아이스버그는 position delete 파일을 재작성할 수 있는데, 이는 두 가지 목적을 제공해요.
- 마이너 컴팩션(Minor Compaction): 작은 position delete 파일을 더 큰 것으로 컴팩션. 이는 매니페스트 파일에 저장되는 메타데이터 크기를 줄이고 작은 삭제 파일을 여는 오버헤드를 줄여요.
- Dangling 삭제 제거(Remove Dangling Deletes): 더 이상 라이브가 아닌 데이터 파일을 참조하는 position delete 레코드를 걸러내요. rewrite_data_files 후, 재작성된 데이터 파일을 가리키는 position delete 레코드가 항상 제거 표시되지는 않아 테이블의 라이브 스냅샷 메타데이터에 추적된 채 남을 수 있어요. 이것을 'dangling delete' 문제라고 해요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 업데이트할 테이블 이름 |
| options | map | 프로시저에 사용할 옵션 | |
| where | string | 파일 필터링에 사용되는 문자열 조건 |
Dangling 삭제는 재작성 중에 항상 걸러져요.
옵션 (Options)
| 이름 | 기본값 | 설명 |
|---|---|---|
| max-concurrent-file-group-rewrites | 5 | 동시에 재작성할 최대 파일 그룹 수 |
| partial-progress.enabled | false | 전체 재작성이 완료되기 전에 파일 그룹 커밋 활성화 |
| partial-progress.max-commits | 10 | 부분 진행이 활성화되면 이 재작성이 만들 수 있는 최대 커밋 수 |
| rewrite-job-order | none | 값에 따라 재작성 작업 순서를 강제. rewrite-job-order=bytes-asc이면 가장 작은 작업 그룹을 먼저 재작성. rewrite-job-order=bytes-desc이면 가장 큰 작업 그룹을 먼저 재작성. rewrite-job-order=files-asc이면 파일이 가장 적은 작업 그룹을 먼저 재작성. rewrite-job-order=files-desc이면 파일이 가장 많은 작업 그룹을 먼저 재작성. rewrite-job-order=none이면 계획된 순서대로 작업 그룹 재작성 (특정 순서 없음). |
| target-file-size-bytes | 67108864 (64MB, 테이블 속성의 write.delete.target-file-size-bytes 기본값) | 목표 출력 파일 크기 |
| min-file-size-bytes | 목표 파일 크기의 75% | 이 임계값 미만의 파일은 다른 기준과 무관하게 재작성 대상으로 고려 |
| max-file-size-bytes | 목표 파일 크기의 180% | 이 임계값보다 큰 파일은 다른 기준과 무관하게 재작성 대상으로 고려 |
| min-input-files | 5 | 이 파일 수를 초과하는 파일 그룹은 다른 기준과 무관하게 재작성됨 |
| rewrite-all | false | 다른 옵션을 재정의해 제공된 모든 파일의 재작성 강제 |
| max-file-group-size-bytes | 107374182400 (100GB) | 단일 파일 그룹에서 재작성해야 하는 가장 큰 데이터 양. 전체 재작성 연산은 파티셔닝을 기준으로, 파티션 안에서는 크기를 기준으로 파일 그룹으로 나뉘어요. 이는 클러스터의 리소스 제약으로 재작성할 수 없는 매우 큰 파티션의 재작성을 나누는 데 도움이 돼요. |
| max-files-to-rewrite | null | 이 옵션은 재작성될 적격 파일 수에 대한 상한을 설정해요. 이 옵션이 지정되지 않으면 모든 적격 파일이 재작성됨 |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| rewritten_delete_files_count | int | 이 명령으로 제거된 삭제 파일 수 |
| added_delete_files_count | int | 이 명령으로 추가된 삭제 파일 수 |
| rewritten_bytes_count | long | 이 명령으로 제거된 삭제 파일 총 바이트 수 |
| added_bytes_count | long | 이 명령으로 추가된 모든 새 삭제 파일 총 바이트 수 |
예시 (Examples)
db.sample 테이블에서 position delete 파일을 재작성. 기본 재작성 기준에 맞는 position delete 파일을 선택하고 목표 크기 target-file-size-bytes의 새 파일을 만들어요. Dangling 삭제는 재작성된 삭제 파일에서 제거돼요.
CALL catalog_name.system.rewrite_position_delete_files('db.sample');
db.sample 테이블의 모든 position delete 파일을 재작성하고 목표 크기 target-file-size-bytes의 새 파일을 만들어요. Dangling 삭제는 재작성된 삭제 파일에서 제거돼요.
CALL catalog_name.system.rewrite_position_delete_files(table => 'db.sample', options => map('rewrite-all', 'true'));
db.sample 테이블에서 position delete 파일을 재작성. 크기 기준으로 2개 이상의 position delete 파일 재작성이 필요한 파티션에서 position delete 파일을 선택해요. Dangling 삭제는 재작성된 삭제 파일에서 제거돼요.
CALL catalog_name.system.rewrite_position_delete_files(table => 'db.sample', options => map('min-input-files','2'));
테이블 마이그레이션 (Table migration)
snapshot과 migrate 프로시저는 기존 Hive 또는 스파크 테이블을 테스트하고 아이스버그로 마이그레이션하는 데 도움을 줘요.
snapshot
소스 테이블을 변경하지 않고 테스트를 위한 가벼운 임시 테이블 복사본을 만들어요.
새로 만들어진 테이블은 소스 테이블에 영향을 주지 않고 변경하거나 쓸 수 있지만, 스냅샷은 원본 테이블의 데이터 파일을 사용해요.
스냅샷에 대해 insert나 overwrite를 실행하면 새 파일은 원본 테이블 위치가 아니라 스냅샷 테이블의 위치에 놓여요.
스냅샷 테이블 테스트를 마치면 DROP TABLE을 실행해서 정리해요.
정보 (Info)
스냅샷으로 만든 테이블은 데이터 파일의 단독 소유자가 아니기 때문에, 데이터 파일을 물리적으로 삭제하는 expire_snapshots 같은 작업이 금지돼요. 메타데이터에만 영향을 주는 아이스버그 삭제는 여전히 허용돼요. 또한 원본 데이터 파일에 영향을 주는 어떤 연산도 스냅샷의 무결성을 깨뜨려요. 원본 Hive 테이블에 대해 실행된 DELETE 문은 원본 데이터 파일을 제거하고 스냅샷 테이블은 더 이상 접근할 수 없게 돼요.
기존 테이블을 아이스버그 테이블로 교체하려면 migrate를 참고해주세요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| source_table | ✔️ | string | 스냅샷할 테이블 이름 |
| table | ✔️ | string | 만들 새 아이스버그 테이블 이름 |
| location | string | 새 테이블의 테이블 위치 (기본적으로 카탈로그에 위임) | |
| properties | map | 새로 만들어진 테이블에 추가할 속성 | |
| parallelism | int | 파일 읽기에 사용할 스레드 수 (기본값: 1) |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| imported_files_count | long | 새 테이블에 추가된 파일 수 |
예시 (Examples)
db.sample 테이블을 참조하는 db.snap이라는 격리된 아이스버그 테이블을 db.snap의 카탈로그 기본 위치에 만들기:
CALL catalog_name.system.snapshot('db.sample', 'db.snap');
수동으로 지정한 위치 /tmp/temptable/의 db.sample 테이블을 참조하는 db.snap이라는 격리된 아이스버그 테이블로 마이그레이션:
CALL catalog_name.system.snapshot('db.sample', 'db.snap', '/tmp/temptable/');
migrate
테이블을 소스의 데이터 파일로 로드된 아이스버그 테이블로 교체해요.
테이블 스키마, 파티셔닝, 속성, 위치가 소스 테이블에서 복사돼요.
테이블 파티션이 지원되지 않는 포맷을 사용하면 migrate가 실패해요. 지원되는 포맷은 Avro, Parquet, ORC예요. 테이블이 버킷화되면 버킷팅이 보존되지 않으므로 migrate도 실패해요. 기존 데이터 파일은 아이스버그 테이블의 메타데이터에 추가되고 원본 테이블 스키마로 만든 name-to-id 매핑을 사용해 읽을 수 있어요.
테스트 중에 원본 테이블을 그대로 두려면 snapshot을 사용해서 소스 데이터 파일과 스키마를 공유하는 새 임시 테이블을 만들어요.
기본적으로 원본 테이블은 table_BACKUP_ 이름으로 보존돼요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 마이그레이션할 테이블 이름 |
| properties | map | 새 아이스버그 테이블의 속성 | |
| drop_backup | boolean | true이면 원본 테이블이 백업으로 보존되지 않음 (기본값: false) | |
| backup_table_name | string | 백업으로 보존될 테이블 이름 (기본값: table_BACKUP_) | |
| parallelism | int | 파일 읽기에 사용할 스레드 수 (기본값: 1) |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| migrated_files_count | long | 아이스버그 테이블에 추가된 파일 수 |
예시 (Examples)
스파크 기본 카탈로그의 db.sample 테이블을 아이스버그 테이블로 마이그레이션하고 'foo'가 'bar'로 설정된 속성 추가:
CALL catalog_name.system.migrate('spark_catalog.db.sample', map('foo', 'bar'));
현재 카탈로그의 db.sample을 추가 속성 없이 아이스버그 테이블로 마이그레이션:
CALL catalog_name.system.migrate('db.sample');
add_files
Hive 또는 파일 기반 테이블의 파일을 주어진 아이스버그 테이블에 직접 추가하려고 시도해요. migrate나 snapshot과 달리 add_files는 특정 파티션이나 파티션들의 파일을 가져올 수 있고 새 아이스버그 테이블을 만들지 않아요. 이 명령은 새 파일의 메타데이터를 만들고 파일을 이동하지 않아요. 이 프로시저는 파일이 실제로 아이스버그 테이블의 스키마와 일치하는지 확인하기 위해 파일의 스키마를 분석하지 않아요. 완료되면 아이스버그 테이블은 이 파일들을 아이스버그가 소유한 파일 집합의 일부로 취급해요. 즉 이후의 expire_snapshot 호출이 추가된 파일을 물리적으로 삭제할 수 있게 돼요. migrate나 snapshot이 가능하면 이 방법을 사용하면 안 돼요.
경고 (Warning)
add_files 프로시저는 추가되는 각 파일의 Parquet 메타데이터를 한 번만 가져와요. 티어드 스토리지(예: Amazon S3 Intelligent-Tiering 스토리지 클래스)를 사용하면 기본 파일이 아카이브에서 검색되고 일정 기간 더 높은 티어에 유지돼요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 파일을 추가할 테이블 |
| source_table | ✔️ | string | 파일이 올 테이블. file_format.path 형태의 경로도 가능 |
| partition_filter | map | 가져올 소스 테이블의 파티션 맵 | |
| check_duplicate_files | boolean | 테이블에 이미 있는 파일이 추가되는 것을 방지할지 여부 (기본값: true) | |
| parallelism | int | 파일 읽기에 사용할 스레드 수 (기본값: 1) |
경고 : 스키마가 검증되지 않아, 다른 스키마의 파일을 아이스버그 테이블에 추가하면 문제가 발생해요.
경고 : 이 메서드로 추가된 파일은 아이스버그 연산으로 물리적으로 삭제될 수 있어요.
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| added_files_count | long | 이 명령으로 추가된 파일 수 |
| changed_partition_count | long | 이 명령으로 변경된 파티션 수 (알려진 경우) |
경고 (Warning)
changed_partition_count는 테이블 속성 compatibility.snapshot-id-inheritance.enabled가 true로 설정됐거나 테이블 포맷 버전이 1보다 크면 NULL이 돼요.
예시 (Examples)
세션 카탈로그에 등록된 Hive 또는 스파크 테이블인 db.src_table의 파일을 아이스버그 테이블 db.tbl에 추가해요. part_col_1이 A와 같은 파티션 안에 존재하는 파일만 추가해요.
CALL spark_catalog.system.add_files(
table => 'db.tbl',
source_table => 'db.src_tbl',
partition_filter => map('part_col_1', 'A')
);
path/to/table 위치의 parquet 파일 기반 테이블에서 아이스버그 테이블 db.tbl로 파일을 추가해요. 어떤 파티션에 속하든 모든 파일을 추가해요.
CALL spark_catalog.system.add_files(
table => 'db.tbl',
source_table => '`parquet`.`path/to/table`'
);
register_table
이미 존재하지만 해당하는 카탈로그 식별자가 없는 metadata.json 파일에 대한 카탈로그 항목을 만들어요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 등록할 테이블 |
| metadata_file | ✔️ | string | 새 카탈로그 식별자로 등록할 메타데이터 파일 |
경고 (Warning)
같은 metadata.json을 하나 이상의 카탈로그에 등록하면 누락된 업데이트, 데이터 손실, 테이블 손상으로 이어질 수 있어요. 이 프로시저는 테이블이 기존 카탈로그에 더 이상 등록되지 않았거나, 카탈로그 간에 테이블을 이동할 때만 사용해주세요.
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| current_snapshot_id | long | 새로 등록된 아이스버그 테이블의 현재 스냅샷 ID |
| total_records_count | long | 새로 등록된 아이스버그 테이블의 총 레코드 수 |
| total_data_files_count | long | 새로 등록된 아이스버그 테이블의 총 데이터 파일 수 |
예시 (Examples)
spark_catalog에 db.tbl로 새 테이블을 등록하고 path/to/metadata/file.json 메타데이터 파일을 가리키게:
CALL spark_catalog.system.register_table(
table => 'db.tbl',
metadata_file => 'path/to/metadata/file.json'
);
메타데이터 정보 (Metadata information)
ancestors_of
지정된 스냅샷의 부모들의 라이브 스냅샷 ID를 보고해요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 라이브 스냅샷 ID를 보고할 테이블 이름 |
| snapshot_id | long | 부모의 라이브 스냅샷 ID를 얻기 위해 지정된 스냅샷 사용 |
팁 : snapshot_id 사용
B로의 롤백과 C' -> D' 추가가 있는 스냅샷 이력을 주어진다고 하면: A -> B - > C -> D \ -> C' -> (D')
스냅샷 ID를 지정하지 않으면 A -> B -> C' -> D'를 반환하고, D의 스냅샷 ID를 인자로 제공하면 A -> B -> C -> D를 반환해요.
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| snapshot_id | long | 조상 스냅샷 id |
| timestamp | long | 스냅샷 생성 시간 |
예시 (Examples)
현재 스냅샷의 모든 스냅샷 조상을 가져오기 (기본):
CALL spark_catalog.system.ancestors_of('db.tbl');
특정 스냅샷으로 모든 스냅샷 조상을 가져오기:
CALL spark_catalog.system.ancestors_of('db.tbl', 1);
CALL spark_catalog.system.ancestors_of(snapshot_id => 1, table => 'db.tbl');
변경 데이터 캡처 (Change Data Capture)
create_changelog_view
주어진 테이블의 변경을 포함하는 뷰를 만들어요.
사용법 (Usage)
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | changelog의 소스 테이블 이름 |
| changelog_view | string | 만들 뷰의 이름 | |
| options | map | 사용할 스파크 읽기 옵션 맵 | |
| net_changes | boolean | 순(net) 변경을 출력할지 여부 (아래 참조). 기본값: false. compute_updates가 true이면 반드시 false여야 해요. | |
| compute_updates | boolean | pre/post update 이미지를 계산할지 여부 (아래 참조). identifier_columns가 제공되면 기본값 true, 그렇지 않으면 기본값 false. | |
| identifier_columns | array | 업데이트를 계산할 식별자 컬럼 목록. compute_updates 인자가 true로 설정되고 identifier_columns가 제공되지 않으면 테이블의 현재 식별자 필드가 사용돼요. |
자주 사용되는 스파크 읽기 옵션 목록은 다음과 같아요.
- start-snapshot-id: 배타적 시작 스냅샷 ID. 제공하지 않으면 테이블의 첫 스냅샷부터 포함해서 읽어요.
- end-snapshot-id: 포함적 끝 스냅샷 id, 기본값은 테이블의 현재 스냅샷.
- start-timestamp: 배타적 시작 타임스탬프. 제공하지 않으면 테이블의 첫 스냅샷부터 포함해서 읽어요.
- end-timestamp: 포함적 끝 타임스탬프, 기본값은 테이블의 현재 스냅샷.
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| changelog_view | string | 만들어진 changelog 뷰의 이름 |
예시 (Examples)
스냅샷 1(배타)과 2(포함) 사이에 일어난 변경을 기반으로 tbl_changes changelog 뷰를 만들어요.
CALL spark_catalog.system.create_changelog_view(
table => 'db.tbl',
options => map('start-snapshot-id','1','end-snapshot-id', '2')
);
타임스탬프 1678335750489(배타)와 1678992105265(포함) 사이에 일어난 변경을 기반으로 my_changelog_view changelog 뷰를 만들어요.
CALL spark_catalog.system.create_changelog_view(
table => 'db.tbl',
options => map('start-timestamp','1678335750489','end-timestamp', '1678992105265'),
changelog_view => 'my_changelog_view'
);
식별자 컬럼 id와 name을 기반으로 업데이트를 계산하는 changelog 뷰를 만들어요.
CALL spark_catalog.system.create_changelog_view(
table => 'db.tbl',
options => map('start-snapshot-id','1','end-snapshot-id', '2'),
identifier_columns => array('id', 'name')
);
changelog 뷰가 만들어지면 뷰를 쿼리해서 스냅샷 사이의 변경을 볼 수 있어요.
SELECT * FROM tbl_changes;
SELECT * FROM tbl_changes where _change_type = 'INSERT' AND id = 3 ORDER BY _change_ordinal;
changelog 뷰는 추적 중인 변경에 대한 추가 정보를 제공하는 CDC 메타데이터 컬럼을 포함한다는 점에 유의해주세요. 이 컬럼들은 다음과 같아요.
- _change_type: 변경의 유형. 다음 값 중 하나를 가져요: INSERT, DELETE, UPDATE_BEFORE, 또는 UPDATE_AFTER.
- _change_ordinal: 변경의 순서.
- _commit_snapshot_id: 변경이 발생한 스냅샷 ID.
다음은 대응 결과의 예시예요. 첫 스냅샷이 2개 레코드를 삽입했고 두 번째 스냅샷이 1개 레코드를 삭제했음을 보여줘요.
| id | name | _change_type | _change_ordinal | _commit_snapshot_id |
|---|---|---|---|---|
| 1 | Alice | INSERT | 0 | 5390529835796506035 |
| 2 | Bob | INSERT | 0 | 5390529835796506035 |
| 1 | Alice | DELETE | 1 | 8764748981452218370 |
순 변경 (Net Changes)
프로시저는 여러 스냅샷에 걸친 중간 변경을 제거하고 순 변경만 출력할 수 있어요. 다음은 순 변경을 계산하는 changelog 뷰를 만드는 예시예요.
CALL spark_catalog.system.create_changelog_view(
table => 'db.tbl',
options => map('end-snapshot-id', '87647489814522183702'),
net_changes => true
);
순 변경에서 위 changelog 뷰는 Alice가 첫 스냅샷에서 삽입되고 두 번째 스냅샷에서 삭제됐으므로 다음 행만 포함해요.
| id | name | _change_type | _change_ordinal | _commit_snapshot_id |
|---|---|---|---|---|
| 2 | Bob | INSERT | 0 | 5390529835796506035 |
캐리오버 행 (Carry-over Rows)
프로시저는 기본적으로 캐리오버 행을 제거해요. 캐리오버 행은 copy-on-write를 사용할 때 행 수준 연산(MERGE, UPDATE, DELETE)의 결과예요. 예를 들어 row1 (id=1, name='Alice')과 row2 (id=2, name='Bob')을 포함하는 파일이 있다고 해볼게요. row2의 copy-on-write 삭제는 이 파일을 지우고 row1을 새 파일에 보존해야 해요. changelog 테이블은 실제 테이블 변경이 아님에도 이 쌍을 다음 행 쌍으로 보고해요.
| id | name | _change_type |
|---|---|---|
| 1 | Alice | DELETE |
| 1 | Alice | INSERT |
캐리오버 행을 보려면 SparkChangelogTable을 다음과 같이 쿼리해요:
SELECT * FROM spark_catalog.db.tbl.changes;
Pre/Post 업데이트 이미지 (Pre/Post Update Images)
프로시저는 구성되면 pre/post 업데이트 이미지를 계산해요. Pre/post 업데이트 이미지는 삭제 행과 삽입 행의 쌍에서 변환돼요. 식별자 컬럼은 삽입과 삭제 레코드가 같은 행을 가리키는지 결정하는 데 사용돼요. 두 레코드가 식별 컬럼에 대해 같은 값을 공유하면 같은 행의 이전·이후 상태로 간주돼요. 테이블 스키마에서 식별자 필드를 설정하거나 프로시저 파라미터로 입력할 수 있어요.
다음 예시는 식별자 컬럼(id)이 있는 pre/post 업데이트 이미지 계산을 보여줘요. 여기서 같은 id를 가진 행 삭제와 삽입이 단일 업데이트 연산으로 취급돼요. 구체적으로 다음 행 쌍이 있다고 가정해요:
| id | name | _change_type |
|---|---|---|
| 3 | Robert | DELETE |
| 3 | Dan | INSERT |
이 경우 프로시저는 업데이트 전의 행을 UPDATE_BEFORE 이미지로, 업데이트 후의 행을 UPDATE_AFTER 이미지로 표시해서 다음 pre/post 업데이트 이미지를 만듭니다.
| id | name | _change_type |
|---|---|---|
| 3 | Robert | UPDATE_BEFORE |
| 3 | Dan | UPDATE_AFTER |
테이블 통계 (Table Statistics)
compute_table_stats
이 프로시저는 특정 테이블에 대한 고유 값 수(NDV) 통계를 계산해요. 기본적으로 모든 컬럼에 대해 테이블의 현재 스냅샷으로 통계가 계산돼요. 이 프로시저는 선택적으로 특정 스냅샷 및/또는 컬럼 하위 집합에 대해 통계를 계산하도록 구성할 수 있어요.
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 테이블 이름 |
| snapshot_id | string | 통계를 수집할 스냅샷의 Id | |
| columns | array | 통계를 수집할 컬럼 |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| statistics_file | string | 이 명령으로 만들어진 통계 파일의 경로 |
예시 (Examples)
테이블 my_table의 최신 스냅샷 통계 수집:
CALL catalog_name.system.compute_table_stats('my_table');
테이블 my_table의 스냅샷 snap1의 통계 수집:
CALL catalog_name.system.compute_table_stats(table => 'my_table', snapshot_id => 'snap1' );
테이블 my_table의 스냅샷 snap1에서 col1과 col2 컬럼의 통계 수집:
CALL catalog_name.system.compute_table_stats(table => 'my_table', snapshot_id => 'snap1', columns => array('col1', 'col2'));
파티션 통계 (Partition Statistics)
compute_partition_stats
이 프로시저는 PartitionStatisticsFile을 가진 마지막 스냅샷부터 주어진 스냅샷(지정하지 않으면 현재 스냅샷 사용)까지 파티션 통계를 증분으로 계산하고, 결합된 결과를 PartitionStatisticsFile에 써요. 이전 파티션 통계 파일이 없으면 전체 계산을 수행해요. 또한 PartitionStatisticsFile을 테이블 메타데이터에 등록해요.
| 인자 이름 | 필수? | 타입 | 설명 |
|---|---|---|---|
| table | ✔️ | string | 테이블 이름 |
| snapshot_id | string | 파티션 통계를 계산할 스냅샷의 Id. 기본값은 현재 스냅샷 id |
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| partition_statistics_file | string | 명령으로 만들어진 파티션 통계 파일의 경로 |
예시 (Examples)
테이블 my_table의 최신 스냅샷 파티션 통계 수집:
CALL catalog_name.system.compute_partition_stats('my_table');
테이블 my_table의 스냅샷 snap1의 파티션 통계 수집:
CALL catalog_name.system.compute_partition_stats(table => 'my_table', snapshot_id => 'snap1');
테이블 복제 (Table Replication)
rewrite_table_path 프로시저는 아이스버그 테이블을 다른 위치로 복사하기 위해 준비해요.
rewrite_table_path
아이스버그 테이블의 메타데이터 파일 복사본을 스테이징하는데, 여기서 모든 절대 경로 소스 프리픽스가 지정된 대상 프리픽스로 교체돼요. 이는 아이스버그 테이블을 새 위치로 완전히 또는 증분으로 복사하는 시작점이 될 수 있어요.
정보 (Info)
이 프로시저는 재작성된 메타데이터 파일만 스테이징하고 복사할 파일 목록을 준비해요. 실제 파일 복사는 이 프로시저에 포함되지 않아요.
| 인자 이름 | 필수? | 기본값 | 타입 | 설명 |
|---|---|---|---|---|
| table | ✔️ | string | 테이블 이름 | |
| source_prefix | ✔️ | string | 교체할 기존 프리픽스 | |
| target_prefix | ✔️ | string | source_prefix의 교체 프리픽스 | |
| start_version | 테이블 메타데이터 로그의 첫 metadata.json | string | 시간순으로 첫 번째로 재작성할 metadata.json의 이름 또는 경로 | |
| end_version | 테이블 메타데이터 로그의 마지막 metadata.json | string | 시간순으로 마지막으로 재작성할 metadata.json의 이름 또는 경로 | |
| staging_location | 테이블 메타데이터 디렉터리 아래의 새 디렉터리 | string | 새로 재작성된 메타데이터 파일의 출력 위치 | |
| create_file_list | true | boolean | 재작성된 메타데이터의 경로를 포함하는 파일 목록을 생성할지 여부 |
작동 모드 (Modes of operation)
- 전체 재작성(Full Rewrite): 전체 재작성은 도달 가능한 모든 메타데이터 파일(metadata.json, 매니페스트 리스트, 매니페스트, position delete 파일 포함)을 재작성하고 file_list_location에 도달 가능한 모든 파일을 반환해요. 이는 이 프로시저의 기본 작동 모드예요.
- 증분 재작성(Incremental Rewrite): 선택적으로 start_version과 end_version을 제공해서 범위를 증분 재작성으로 제한할 수 있어요. 증분 재작성은 start_version과 end_version 사이에 추가된 메타데이터 파일만 재작성하고, file_list_location에 이 범위에서 추가된 파일만 반환해요.
출력 (Output)
| 출력 이름 | 타입 | 설명 |
|---|---|---|
| latest_version | string | 이 프로시저가 재작성한 최신 메타데이터 파일 이름 |
| file_list_location | string | 소스에서 대상 경로로의 매핑을 포함하는 CSV 파일 경로 |
| rewritten_manifest_file_paths_count | int | 재작성된 경로를 가진 매니페스트 파일 수 |
| rewritten_delete_file_paths_count | int | 재작성된 경로를 가진 삭제 파일 수 |
파일 목록 (File List)
파일은 start_version과 end_version 사이에 테이블에 추가된 모든 파일의 복사 계획을 포함해요.
각 파일에 대해 다음을 지정해요.
- Source Path: 테이블의 원본 파일 경로, 또는 파일이 재작성된 경우 스테이징 위치
- Target Path: 교체 프리픽스를 가진 경로
다음 예시는 3개 파일의 복사 계획을 보여줘요.
sourcepath/datafile1.parquet,targetpath/datafile1.parquet
sourcepath/datafile2.parquet,targetpath/datafile2.parquet
stagingpath/manifest.avro,targetpath/manifest.avro
예시 (Examples)
이 예시는 my_table의 메타데이터 경로를 HDFS의 소스 위치에서 S3의 대상 위치로 완전히 재작성해요. 테이블 메타데이터 디렉터리 아래의 기본 스테이징 위치에 새 메타데이터 집합을 만들어요.
CALL catalog_name.system.rewrite_table_path(
table => 'db.my_table',
source_prefix => 'hdfs://nn:8020/path/to/source_table',
target_prefix => 's3a://bucket/prefix/db.db/my_table'
);
이 예시는 my_table의 메타데이터 경로를 v2.metadata.json과 v20.metadata.json 메타데이터 버전 사이에서 증분으로 재작성하고, 새 메타데이터 파일을 명시적 스테이징 위치에 써요.
CALL catalog_name.system.rewrite_table_path(
table => 'db.my_table',
source_prefix => 's3a://bucketOne/prefix/db.db/my_table',
target_prefix => 's3a://bucketTwo/prefix/db.db/my_table',
start_version => 'v2.metadata.json',
end_version => 'v20.metadata.json',
staging_location => 's3a://bucketStaging/my_table'
);
재작성이 완료되면 서드파티 도구(예: Distcp)가 새로 만든 메타데이터 파일과 데이터 파일을 대상 위치로 복사할 수 있어요.
마지막으로 register_table 프로시저로 복사된 테이블을 대상 위치의 카탈로그에 등록할 수 있어요.
경고 (Warning)
파티션 통계 파일이 있는 아이스버그 테이블은 현재 경로 재작성이 지원되지 않아요.