타겟 레이크하우스 Apache Iceberg 테이블을 자주 업데이트되는 Apache Iceberg 소스 테이블과 동기화하는 것은 어려울 수 있습니다. 데이터 세트 변경사항을 캡처하기 위해 전체 테이블을 다시 로드하면 컴퓨팅 비용과 처리 지연 시간이 증가하고, 변경 데이터 캡처 (CDC) 로그를 처리하기 위해 맞춤 파이프라인을 빌드하면 복잡한 운영 오버헤드가 발생합니다.
Dataflow 일괄 또는 스트리밍 CDC 파이프라인을 사용하면 Apache Iceberg 소스 테이블을 Lakehouse Apache Iceberg 테이블과 지속적으로 또는 주기적으로 동기화할 수 있습니다. 전체 테이블을 다시 로드하는 대신 CDC 파이프라인은 소스 Apache Iceberg 변경 로그에서 증분 업데이트를 직접 추출하여 대상 테이블에 적용합니다.
CDC 파이프라인은 다음과 같은 기능을 제공합니다.
- 비용 효율적인 증분 동기화: 변경된 데이터(
INSERT,UPDATE,DELETE작업)만 처리하여 컴퓨팅 오버헤드와 데이터 지연 시간을 줄입니다. - 유연한 실행 모드: 비즈니스 요구사항에 맞게 일괄 모드와 스트리밍 모드를 모두 지원합니다.
CDC 파이프라인은 두 가지 실행 모드를 지원합니다.
- 일괄 CDC 모드: 지정된 시작 및 종료 스냅샷 경계 사이에서 소스 Apache Iceberg 테이블의 증분 업데이트를 주기적으로 읽습니다.
- 스트리밍 CDC 모드: 구성 가능한 빈도 (기본값 1분, 최소 1초)로 소스 Apache Iceberg 테이블에서 새 변경사항 커밋을 지속적으로 폴링하고 거의 실시간으로 Lakehouse 테이블에 적용합니다.
시작하기 전에
Apache Iceberg에서 Lakehouse로의 CDC 마이그레이션을 설정하려면 다음이 필요합니다.
Dataflow, BigQuery, Lakehouse API가 아직 사용 설정되지 않은 경우 사용 설정합니다.
API 사용 설정에 필요한 역할
API를 사용 설정하려면
serviceusage.services.enable권한이 필요합니다. 프로젝트를 만든 경우 소유자 역할 (roles/owner)을 통해 이 권한이 이미 있을 수 있습니다. 그렇지 않은 경우 서비스 사용량 관리자 역할 (roles/serviceusage.serviceUsageAdmin)을 통해 이 권한을 얻을 수 있습니다. 역할을 부여하는 방법 알아보기리소스를 만드는 데 필요한 권한을 얻으려면 관리자에게 프로젝트에 필요한 Identity and Access Management (IAM) 역할을 부여해 달라고 요청하세요.
Cloud Storage에 저장되거나 Iceberg REST 카탈로그 또는 metastore에 등록된 유효한 소스 Apache Iceberg 테이블
동기화된 데이터를 수신할 기존 Lakehouse Iceberg 카탈로그
타겟 네임스페이스가 없으면 CDC 파이프라인에서 자동으로 생성합니다. 타겟 Lakehouse 테이블이 없으면 동등성 열의 유효한 목록을 제공하는 경우 CDC 파이프라인에서 자동으로 생성합니다.
지원 및 제한 사항
Apache Iceberg에서 Lakehouse로의 CDC 마이그레이션에는 다음과 같은 고려사항이 있습니다.
- SDK 요구사항: Apache Beam SDK 버전 2.77.0 이상이 필요합니다.
- 대상 테이블 생성: 대상 Lakehouse 테이블 생성이 지원됩니다. 타겟 테이블이 없으면 유효한 동등성 열 목록 (
equality_columns)을 제공하는 경우 CDC 파이프라인에서 자동으로 생성합니다. 기존 Lakehouse 카탈로그만 필요하며, CDC 파이프라인에서 네임스페이스가 없는 경우 자동으로 생성합니다. - 스키마 변경: 자동 스키마 변경은 지원되지 않습니다. 소스 Apache Iceberg 테이블의 스키마가 변경되면 CDC 동기화 작업을 실행하기 전에 대상 Lakehouse 테이블 스키마를 수동으로 업데이트하여 일치시켜야 합니다.
- 실행 모드:
- 일괄 CDC: 지정된 시작 및 종료 스냅샷 경계 (스냅샷 ID 또는 타임스탬프) 사이에서 작업 실행당 한 번 트랜잭션 로그 변경사항을 가져옵니다.
- CDC 스트리밍: 새 커밋 로그가 있는지 소스 Apache Iceberg 테이블을 지속적으로 폴링합니다. 폴링 빈도는 구성할 수 있습니다 (기본값은 1분, 최소 1초).
Java를 사용하여 CDC 파이프라인 만들기
Apache Beam Java SDK를 사용하여 Apache Iceberg CDC 배치 또는 스트리밍 이전 파이프라인을 프로그래매틱 방식으로 빌드하고 실행하려면 Managed.ICEBERG_CDC와 함께 Managed.read 변환을 사용하세요.
다음 Java 예시에서는 배치 및 스트리밍 워크로드에 Apache Iceberg CDC 파이프라인을 구성하고 실행하는 방법을 보여줍니다.
일괄 CDC 파이프라인
다음 Java 예에서는 지정된 시작 및 종료 스냅샷 범위 내에서 Apache Iceberg CDC 이벤트를 읽는 일괄 파이프라인을 구성합니다.
자바
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 스니펫은 Iceberg 테이블에서 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.io.iceberg.cdc.IcebergCdcMetadataColumns;
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 IcebergBatchCdcToLakehouse {
public static void main(String[] args) {
Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());
// Configure source catalog properties
Map<String, String> sourceCatalogProps = new HashMap<>();
sourceCatalogProps.put("type", "SOURCE_CATALOG_TYPE");
sourceCatalogProps.put("warehouse", "SOURCE_WAREHOUSE_LOCATION");
// Add additional properties if required by your catalog type (for example, uri)
// Configure batch Iceberg CDC reader with snapshot bounds
Map<String, Object> readConfig = new HashMap<>();
readConfig.put("table", "SOURCE_NAMESPACE.SOURCE_TABLE");
readConfig.put("catalog_properties", sourceCatalogProps);
readConfig.put("start_snapshot_id", START_SNAPSHOT_ID);
readConfig.put("end_snapshot_id", END_SNAPSHOT_ID);
readConfig.put(
"include_metadata_columns",
Arrays.asList(
IcebergCdcMetadataColumns.CHANGE_TYPE,
IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER));
// Read batch CDC events from source Iceberg table
PCollection<Row> cdcRows =
pipeline.apply(
"ReadFromIcebergCDC",
Managed.read(Managed.ICEBERG_CDC).withConfig(readConfig)).getSinglePCollection();
// Configure destination Lakehouse REST Catalog properties for Lakehouse Iceberg
Map<String, String> destinationCatalogProps = new HashMap<>();
destinationCatalogProps.put("type", "rest");
destinationCatalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
destinationCatalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
destinationCatalogProps.put("header.x-goog-user-project", "PROJECT_ID");
destinationCatalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
destinationCatalogProps.put("rest.auth.type", "google");
// Configure Lakehouse Iceberg writer
Map<String, Object> writeConfig = new HashMap<>();
writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
writeConfig.put("catalog_properties", destinationCatalogProps);
writeConfig.put("mode", "merge-on-read");
writeConfig.put("change_type_column", IcebergCdcMetadataColumns.CHANGE_TYPE);
// Write CDC rows to target Lakehouse Iceberg table
cdcRows.apply(
"WriteToLakehouse",
Managed.write(Managed.ICEBERG).withConfig(writeConfig));
pipeline.run();
}
}
다음을 바꿉니다.
SOURCE_CATALOG_TYPE: 소스 Iceberg 카탈로그의 유형입니다 (예:hadoop또는hive).SOURCE_WAREHOUSE_LOCATION: 소스 Iceberg 카탈로그의 웨어하우스 위치입니다 (예:s3://source-warehouse또는gs://source-warehouse).SOURCE_NAMESPACE: 소스 Iceberg 테이블의 네임스페이스SOURCE_TABLE: 소스 Iceberg 테이블의 이름입니다.START_SNAPSHOT_ID: 일괄 CDC 읽기의 시작 스냅샷 ID입니다.END_SNAPSHOT_ID: 일괄 CDC 읽기의 종료 스냅샷 ID입니다.WAREHOUSE_BUCKET: Lakehouse 카탈로그 웨어하우스로 사용되는 Cloud Storage 버킷의 이름입니다.PROJECT_ID: Google Cloud 프로젝트 ID입니다.TARGET_NAMESPACE: 대상 Lakehouse 테이블의 네임스페이스입니다.TARGET_TABLE: 대상 Lakehouse 테이블의 이름입니다.
스트리밍 CDC 파이프라인
다음 Java 예시에서는 새 CDC 변경사항이 있는지 소스 Apache Iceberg 테이블을 지속적으로 폴링하는 스트리밍 파이프라인을 구성합니다.
자바
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 스니펫은 Iceberg 테이블에서 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.io.iceberg.cdc.IcebergCdcMetadataColumns;
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 IcebergStreamingCdcToLakehouse {
public static void main(String[] args) {
Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());
// Configure source catalog properties
Map<String, String> sourceCatalogProps = new HashMap<>();
sourceCatalogProps.put("type", "SOURCE_CATALOG_TYPE");
sourceCatalogProps.put("warehouse", "SOURCE_WAREHOUSE_LOCATION");
// Add additional properties if required by your catalog type (for example, uri)
// Configure streaming Iceberg CDC reader with continuous polling
Map<String, Object> readConfig = new HashMap<>();
readConfig.put("table", "SOURCE_NAMESPACE.SOURCE_TABLE");
readConfig.put("catalog_properties", sourceCatalogProps);
readConfig.put("streaming", true);
readConfig.put("poll_interval_seconds", POLL_INTERVAL_SECONDS);
readConfig.put(
"include_metadata_columns",
Arrays.asList(
IcebergCdcMetadataColumns.CHANGE_TYPE,
IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER));
// Read streaming CDC events from source Iceberg table
PCollection<Row> cdcRows =
pipeline.apply(
"ReadFromIcebergCDC",
Managed.read(Managed.ICEBERG_CDC).withConfig(readConfig)).getSinglePCollection();
// Configure destination BigLake REST Catalog properties for Lakehouse Iceberg
Map<String, String> destinationCatalogProps = new HashMap<>();
destinationCatalogProps.put("type", "rest");
destinationCatalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
destinationCatalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
destinationCatalogProps.put("header.x-goog-user-project", "PROJECT_ID");
destinationCatalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
destinationCatalogProps.put("rest.auth.type", "google");
// Configure Lakehouse Iceberg writer
Map<String, Object> writeConfig = new HashMap<>();
writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
writeConfig.put("catalog_properties", destinationCatalogProps);
writeConfig.put("mode", "merge-on-read");
writeConfig.put("change_type_column", IcebergCdcMetadataColumns.CHANGE_TYPE);
// Write CDC rows to target Lakehouse Iceberg table
cdcRows.apply(
"WriteToLakehouse",
Managed.write(Managed.ICEBERG).withConfig(writeConfig));
pipeline.run();
}
}
다음을 바꿉니다.
SOURCE_CATALOG_TYPE: 소스 Iceberg 카탈로그의 유형입니다 (예:hadoop또는hive).SOURCE_WAREHOUSE_LOCATION: 소스 Iceberg 카탈로그의 웨어하우스 위치입니다 (예:s3://source-warehouse또는gs://source-warehouse).SOURCE_NAMESPACE: 소스 Iceberg 테이블의 네임스페이스SOURCE_TABLE: 소스 Iceberg 테이블의 이름입니다.POLL_INTERVAL_SECONDS: 연속 폴링 빈도(초)(기본값60, 최소값1).WAREHOUSE_BUCKET: Lakehouse 카탈로그 웨어하우스로 사용되는 Cloud Storage 버킷의 이름입니다.PROJECT_ID: Google Cloud 프로젝트 ID입니다.TARGET_NAMESPACE: 대상 Lakehouse 테이블의 네임스페이스입니다.TARGET_TABLE: 대상 Lakehouse 테이블의 이름입니다.
작업 출력 검사
CDC 데이터가 Lakehouse 테이블에 성공적으로 병합되었는지 확인합니다.
Google Cloud 콘솔에서 BigQuery Studio 페이지로 이동합니다.
쿼리 편집기에서 SQL 쿼리를 실행하여 동기화된 데이터를 확인합니다.
SELECT * FROM `PROJECT_ID`.`CATALOG`.`NAMESPACE`.`TABLE_NAME` LIMIT 10;다음을 바꿉니다.
PROJECT_ID: Google Cloud 프로젝트 ID입니다.CATALOG: Lakehouse 카탈로그의 이름입니다.NAMESPACE: Lakehouse 테이블의 네임스페이스입니다.TABLE_NAME: 대상 Lakehouse 테이블의 이름입니다.
실행을 클릭하고 결과를 확인합니다.
다음 단계
- Iceberg 테이블을 레이크하우스로 가져오기에 대해 자세히 알아보세요.
- Lakehouse의 변경 데이터 캡처 수집에 대해 알아봅니다.
- Dataflow의 관리형 I/O에 대해 자세히 알아보세요.