Dataflow를 사용하여 Delta Lake 변경 데이터 캡처 (CDC) 데이터를 레이크하우스로 마이그레이션

타겟 Lakehouse Apache Iceberg 테이블을 자주 업데이트되는 Delta Lake 소스 테이블과 동기화하는 것은 어려울 수 있습니다. 데이터 세트 변경사항을 캡처하기 위해 전체 테이블을 다시 로드하면 컴퓨팅 비용과 처리 지연 시간이 증가하고, 트랜잭션 로그를 처리하기 위해 맞춤 파이프라인을 빌드하면 복잡한 운영 오버헤드가 발생합니다.

Dataflow 예약 일괄 변경 데이터 캡처 (CDC) 파이프라인을 사용하면 Delta Lake 소스 테이블을 Lakehouse Apache Iceberg 테이블과 지속적으로 동기화할 수 있습니다. 전체 테이블을 다시 로드하는 대신 예약된 일괄 작업은 Delta Lake 트랜잭션 로그에서 증분 업데이트를 직접 추출하고 구성 가능한 일정에 따라 대상 테이블에 적용합니다.

예약된 일괄 CDC 파이프라인은 다음과 같은 기능을 제공합니다.

  • 비용 효율적인 증분 동기화: 변경된 데이터 (INSERT, UPDATE, DELETE 작업)만 주기적으로 처리하여 컴퓨팅 오버헤드와 데이터 지연 시간을 줄입니다.
  • 자동 변경사항 매핑: Delta Lake 변경 데이터 피드 로그를 파싱하고 변경 작업을 타겟 Lakehouse Iceberg 테이블에 매핑합니다.
  • 유연한 작업 예약: 비즈니스 요구사항에 맞는 일정 (예: 15분 간격)에 따라 일괄 모드 Dataflow 작업을 실행합니다.

지속적인 마이그레이션 워크플로

Delta Lake에서 Lakehouse로의 엔드 투 엔드 연속 마이그레이션을 설정하려면 다음 작업 순서를 완료하세요.

  1. 입력 Delta Lake 테이블에서 delta.enableChangeDataFeed 속성 (true)을 사용 설정합니다.
  2. delta.enableChangeDataFeed 속성을 사용 설정한 후 커밋 ID 또는 타임스탬프로 초기 테이블 마이그레이션을 실행합니다. 자세한 내용은 Dataflow를 사용하여 Delta Lake 테이블을 Lakehouse로 가져오기를 참고하세요.
  3. 워크로드 요구사항에 맞게 선택한 빈도(커밋 ID 또는 타임스탬프 범위)로 후속 데이터 업데이트를 위해 예약된 일괄 CDC 작업을 실행합니다.

시작하기 전에

예약된 일괄 CDC 마이그레이션을 설정하려면 다음이 필요합니다.

  1. Dataflow, BigQuery, Lakehouse API가 아직 사용 설정되지 않은 경우 사용 설정합니다.

    API 사용 설정에 필요한 역할

    API를 사용 설정하려면 serviceusage.services.enable 권한이 필요합니다. 프로젝트를 만든 경우 소유자 역할 (roles/owner)을 통해 이 권한이 이미 있을 수 있습니다. 그렇지 않은 경우 서비스 사용량 관리자 역할 (roles/serviceusage.serviceUsageAdmin)을 통해 이 권한을 얻을 수 있습니다. 역할을 부여하는 방법 알아보기

    API 사용 설정

  2. 리소스를 만드는 데 필요한 권한을 얻으려면 관리자에게 프로젝트에 필요한 Identity and Access Management (IAM) 역할을 부여해 달라고 요청하세요.

  3. Cloud Storage 버킷에 저장된 유효한 Delta Lake 테이블입니다. 표 디렉터리에는 데이터 파일과 _delta_log/ 트랜잭션 로그 디렉터리가 포함되어야 합니다. 또한 표는 다음 요구사항을 충족해야 합니다.

    • Delta Lake 테이블에서 delta.enableChangeDataFeed 속성이 사용 설정 (true)되어 있어야 합니다.
    • CDC 파이프라인에 구성된 시작 커밋 ID 또는 타임스탬프는 delta.enableChangeDataFeed 속성이 사용 설정된 후여야 합니다.
  4. 동기화된 데이터를 수신할 기존 Lakehouse Iceberg 카탈로그 타겟 네임스페이스가 없으면 CDC 파이프라인에서 자동으로 생성합니다. 타겟 Lakehouse 테이블이 없으면 동등성 열의 유효한 목록을 제공하는 경우 CDC 파이프라인에서 자동으로 생성합니다.

지원 및 제한 사항

Delta Lake에서 Lakehouse로의 예약된 일괄 CDC 마이그레이션에는 다음과 같은 고려사항이 있습니다.

  • SDK 요구사항: Apache Beam SDK 버전 2.77.0 이상이 필요합니다.
  • 대상 테이블 생성: 대상 Lakehouse 테이블 생성이 지원됩니다. 타겟 테이블이 없으면 CDC 파이프라인에서 동등성 열 (equality_columns)의 유효한 목록을 제공하는 경우 자동으로 생성합니다. 기존 Lakehouse 카탈로그만 필요하며, CDC 파이프라인에서 네임스페이스가 없는 경우 자동으로 생성합니다.
  • 스키마 변경: 자동 스키마 변경은 지원되지 않습니다. 소스 Delta Lake 테이블의 스키마가 변경되면 CDC 동기화 작업을 실행하기 전에 대상 Lakehouse 테이블 스키마를 수동으로 업데이트하여 일치시켜야 합니다.
  • 일괄 실행 모드: Delta Lake CDC는 일괄 모드로 작동합니다. 파이프라인은 지정된 시작 및 종료 커밋 경계 (커밋 ID 또는 타임스탬프) 사이에서 작업 실행당 한 번 트랜잭션 로그를 가져옵니다. 종료 커밋 바운드 (end_version 또는 end_timestamp)가 지정되지 않으면 파이프라인은 최신 커밋 버전까지 읽습니다.
  • 지원되는 소스 스토리지: 소스 데이터는 Cloud Storage (gs://)에 저장된 유효한 Delta Lake 테이블이어야 합니다. Amazon S3는 Delta Lake 테이블 소스에 지원되지 않습니다.
  • 카탈로그 기반 테이블: 카탈로그 기반 Delta Lake 테이블 (예: Unity Catalog)은 지원되지 않습니다.
  • 지원되지 않는 Delta Lake 버전: Delta Lake 버전 1.2.1 이하에서는 변경 데이터 피드 (CDF)가 지원되지 않으므로 지원되지 않습니다.
  • 자동 변경 데이터 피드: CDF 변경 로그 파일을 작성하지 않으므로 자동 변경 데이터 피드는 지원되지 않습니다.

프로그래매틱 방식으로 CDC 파이프라인 만들기

프로그래매틱 방식으로 Delta Lake CDC 일괄 마이그레이션 파이프라인을 빌드하고 실행하려면 Apache Beam 파이프라인에서 Managed.read 변환을 사용하세요.

자바

pom.xml 파일에 다음 종속 항목을 추가합니다.

<dependency>
  <groupId>org.apache.beam</groupId>
  <artifactId>beam-sdks-java-managed</artifactId>
  <version>${beam.version}</version>
</dependency>
<dependency>
  <groupId>org.apache.beam</groupId>
  <artifactId>beam-sdks-java-io-iceberg</artifactId>
  <version>${beam.version}</version>
</dependency>

다음 Java 스니펫은 Delta Lake에서 Lakehouse로 배치 CDC 파이프라인을 구성하고 실행하는 방법을 보여줍니다.

import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.managed.Managed;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.Row;

public class DeltaLakeCdcToLakehouse {
  public static void main(String[] args) {
    Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());

    // 1. Configure Delta Lake CDC reader
    Map<String, Object> readConfig = new HashMap<>();
    readConfig.put("table", "gs://BUCKET_NAME/TABLE_NAME");
    // Starting commit bound (either start_version or start_timestamp must be provided):
    readConfig.put("start_version", START_COMMIT_ID);
    // Optional ending commit bound (end_version or end_timestamp, matching the start bound type):
    // readConfig.put("end_version", END_COMMIT_ID);
    readConfig.put("include_metadata_columns", Arrays.asList("_commit_version", "_change_type"));

    // Read CDC events from Delta Lake
    PCollection<Row> cdcRows =
        pipeline.apply(
            "ReadFromDeltaCDC",
            Managed.read(Managed.DELTA_LAKE_CDC).withConfig(readConfig)).getSinglePCollection();

    // 2. Configure BigLake REST Catalog properties for Lakehouse Iceberg
    Map<String, String> catalogProps = new HashMap<>();
    catalogProps.put("type", "rest");
    catalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
    catalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
    catalogProps.put("header.x-goog-user-project", "PROJECT_ID");
    catalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
    catalogProps.put("rest.auth.type", "org.apache.iceberg.gcp.auth.GoogleAuthManager");

    // 3. Map Delta Lake change types to Iceberg CDC change types
    Map<String, String> changeTypeMap = new HashMap<>();
    changeTypeMap.put("insert", "INSERT");
    changeTypeMap.put("delete", "DELETE");
    changeTypeMap.put("update_preimage", "UPDATE_BEFORE");
    changeTypeMap.put("update_postimage", "UPDATE_AFTER");

    // 4. Configure Lakehouse Iceberg CDC writer
    Map<String, Object> writeConfig = new HashMap<>();
    writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
    writeConfig.put("catalog_name", "lakehouse");
    writeConfig.put("catalog_properties", catalogProps);
    writeConfig.put("mode", "merge-on-read");
    writeConfig.put("sequence_number_column", "_commit_version");
    writeConfig.put("equality_columns", Arrays.asList("PRIMARY_KEY_COLUMN"));
    // change_type_column and change_type_map are only required for Portable Runner
    writeConfig.put("change_type_column", "_change_type");
    writeConfig.put("change_type_map", changeTypeMap);

    // Write CDC rows to target Lakehouse Iceberg table
    cdcRows.apply(
        "WriteToLakehouse",
        Managed.write(Managed.ICEBERG).withConfig(writeConfig));

    pipeline.run();
  }
}

다음을 바꿉니다.

  • BUCKET_NAME: Delta Lake 테이블이 포함된 Cloud Storage 버킷의 이름입니다.
  • TABLE_NAME: 소스 Delta Lake 테이블 디렉터리의 이름입니다.
  • START_COMMIT_ID: Delta Lake CDC 이벤트를 읽기 위한 시작 커밋 버전 번호입니다. start_version 또는 start_timestamp를 제공해야 합니다. 시작 커밋 ID 또는 타임스탬프는 소스 테이블에서 delta.enableChangeDataFeed 속성이 사용 설정된 후여야 합니다.
  • END_COMMIT_ID: Delta Lake CDC 이벤트를 읽기 위한 종료 커밋 버전 번호입니다. 시작에 사용된 경계 유형(start_version 또는 start_timestamp)과 일치하는 종료 경계 (end_version 또는 end_timestamp)를 지정하는 것이 좋습니다. 종료 경계를 생략하면 파이프라인은 최신 커밋 버전까지 읽습니다.
  • WAREHOUSE_BUCKET: Lakehouse 카탈로그 웨어하우스로 사용되는 Cloud Storage 버킷의 이름입니다.
  • PROJECT_ID: Google Cloud 프로젝트 ID입니다.
  • TARGET_NAMESPACE: 대상 Lakehouse 테이블의 네임스페이스입니다.
  • TARGET_TABLE: 대상 Lakehouse 테이블의 이름입니다.
  • PRIMARY_KEY_COLUMN: CDC 업데이트 및 삭제를 위해 행을 식별하는 데 사용되는 기본 키 열 이름입니다.

작업 출력 검사

CDC 데이터가 Lakehouse 테이블에 성공적으로 병합되었는지 확인합니다.

  1. Google Cloud 콘솔에서 BigQuery Studio 페이지로 이동합니다.

    BigQuery로 이동

  2. 쿼리 편집기에서 SQL 쿼리를 실행하여 동기화된 데이터를 확인합니다.

    SELECT * FROM `PROJECT_ID`.`CATALOG`.`NAMESPACE`.`TABLE_NAME` LIMIT 10;
    

    다음을 바꿉니다.

    • PROJECT_ID: Google Cloud 프로젝트 ID입니다.
    • CATALOG: Lakehouse 카탈로그의 이름입니다.
    • NAMESPACE: Lakehouse 테이블의 네임스페이스입니다.
    • TABLE_NAME: 대상 Lakehouse 테이블의 이름입니다.
  3. 실행을 클릭하고 결과를 확인합니다.

다음 단계