델타 레이크 테이블 마이그레이션

델타 레이크 테이블 마이그레이션 (Delta Lake Table Migration)

델타 레이크(Delta Lake)는 Parquet 파일 포맷을 지원하고 타임 트래블과 버전 관리 기능을 제공하는 테이블 포맷이에요. 이 문서에서는 델타 레이크 테이블에서 아이스버그로 데이터를 옮기는 방법을 알려드릴게요. 델타 레이크에서 마이그레이션할 때는 데이터의 이력을 유지하기 위해 모든 스냅샷을 옮기는 것이 일반적이에요.

출처: 문서

본문

델타 레이크는 Parquet 파일 포맷을 지원하고 타임 트래블과 버전 관리 기능을 제공하는 테이블 포맷이에요. 델타 레이크에서 아이스버그로 데이터를 마이그레이션할 때는 데이터의 이력을 유지하기 위해 모든 스냅샷을 마이그레이션하는 것이 일반적이에요.

현재 아이스버그는 델타 레이크에서 아이스버그 테이블로 마이그레이션하기 위해 스냅샷 테이블(Snapshot Table) 작업을 지원해요. 델타 레이크 테이블은 트랜잭션을 유지하므로, 사용 가능한 모든 트랜잭션이 순서대로 새 아이스버그 테이블에 트랜잭션으로 커밋돼요. 델타 레이크 테이블의 경우, 초기 마이그레이션 이후에 추가된 데이터 파일들은 해당하는 트랜잭션에 포함되고, 이후 파일 추가(Add Transaction) 작업으로 새 아이스버그 테이블에 추가돼요. 파일 추가(Add File) 작업의 변형인 파일 추가 트랜잭션(Add Transaction) 작업은 아직 개발 중이에요.

델타 레이크에서 아이스버그로의 마이그레이션 활성화 (Enabling Migration from Delta Lake to Iceberg)

iceberg-delta-lake 모듈은 스파크와 플링크 엔진 런타임에 번들되지 않아요. 델타 레이크 기능에서 마이그레이션을 활성화하려면 최소한 다음 의존성이 필요해요.

  • iceberg-delta-lake
  • delta-standalone-0.6.0
  • delta-storage-2.2.0

호환성 (Compatibilities)

이 모듈은 Delta Standalone:0.6.0으로 빌드·테스트됐고, 다음 프로토콜 버전을 가진 델타 레이크 테이블을 지원해요.

  • minReaderVersion: 1
  • minWriterVersion: 2

델타 레이크 프로토콜 버전에 대한 자세한 내용은 Delta Lake Table Protocol Versioning 문서를 참고해주세요.

API

iceberg-delta-lake 모듈은 DeltaLakeToIcebergMigrationActionsProvider라는 이름의 인터페이스를 제공해요. 이 인터페이스에는 델타 레이크를 아이스버그로 변환하는 데 도움이 되는 작업(actions)이 포함돼 있어요. 지원되는 작업은 다음과 같아요.

  • snapshotDeltaLakeTable: 기존 델타 레이크 테이블을 아이스버그 테이블로 스냅샷

기본 구현 (Default Implementation)

iceberg-delta-lake 모듈은 인터페이스의 기본 구현도 제공하며, 다음처럼 접근할 수 있어요.

DeltaLakeToIcebergMigrationActionsProvider defaultActions = DeltaLakeToIcebergMigrationActionsProvider.defaultActions()

델타 레이크 테이블을 아이스버그로 스냅샷 (Snapshot Delta Lake Table to Iceberg)

snapshotDeltaLakeTable 작업은 델타 레이크 테이블의 트랜잭션을 읽고, 같은 스키마와 파티셔닝을 가진 새 아이스버그 테이블로 하나의 아이스버그 트랜잭션 안에서 변환해요. 원본 델타 레이크 테이블은 그대로 유지돼요.

새로 만들어진 테이블은 소스 테이블에 영향을 주지 않고 변경하거나 쓸 수 있지만, 스냅샷은 원본 테이블의 데이터 파일을 사용해요. 기존 데이터 파일은 아이스버그 테이블의 메타데이터에 추가되고, 원본 테이블 스키마로 만든 name-to-id 매핑을 사용해 읽을 수 있어요.

스냅샷에서 insert나 overwrite를 실행하면 새 파일이 스냅샷 테이블의 위치에 놓여요. 이 위치는 기본적으로 소스 델타 레이크 테이블의 위치와 같아요. 사용자는 스냅샷 테이블에 다른 위치를 지정할 수도 있어요.

정보 (Info)

snapshotDeltaLakeTable로 만들어진 테이블은 데이터 파일의 단독 소유자가 아니기 때문에, 데이터 파일을 물리적으로 삭제하는 expire_snapshots 같은 작업이 금지돼요. 메타데이터에만 영향을 주는 아이스버그 삭제는 여전히 허용돼요. 또한 원본 데이터 파일에 영향을 주는 어떤 연산도 스냅샷의 무결성을 깨뜨려요. 원본 델타 레이크 테이블에 대해 실행된 DELETE 문은 원본 데이터 파일을 제거하고, snapshotDeltaLakeTable 테이블은 더 이상 그 파일들에 접근할 수 없게 돼요.

사용법 (Usage)

필수 입력 설정 방식 설명
소스 테이블 위치 인자 sourceTableLocation 소스 델타 레이크 테이블의 위치
새 아이스버그 테이블 식별자 구성 API로 새 아이스버그 테이블의 네임스페이스와 테이블 이름을 지정하는 식별자
아이스버그 카탈로그 구성 API icebergCatalog 새 아이스버그 테이블을 만드는 데 사용되는 카탈로그
하둡 구성 구성 API deltaLakeConfiguration 소스 델타 레이크 테이블을 읽는 데 사용되는 하둡 구성

자세한 사용법과 다른 선택적 구성은 SnapshotDeltaLakeTable API를 참고해주세요.

출력 (Output)

출력 이름 타입 설명
imported_files_count long 새 테이블에 추가된 파일의 개수

추가되는 테이블 속성 (Added Table Properties)

다음 테이블 속성이 기본으로 만들어질 아이스버그 테이블에 추가돼요.

속성 이름 설명
snapshot_source delta 테이블이 델타 레이크 테이블에서 스냅샷됐음을 나타내요
original_location 델타 레이크 테이블의 위치 원본 델타 레이크 테이블 위치의 절대 경로
schema.name-mapping.default 스키마에서 파생된 JSON 이름 매핑 델타 레이크 테이블의 데이터 파일을 읽는 데 사용되는 이름 매핑 문자열

예시 (Examples)

import org.apache.iceberg.catalog.TableIdentifier;
import org.apache.iceberg.catalog.Catalog;
import org.apache.hadoop.conf.Configuration;
import org.apache.iceberg.delta.DeltaLakeToIcebergMigrationActionsProvider;

String sourceDeltaLakeTableLocation = "s3://my-bucket/delta-table";
String destTableLocation = "s3://my-bucket/iceberg-table";
TableIdentifier destTableIdentifier = TableIdentifier.of("my_db", "my_table");
Catalog icebergCatalog = ...; // Iceberg Catalog fetched from engines like Spark or created via CatalogUtil.loadCatalog
Configuration hadoopConf = ...; // Hadoop Configuration fetched from engines like Spark and have proper file system configuration to access the Delta Lake table.

DeltaLakeToIcebergMigrationActionsProvider.defaultActions()
    .snapshotDeltaLakeTable(sourceDeltaLakeTableLocation)
    .as(destTableIdentifier)
    .icebergCatalog(icebergCatalog)
    .tableLocation(destTableLocation)
    .deltaLakeConfiguration(hadoopConf)
    .tableProperty("my_property", "my_value")
    .execute();

더 알아보기 (Learn more)