플링크 쓰기

아이스버그는 아파치 플링크의 DataStream API와 Table API로 배치 및 스트리밍 쓰기를 지원해요. 이 문서에서는 Flink SQL과 DataStream API로 아이스버그 테이블에 데이터를 쓰는 방법을 알려드릴게요. INSERT INTO, INSERT OVERWRITE, UPSERT, 분포 모드(HASH/RANGE), Flink 싱크 메트릭, Sink V2 구현, 그리고 동적 아이스버그 싱크(Dynamic Iceberg Sink)까지 폭넓게 살펴볼게요.

출처: 문서

본문

아이스버그는 아파치 플링크의 DataStream API와 Table API로 배치 및 스트리밍 쓰기를 지원해요.

Flink 아이스버그 싱크는 정확히 한 번(exactly-once) 의미론을 보장해요.

SQL로 쓰기 (Writing with SQL)

아이스버그는 INSERT INTO와 INSERT OVERWRITE를 모두 지원해요.

INSERT INTO

Flink 스트리밍 작업으로 테이블에 새 데이터를 추가하려면 INSERT INTO를 사용해요.

INSERT INTO `hive_catalog`.`default`.`sample` VALUES (1, 'a');
INSERT INTO `hive_catalog`.`default`.`sample` SELECT id, data from other_kafka_table;

INSERT OVERWRITE

테이블의 데이터를 쿼리 결과로 바꾸려면 배치 작업에서 INSERT OVERWRITE를 사용해요(Flink 스트리밍 작업은 INSERT OVERWRITE를 지원하지 않아요). Overwrite는 아이스버그 테이블에서 원자적 연산이에요.

SELECT 쿼리가 만든 행을 가진 파티션이 교체돼요. 예를 들어:

INSERT OVERWRITE sample VALUES (1, 'a');

아이스버그는 select 값으로 주어진 파티션을 덮어쓰는 것도 지원해요.

INSERT OVERWRITE `hive_catalog`.`default`.`sample` PARTITION(data='a') SELECT 6;

파티셔닝된 아이스버그 테이블의 경우, PARTITION 절의 모든 파티션 컬럼에 값이 설정되면 정적 파티션(static partition)에 삽입하고, PARTITION 절의 일부 파티션 컬럼(모든 파티션 컬럼의 프리픽스 부분)에 값이 설정되면 쿼리 결과를 동적 파티션(dynamic partition)에 쓰는 것이에요. 파티셔닝되지 않은 아이스버그 테이블의 경우 INSERT OVERWRITE로 데이터가 완전히 덮어써져요.

UPSERT

아이스버그는 v2 테이블 포맷으로 데이터를 쓸 때 기본 키 기반의 UPSERT를 지원해요. upsert를 활성화하는 방법은 두 가지가 있어요.

  1. 테이블 수준 속성 write.upsert.enabled로 UPSERT 모드를 활성화. 다음은 테이블을 만들 때 테이블 속성을 설정하는 예시 SQL 문이에요. 나중에 설명할 쓰기 옵션으로 덮어쓰지 않는 한 이 테이블의 모든 쓰기 경로(배치 또는 스트리밍)에 적용돼요.
    CREATE TABLE `hive_catalog`.`default`.`sample` (
        `id` INT COMMENT 'unique id',
        `data` STRING NOT NULL,
        PRIMARY KEY(`id`) NOT ENFORCED
    ) with ('format-version'='2', 'write.upsert.enabled'='true');
    
  2. 쓰기 옵션에서 upsert-enabled를 사용해 UPSERT 모드를 활성화하면 테이블 수준 구성보다 더 유연해요. 테이블을 만들 때 여전히 v2 테이블 포맷을 사용하고 기본 키나 식별자(identifier) 필드를 지정해야 한다는 점에 유의해주세요.
    INSERT INTO tableName /*+ OPTIONS('upsert-enabled'='true') */ ...
    

정보 (Info)

OVERWRITE와 UPSERT 모드는 상호 배타적이라 동시에 활성화할 수 없어요. 파티션 테이블에서 UPSERT 모드를 사용할 때는 해당 파티션 필드의 소스 컬럼이 equality 필드에 포함되어야 해요. 예를 들어 파티션 필드가 days(ts)라면 ts가 equality 필드의 일부여야 해요.

DataStream으로 쓰기 (Writing with DataStream)

아이스버그는 다양한 DataStream 입력에서 아이스버그 테이블로 쓰는 것을 지원해요.

데이터 추가 (Appending data)

Flink는 DataStream와 DataStream을 싱크 아이스버그 테이블에 네이티브로 쓰는 것을 지원해요.

StreamExecutionEnvironment env = ...;

DataStream<RowData> input = ... ;
Configuration hadoopConf = new Configuration();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path", hadoopConf);

FlinkSink.forRowData(input)
    .tableLoader(tableLoader)
    .append();

env.execute("Test Iceberg DataStream");

데이터 덮어쓰기 (Overwrite data)

FlinkSink 빌더에서 overwrite 플래그를 설정해서 기존 아이스버그 테이블의 데이터를 덮어써요.

StreamExecutionEnvironment env = ...;

DataStream<RowData> input = ... ;
Configuration hadoopConf = new Configuration();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path", hadoopConf);

FlinkSink.forRowData(input)
    .tableLoader(tableLoader)
    .overwrite(true)
    .append();

env.execute("Test Iceberg DataStream");

데이터 업서트 (Upsert data)

FlinkSink 빌더에서 upsert 플래그를 설정해서 기존 아이스버그 테이블의 데이터를 업서트해요. 테이블은 v2 테이블 포맷을 사용하고 기본 키가 있어야 해요.

StreamExecutionEnvironment env = ...;

DataStream<RowData> input = ... ;
Configuration hadoopConf = new Configuration();
TableLoader tableLoader = TableLoader.fromHadoopTable("hdfs://nn:8020/warehouse/path", hadoopConf);

FlinkSink.forRowData(input)
    .tableLoader(tableLoader)
    .upsert(true)
    .append();

env.execute("Test Iceberg DataStream");

정보 (Info)

OVERWRITE와 UPSERT 모드는 상호 배타적이라 동시에 활성화할 수 없어요. 파티션 테이블에서 UPSERT 모드를 사용할 때는 해당 파티션 필드의 소스 컬럼이 equality 필드에 포함되어야 해요. 예를 들어 파티션 필드가 days(ts)라면 ts가 equality 필드의 일부여야 해요.

Avro GenericRecord로 쓰기 (Write with Avro GenericRecord)

Flink 아이스버그 싱크는 Avro GenericRecord를 Flink RowData로 변환하는 AvroGenericRecordToRowDataMapper를 제공해요. 매퍼를 사용해서 Avro GenericRecord DataStream을 아이스버그에 쓸 수 있어요.

flink-avro jar가 클래스패스에 포함돼 있는지 확인해주세요. 또한 iceberg-flink-runtime shaded bundle jar는 런타임 jar가 avro 패키지를 쉐이딩하기 때문에 사용할 수 없어요. 대신 non-shaded iceberg-flink jar를 사용해주세요.

DataStream<org.apache.avro.generic.GenericRecord> dataStream = ...;

Schema icebergSchema = table.schema();

// The Avro schema converted from Iceberg schema can't be used
// due to precision difference between how Iceberg schema (micro)
// and Flink AvroToRowDataConverters (milli) deal with time type.
// Instead, use the Avro schema defined directly.
// See AvroGenericRecordToRowDataMapper Javadoc for more details.
org.apache.avro.Schema avroSchema = AvroSchemaUtil.convert(icebergSchema, table.name());

GenericRecordAvroTypeInfo avroTypeInfo = new GenericRecordAvroTypeInfo(avroSchema);
RowType rowType = FlinkSchemaUtil.convert(icebergSchema);

FlinkSink.builderFor(
    dataStream,
    AvroGenericRecordToRowDataMapper.forAvroSchema(avroSchema),
    FlinkCompatibilityUtil.toTypeInfo(rowType))
  .table(table)
  .tableLoader(tableLoader)
  .append();

브랜치 쓰기 (Branch Writes)

아이스버그 테이블의 브랜치에 쓰는 것도 FlinkSink의 toBranch API로 지원돼요. 브랜치에 대한 자세한 내용은 branches 문서를 참고해주세요.

FlinkSink.forRowData(input)
    .tableLoader(tableLoader)
    .toBranch("audit-branch")
    .append();

메트릭 (Metrics)

Flink 아이스버그 싱크는 다음 Flink 메트릭을 제공해요.

병렬 작성자(parallel writer) 메트릭은 IcebergStreamWriter의 하위 그룹 아래에 추가돼요. 다음 키-값 태그를 가져야 해요.

  • table: 전체 테이블 이름 (예: iceberg.my_db.my_table)
  • subtask_index: 0부터 시작하는 작성자 서브태스크 인덱스
메트릭 이름 메트릭 타입 설명
lastFlushDurationMs Gauge 체크포인트 중 작성자 서브태스크가 파일을 플러시하고 업로드하는 데 걸리는 시간(밀리초)
flushedDataFiles Counter 플러시되고 업로드된 데이터 파일의 수
flushedDeleteFiles Counter 플러시되고 업로드된 삭제 파일의 수
flushedReferencedDataFiles Counter 플러시된 삭제 파일이 참조하는 데이터 파일의 수
dataFilesSizeHistogram Histogram 데이터 파일 크기의 히스토그램 분포(바이트)
deleteFilesSizeHistogram Histogram 삭제 파일 크기의 히스토그램 분포(바이트)

위의 Histogram 메트릭은 org.apache.flink:flink-metrics-dropwizard가 클래스패스에 필요해요. 이 아티팩트는 기본적으로 Flink에 포함되지 않아요. 히스토그램 메트릭을 보려면 이 아티팩트를 클래스패스에 추가해주세요. 없으면 히스토그램 메트릭만 누락되고, 다른 모든 메트릭 타입은 계속 게시돼요.

커미터(committer) 메트릭은 IcebergFilesCommitter의 하위 그룹 아래에 추가돼요. 다음 키-값 태그를 가져야 해요.

  • table: 전체 테이블 이름 (예: iceberg.my_db.my_table)
메트릭 이름 메트릭 타입 설명
lastCheckpointDurationMs Gauge 커미터 연산자가 자신의 상태를 체크포인트하는 데 걸리는 시간(밀리초)
lastCommitDurationMs Gauge 아이스버그 테이블 커밋이 걸리는 시간(밀리초)
committedDataFilesCount Counter 커밋된 데이터 파일의 수
committedDataFilesRecordCount Counter 커밋된 데이터 파일에 포함된 레코드의 수
committedDataFilesByteCount Counter 커밋된 데이터 파일에 포함된 바이트의 수
committedDeleteFilesCount Counter 커밋된 삭제 파일의 수
committedDeleteFilesRecordCount Counter 커밋된 삭제 파일에 포함된 레코드의 수
committedDeleteFilesByteCount Counter 커밋된 삭제 파일에 포함된 바이트의 수
elapsedSecondsSinceLastSuccessfulCommit Gauge 마지막 성공적인 아이스버그 커밋 이후 경과 시간(초)

elapsedSecondsSinceLastSuccessfulCommit은 실패했거나 누락된 아이스버그 커밋을 감지하기에 이상적인 알림 메트릭이에요.

  • Iceberg 커밋은 성공적인 Flink 체크포인트 후 notifyCheckpointComplete 콜백에서 발생해요. Flink 체크포인트는 성공하는데(어떤 이유로든) 아이스버그 커밋이 실패할 수 있어요.
  • notifyCheckpointComplete가 트리거되지 않을 수도 있어요(어떤 버그로 인해). 결과적으로 어떤 아이스버그 커밋도 시도되지 않을 수 있어요.

체크포인트 간격(그리고 예상 아이스버그 커밋 간격)이 5분이라면, elapsedSecondsSinceLastSuccessfulCommit > 60 minutes 같은 규칙으로 알림을 설정해서 지난 1시간의 실패했거나 누락된 아이스버그 커밋을 감지해요.

옵션 (Options)

쓰기 옵션 (Write options)

Flink 쓰기 옵션은 FlinkSink를 구성할 때 전달돼요.

FlinkSink.Builder builder = FlinkSink.forRow(dataStream, SimpleDataUtil.FLINK_SCHEMA)
    .table(table)
    .tableLoader(tableLoader)
    .set("write-format", "orc")
    .set(FlinkWriteOptions.OVERWRITE_MODE, "true");

Flink SQL의 경우 쓰기 옵션은 SQL 힌트로 전달할 수 있어요.

INSERT INTO tableName /*+ OPTIONS('upsert-enabled'='true') */
...

모든 옵션은 여기(write-options)에서 확인할 수 있어요.

분포 모드 (Distribution mode)

Flink 스트리밍 작성자는 HASH와 RANGE 분포 모드를 모두 지원해요. FlinkSink#Builder#distributionMode(DistributionMode) 또는 write-options를 통해 활성화할 수 있어요.

해시 분포 (Hash distribution)

HASH 분포는 파티션 키(파티션 테이블) 또는 equality 필드(비파티션 테이블)로 데이터를 셔플해요. 이는 단순히 Flink의 DataStream#keyBy를 이용해서 데이터를 분포시키는 것이에요.

HASH 분포에는 몇 가지 제한이 있어요.

  • 스큐된(skewed) 데이터를 잘 처리하지 못해요. 예를 들어 어떤 파티션은 다른 파티션보다 데이터가 훨씬 많을 수 있어요. PR 4228에서 입증된 것처럼 파티션 키나 equality 필드의 카디널리티가 낮으면 트래픽 분포가 불균형해질 수 있어요. 작성자 병렬도는 해시 키의 카디널리티로 제한돼요. 카디널리티가 10이면 최대 10개의 작성자 작업만 트래픽을 받아요. 더 높은 작성자 병렬도를 설정해도(트래픽 볼륨이 필요하더라도) 도움이 되지 않아요.

범위 분포 (실험적) (Range distribution (experimental))

RANGE 분포는 커스텀 범위 파티셔너(range partitioner)를 통해 파티션 키나 정렬 순서로 데이터를 셔플해요. 범위 분포는 트래픽 통계를 수집해서 범위 파티셔너가 트래픽을 작성자 작업에 고르게 분포시키도록 안내해요.

범위 분포는 범위 파티셔너를 통해서만 데이터를 셔플해요. 데이터 파일 내의 행이 정렬되지는 않는데, Flink 스트리밍 작성자는 아직 정렬을 지원하지 않아요.

사용 사례 (Use cases)

RANGE 분포는 파티셔닝됐거나 SortOrder가 정의된 아이스버그 테이블에 적용할 수 있어요. SortOrder가 없는 파티션 테이블의 경우 파티션 컬럼이 정렬 순서로 사용돼요. 테이블에 SortOrder가 명시적으로 정의돼 있으면 범위 파티셔너가 그것을 사용해요.

범위 분포는 스큐된 데이터를 처리할 수 있어요. 예를 들어:

  • 테이블이 이벤트 시간으로 파티셔닝됨. 전형적으로 최근 시간대는 데이터가 많고, 긴 꼬리(long-tail) 시간대는 데이터가 점점 더 적어짐.
  • 테이블이 국가 코드로 파티셔닝됨. 어떤 국가(예: US)는 트래픽이 훨씬 많고 작은 국가는 데이터가 훨씬 적음.
  • 테이블이 이벤트 타입으로 파티셔닝됨. 어떤 타입은 다른 타입보다 데이터가 훨씬 많음.

범위 분포는 비파티션 컬럼에서도 데이터를 클러스터링할 수 있어요. 예를 들어 테이블이 수집 시간 기준으로 매시간 파티셔닝됨. 쿼리는 종종 device_id나 country_code 같은 비파티션 컬럼에 조건을 포함해요. 테이블 SortOrder가 비파티션 컬럼으로 정의되면 범위 파티셔닝이 비파티션 컬럼에서 클러스터링해서 쿼리 성능을 개선해요.

트래픽 통계 (Traffic statistics)

통계는 모든 셔플 연산자 서브태스크에 의해 수집되고, 체크포인트 주기마다 코디네이터에 의해 집계돼요. 집계된 통계는 모든 서브태스크에 브로드캐스트되고 다음 체크포인트에서 범위 파티셔너에 적용돼요. 그래서 트래픽 분포 변화를 감지하고 새 통계를 범위 파티셔너에 적용하는 데 최대 두 번의 체크포인트 주기가 걸릴 수 있어요.

범위 분포는 저카디널리티(예: country_code) 또는 고카디널리티(예: device_id) 시나리오에서 동작할 수 있어요.

  • 저카디널리티 시나리오(수백수천)에서는 HashMap이 모든 키의 트래픽 분포를 추적해요. 새 정렬 키 값이 나타나면 범위 파티셔너는 새 키에 대한 트래픽 분포를 학습하기 전에 라운드로빈으로 작성자 작업에 분배해요. 고카디널리티 시나리오(수백만수십억)에서는 균일 무작위 샘플링(reservoir sampling)을 사용해 정렬 키 공간을 균등하게 나누는 범위 경계를 계산해요. 메모리 사용량과 네트워크 교환을 낮게 유지해요. 키 분포가 비교적 균등하면 reservoir sampling이 잘 동작해요. 단일 핫 키가 트래픽에서 불균형하게 큰 몫을 차지하면 균등 샘플링에 의한 범위 분할은 잘 동작하지 않을 수 있어요.

사용법 (Usage)

다음은 Java에서 범위 분포를 활성화하는 방법이에요. 선택적인 고급 구성 두 가지가 있어요. 대부분의 경우 기본값이 잘 동작해요. 자세한 내용은 write-options를 참고해주세요.

FlinkSink.forRowData(input)
    ...
    .distributionMode(DistributionMode.RANGE)
    .rangeDistributionStatisticsType(StatisticsType.Auto)
    .rangeDistributionSortKeyBaseWeight(0.0d)
    .append();

오버헤드 (Overhead)

데이터 셔플(해시 또는 범위)은 직렬화/역직렬화와 네트워크 I/O의 계산 오버헤드가 있어요. CPU 사용률의 일부 증가를 예상해요.

범위 분포는 데이터 분포 통계도 수집하고 집계해요. 이 역시 일부 CPU 오버헤드를 발생시켜요. 기본 통계 타입 Auto를 사용하면 메모리 오버헤드는 전형적으로 작아요. 키 카디널리티가 높을 때는 Map 통계 타입을 사용하지 마세요. 이는 상당한 메모리 사용량과 통계 집계를 위한 큰 네트워크 교환을 초래할 수 있어요.

참고 (Notes)

Flink 스트리밍 쓰기 작업은 스냅샷 요약을 사용해 마지막 커밋된 체크포인트 ID를 유지하고, 커밋되지 않은 데이터를 임시 파일로 저장해요. 따라서 스냅샷 만료와 고아 파일 삭제는 Flink 작업의 상태를 손상시킬 수 있어요. 이를 피하려면 Flink 작업이 만든 마지막 스냅샷(요약의 flink.job-id 속성으로 식별 가능)을 유지하고, 충분히 오래된 고아 파일만 삭제해야 해요.

Sink V2 기반 구현 (Sink V2 based implementation)

현재 기본인 FlinkSink 구현이 만들어질 당시, Flink Sink의 인터페이스는 아이스버그 테이블의 목적에 맞지 않는 몇 가지 제한이 있었어요. 이런 제한 때문에 FlinkSink는 DiscardingSink으로 끝나는 커스텀 StreamOperator 체인에 기반을 두고 있어요.

Flink SinkV2 인터페이스의 1.15 버전에서 이 인터페이스가 도입됐어요. 이 인터페이스는 iceberg-flink 모듈에서 사용할 수 있는 새 IcebergSink 구현에 사용돼요. 새 구현은 테이블 유지보수 같은 기능에 대한 향후 작업의 기반이에요. SinkV2 기반 구현은 현재 실험적 기능이므로 주의해서 사용해주세요.

SQL로 쓰기 (Writing with SQL)

SQL에서 SinkV2 기반 구현을 켜려면 이 구성 옵션을 설정해요.

SET table.exec.iceberg.use-v2-sink = true;

DataStream으로 쓰기 (Writing with DataStream)

SinkV2 기반 구현을 사용하려면 제공된 스니펫에서 FlinkSink를 IcebergSink로 바꿔주세요.

경고 (Warning)

이 구현들 사이에는 약간의 차이가 있어요:

  • RANGE 분포 모드는 아직 IcebergSink에서 사용할 수 없어요.
  • IcebergSink를 사용할 때는 uidPrefix 대신 uidSuffix를 사용해요.

Flink 동적 아이스버그 싱크(Dynamic Sink)는 다음을 허용해요.

  1. 원하는 수의 테이블에 쓰기 — 단일 싱크가 레코드를 여러 아이스버그 테이블로 동적으로 라우팅할 수 있어요.
  2. 동적 테이블 생성과 업데이트 — 사용자 정의 라우팅 로직에 따라 테이블이 생성되고 업데이트돼요.
  3. 동적 스키마와 파티션 진화 — 스트리밍 실행 중 테이블 스키마와 파티션 스펙이 업데이트돼요.

모든 구성은 DynamicRecord 클래스를 통해 제어되고, 요구사항이 바뀌어도 Flink 작업 재시작이 필요 없어요.

    DynamicIcebergSink.forInput(dataStream)
        .generator((inputRecord, out) -> out.collect(
                new DynamicRecord(
                        TableIdentifier.of("db", "table"),
                        "branch",
                        SCHEMA,
                        (RowData) inputRecord,
                        PartitionSpec.unpartitioned(),
                        DistributionMode.HASH,
                        2)))
        .catalogLoader(CatalogLoader.hive("hive", new Configuration(), Map.of()))
        .writeParallelism(10)
        .immediateTableUpdate(true)
        .append();

구성 예시 (Configuration Example)

DynamicIcebergSink.Builder<RowData> builder = DynamicIcebergSink.forInput(inputStream);

// Set common properties
builder
    .set("write.parquet.compression-codec", "gzip");

// Set Dynamic Sink specific options
builder
    .writeParallelism(4)
    .uidPrefix("dynamic-sink")
    .cacheMaxSize(500)
    .cacheRefreshMs(5000);

// Add generator and append sink
builder.generator(new CustomRecordGenerator());
builder.append();

동적 라우팅 구성 (Dynamic Routing Configuration)

동적 테이블 라우팅은 DynamicRecordGenerator 인터페이스를 구현해서 커스터마이즈할 수 있어요.

public class CustomRecordGenerator implements DynamicRecordGenerator<RowData> {
    @Override
    public DynamicRecord generate(RowData row) {
        DynamicRecord record = new DynamicRecord();
        // Set table name based on business logic
        TableIdentifier tableIdentifier = TableIdentifier.of(database, tableName);
        record.setTableIdentifier(tableIdentifier);
        record.setData(row);
        // Set the maximum number of parallel writers for a given table/branch/schema/spec
        record.writeParallelism(2);
        return record;
    }
}

// Set custom record generator when building the sink
DynamicIcebergSink.Builder<RowData> builder = DynamicIcebergSink.forInput(inputStream);
builder.generator(new CustomRecordGenerator());
// ... other config ...
builder.append();

사용자는 입력 레코드를 DynamicRecord로 변환하는 컨버터를 제공해야 해요. 모든 레코드에 대해 다음 정보(DynamicRecord)가 필요해요.

속성 설명
TableIdentifier 레코드가 쓰여질 대상 테이블
Branch 레코드를 쓸 대상 브랜치 (선택)
Schema 레코드의 스키마
Spec 레코드에 대한 예상 파티셔닝 스펙
RowData 쓰여질 실제 행 데이터
DistributionMode 레코드 쓰기의 분포 모드 (NONE, HASH 또는 null). null이면 레코드를 전혀 셔플하지 않아요.
Parallelism 주어진 테이블/브랜치/스키마/스펙(WriteTarget)에 대한 최대 병렬 작성자 수
UpsertMode 이 테이블의 write.upsert.enabled를 재정의 (선택)
EqualityFields 테이블의 equality 필드 (선택)

스키마 진화 (Schema Evolution)

동적 싱크는 DynamicRecord에 제공된 스키마를 기존 테이블 스키마와 일치시키려고 해요.

  • 기존 테이블 스키마 중 하나와 직접 일치하면, 그 테이블 스키마가 테이블 쓰기에 사용돼요.
  • 직접 일치하지 않으면 DynamicSink는 제공된 스키마를 테이블 스키마 중 하나와 일치하도록 적응시키려 해요. 예를 들어 테이블 스키마에 추가적인 선택(optional) 컬럼이 있으면 DynamicRecord를 통해 제공된 RowData에 null 값이 추가돼요.
  • 그렇지 않으면 아래에서 설명하는 제약 안에서 테이블 스키마를 입력 스키마와 일치하도록 진화시켜요.

동적 싱크는 테이블 메타데이터와 들어오는 스키마 둘 다에 대해 LRU 캐시를 유지하고, 크기와 시간 제약에 기반해 제거해요. DynamicRecord가 현재 테이블 스키마와 호환되지 않는 스키마를 포함하면 스키마 업데이트가 트리거돼요. 이 업데이트는 immediateTableUpdate 구성에 따라 즉시 또는 중앙 집중 실행자를 통해 발생할 수 있어요. 중앙 집중 업데이트는 카탈로그의 부하를 줄이지만 싱크에 백프레셔를 도입할 수 있어요.

지원되는 스키마 업데이트 (Supported schema updates)

  • 새 컬럼 추가
  • 기존 컬럼 타입 확장 (예: Integer → Long, Float → Double)
  • 필수 컬럼을 선택적으로 만들기
  • 컬럼 드롭 (기본적으로 비활성화)

컬럼 드롭은 제거된 필드는 데이터 손실 없이 쉽게 복원할 수 없으므로, 지연되거나 순서가 잘못된 데이터로 인한 문제를 방지하기 위해 기본적으로 비활성화돼 있어요.

컬럼 드롭을 허용하도록 옵트인할 수 있어요(아래 구성 옵션 참조). 컬럼이 드롭된 후에도 아이스버그가 모든 과거 테이블 스키마를 유지하므로 기술적으로는 여전히 그 컬럼에 데이터를 쓸 수 있어요. 하지만 일반 쿼리는 그 컬럼을 참조할 수 없어요. 필드가 새 스키마의 일부로 다시 나타나면 완전히 새로운 컬럼이 추가되는데, 이름 외에는 이전 컬럼과 공통점이 없어요. 즉 새 컬럼에 대한 쿼리는 이전 컬럼의 데이터를 절대 반환하지 않아요.

지원되지 않는 스키마 업데이트 (Unsupported schema updates)
  • 컬럼 이름 변경

이름 변경은 스키마 비교가 이름 기반이기 때문에 지원되지 않아요. 이름 변경을 해결하려면 추가 메타데이터나 힌트가 필요해요.

캐싱 (Caching)

관련된 두 가지 별개의 캐시가 있어요: 테이블 메타데이터 캐시와 입력 스키마 캐시.

  • 테이블 메타데이터 캐시는 스키마 정의와 파티션 스펙 같은 메타데이터를 보관해서 반복적인 카탈로그 조회를 줄여요. 크기는 cacheMaxSize 설정에 의해 결정돼요.
  • 입력 스키마 캐시는 테이블별로 들어오는 스키마와 그 호환성 해결 결과를 저장해요. 크기는 inputSchemasPerTableCacheMaxSize로 제어돼요.

캐시 적중률과 성능을 개선하려면 레코드 스키마가 변경되지 않았다면 같은 DynamicRecord.schema 인스턴스를 재사용해요.

동적 싱크 구성 (Dynamic Sink Configuration)

동적 아이스버그 Flink 싱크는 Builder 패턴으로 구성돼요. 주요 구성 메서드는 다음과 같아요.

메서드 설명
overwrite(boolean enabled) overwrite 모드 활성화
writeParallelism(int parallelism) 작성자 병렬도 설정
uidPrefix(String prefix) 연산자 UID 프리픽스 설정
snapshotProperties(Map<String, String> properties) 스냅샷 메타데이터 속성 설정
toBranch(String branch) 특정 브랜치에 쓰기
cacheMaxSize(int maxSize) 테이블 메타데이터용 캐시 크기 설정
cacheRefreshMs(long refreshMs) 캐시 갱신 간격 설정
inputSchemasPerTableCacheMaxSize(int size) 테이블당 캐시할 최대 입력 스키마 수 설정
immediateTableUpdate(boolean enabled) 테이블 메타데이터(스키마/파티션 스펙) 업데이트가 즉시 일어나는지 제어 (기본값: false)
set(String property, String value) 어떤 아이스버그 쓰기 속성이든 설정 (예: "write.format", "write.upsert.enabled"). 모든 옵션은 여기(write-options)에서 확인할 수 있어요
setAll(Map<String, String> properties) 여러 속성을 한 번에 설정
tableCreator(TableCreator creator) DynamicIcebergSink가 새 아이스버그 테이블을 만들 때, 테이블 이름을 기준으로 커스텀 테이블 속성과 위치를 설정해서 만들 방식을 재정의할 수 있게 해요
dropUnusedColumns(boolean enabled) 활성화하면 현재 테이블 스키마에서 입력 스키마에 포함되지 않은 모든 컬럼을 드롭해요 (위의 컬럼 드롭 주의사항 참조)

분포 모드 (Distribution Modes)

각 DynamicRecord에 설정된 DistributionMode는 그 레코드가 프로세서에서 작성자로 어떻게 라우팅되는지 제어해요.

모드 동작
NONE 레코드가 라운드로빈 방식(또는 설정된 경우 equality 필드에 따라)으로 작성자 서브태스크에 분산돼요.
HASH 레코드가 파티션 키(파티션 테이블) 또는 equality 필드(비파티션 테이블)로 분산돼요. 같은 파티션의 레코드가 같은 작성자 서브태스크에서 처리되도록 보장해요.
null 포워드(Forward) 모드: 분포를 완전히 건너뛰고 포워드 엣지를 통해 레코드를 직접 보내요 (아래 참조).

포워드 모드 (Forward Mode)

distributionMode 파라미터가 없는 DynamicRecord 생성자 오버로드를 사용하면 분포를 완전히 건너뛰어요. 이는 모든 파티션이 이미 많은 양의 데이터를 가져서 직렬화와 네트워크 셔플 비용이 과도한 고처리량 파이프라인을 위해 설계됐어요. 레코드가 포워드 엣지를 사용해 프로세서에서 작성자로 직접 보내져 Flink 연산자 체이닝을 가능하게 해요. 테이블 메타데이터 업데이트는 의도적으로 추가 데이터 셔플을 피하기 위해 전용 테이블 업데이트 연산자를 생략했으므로, 프로세서 안에서 항상 즉시 수행돼요(immediateTableUpdate 설정과 무관).

포워드 레코드와 일반 레코드는 같은 파이프라인에서 섞일 수 있어요. 프로세서는 레코드를 두 개의 별도 싱크 출력으로 라우팅해요.

  • 셔플 싱크: 셔플링 레코드를 받아요. 작성자에 도달하기 전에 일반 분포 토폴로지(해시/라운드로빈)를 통과해요.
  • 포워드 싱크: distributionMode가 없는 레코드를 받아요. 분포를 완전히 건너뛰고 프로세서에서 포워드 엣지를 통해 흘러가 Flink 연산자 체이닝을 허용해요. 셔플 오버헤드를 피하는 것이 중요한 고처리량 테이블에 적합해요. 싱크의 writeParallelism 구성은 이 경로에 적용되지 않아요.

경고 (Warning)

  1. 포워드 경로에서 스키마 변경은 레코드가 포워드 엣지를 통해 직통으로 통과해야 하므로 항상 즉시 적용돼요. 의도된 대용량 사용 사례에서 이로 인해 아이스버그 카탈로그에 많은 충돌 커밋이 발생하고 데이터 처리가 일시적으로 지연될 수 있어요. 새 스키마의 레코드를 게시하기 전에 스키마를 외부에서 업데이트하거나, 업스트림에서 새 스키마가 도입될 때 처리량의 일시적 중단을 계획하는 것을 고려해주세요.
  2. 포워드 경로는 분포를 완전히 건너뛰기 때문에, 레코드가 동적 아이스버그 싱크에 도달하기 전에 업스트림에서 데이터를 올바르게 분포시킬 책임은 사용자에게 있어요. 그렇지 않으면 쓰기가 불균형해질 수 있어요.

참고 (Notes)

  • 범위 분포 모드: 현재 동적 싱크는 RANGE 분포 모드를 지원하지 않아요. 설정하면 HASH로 폴백해요.
  • 속성 우선순위 참고: 테이블 속성과 싱크 속성 사이에 충돌이 발생하면 싱크 속성이 테이블 속성 구성을 재정의해요.
  • 테이블 포맷 버전 업그레이드: 동적 싱크는 동적 레코드로 테이블을 업그레이드하는 것을 지원하지 않아요. V2에서 V3로의 업그레이드가 진행되는 동안 작업이 실행 중이면 안 돼요.

더 알아보기 (Learn more)