플링크 구성

이 문서에서는 Flink에서 아이스버그 카탈로그를 구성하는 방법과 런타임 읽기·쓰기 옵션을 알려드릴게요. CREATE CATALOG로 카탈로그를 만들고 설정하는 방법부터, Flink SQL 힌트나 Flink 구성, 테이블 속성으로 전달하는 각종 읽기·쓰기 옵션의 기본값과 설명까지 표로 정리해드릴게요.

출처: 문서

본문

카탈로그 구성 (Catalog Configuration)

다음 쿼리를 실행해서 카탈로그를 만들고 이름을 붙여요(<catalog_name>은 카탈로그 이름으로, <config_key>=<config_value>는 카탈로그 구현 구성으로 바꿔주세요).

CREATE CATALOG <catalog_name> WITH (
  'type'='iceberg',
  `<config_key>`=`<config_value>`
); 

다음 속성들은 전역적으로 설정할 수 있고 특정 카탈로그 구현에 국한되지 않아요.

속성 필수 설명
type ✔️ iceberg iceberg여야 해요.
catalog-type hive, hadoop, rest, glue, jdbc 또는 nessie 기본 아이스버그 카탈로그 구현, HiveCatalog, HadoopCatalog, RESTCatalog, GlueCatalog, JdbcCatalog, NessieCatalog 또는 catalog-impl로 커스텀 카탈로그 구현을 사용할 때는 설정하지 않음
catalog-impl 커스텀 카탈로그 구현의 완전한 클래스 이름. catalog-type이 설정되지 않으면 반드시 설정해야 해요.
property-version 속성 버전을 설명하는 버전 번호. 속성 형식이 바뀔 때 하위 호환을 위해 사용할 수 있어요. 현재 속성 버전은 1이에요.
cache-enabled true 또는 false 카탈로그 캐시를 활성화할지 여부, 기본값은 true.
cache.expiration-interval-ms 카탈로그 항목이 로컬에 캐시되는 시간(밀리초); -1 같은 음수 값은 만료를 비활성화하고, 0은 설정할 수 없어요. 기본값은 -1.

Hive 카탈로그를 사용할 때 다음 속성들을 설정할 수 있어요.

속성 필수 설명
uri ✔️ Hive metastore의 thrift URI.
clients Hive metastore 클라이언트 풀 크기, 기본값은 2.
warehouse Hive warehouse 위치. hive-conf-dir을 설정해서 hive-site.xml 구성 파일이 들어있는 위치를 지정하지 않았거나 올바른 hive-site.xml을 classpath에 추가하지 않았다면 사용자가 이 경로를 지정해야 해요.
hive-conf-dir 커스텀 Hive 구성 값을 제공하는 데 사용될 hive-site.xml 구성 파일이 들어있는 디렉터리의 경로. 아이스버그 카탈로그를 만들 때 hive-conf-dir와 warehouse를 모두 설정하면 /hive-site.xml(또는 classpath의 hive 구성 파일)의 hive.metastore.warehouse.dir 값이 warehouse 값으로 덮어써져요.
hadoop-conf-dir 커스텀 Hadoop 구성 값을 제공하는 데 사용될 core-site.xml과 hdfs-site.xml 구성 파일이 들어있는 디렉터리의 경로.

Hadoop 카탈로그를 사용할 때 다음 속성들을 설정할 수 있어요.

속성 필수 설명
warehouse ✔️ 메타데이터 파일과 데이터 파일을 저장할 HDFS 디렉터리.

REST 카탈로그를 사용할 때 다음 속성들을 설정할 수 있어요.

속성 필수 설명
uri ✔️ REST 카탈로그의 URL.
credential OAuth2 클라이언트 자격증명 흐름에서 토큰과 교환할 자격증명.
token 서버와 상호작용하는 데 사용될 토큰.

런타임 구성 (Runtime configuration)

읽기 옵션 (Read options)

Flink 읽기 옵션은 Flink IcebergSource를 구성할 때 전달돼요.

IcebergSource.forRowData()
    .tableLoader(TableLoader.fromCatalog(...))
    .assignerFactory(new SimpleSplitAssignerFactory())
    .streaming(true)
    .streamingStartingStrategy(StreamingStartingStrategy.INCREMENTAL_FROM_SNAPSHOT_ID)
    .startSnapshotId(3821550127947089987L)
    .monitorInterval(Duration.ofMillis(10L)) // or .set("monitor-interval", "10s") \ set(FlinkReadOptions.MONITOR_INTERVAL, "10s")
    .build()

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

SELECT * FROM tableName /*+ OPTIONS('monitor-interval'='10s') */
...

옵션은 Flink 구성을 통해서도 전달할 수 있고, 이 경우 현재 세션에 적용돼요. 모든 옵션이 이 모드를 지원하는 것은 아니라는 점에 유의해주세요.

env.getConfig()
    .getConfiguration()
    .set(FlinkReadOptions.SPLIT_FILE_OPEN_COST_OPTION, 1000L);
...

읽기 옵션이 가장 높은 우선순위를 갖고, 그 다음 Flink 구성, 그다음 테이블 속성이에요.

읽기 옵션 Flink 구성 테이블 속성 기본값 설명
snapshot-id N/A N/A null 배치 모드에서 타임 트래블용. 지정된 snapshot-id의 데이터를 읽어요.
case-sensitive connector.iceberg.case-sensitive N/A false true이면 컬럼 이름을 대소문자 구분 방식으로 매칭해요.
as-of-timestamp N/A N/A null 배치 모드에서 타임 트래블용. 주어진 시간(밀리초)을 기준으로 가장 최근 스냅샷의 데이터를 읽어요.
starting-strategy connector.iceberg.starting-strategy N/A INCREMENTAL_FROM_LATEST_SNAPSHOT 스트리밍 실행의 시작 전략. TABLE_SCAN_THEN_INCREMENTAL: 일반 테이블 스캔을 수행한 후 증분 모드로 전환. 증분 모드는 현재 스냅샷(배타적)에서 시작해요. INCREMENTAL_FROM_LATEST_SNAPSHOT: 최신 스냅샷(포함)부터 증분 모드 시작. 빈 테이블이면 모든 미래 append 스냅샷이 발견돼요. INCREMENTAL_FROM_LATEST_SNAPSHOT_EXCLUSIVE: 최신 스냅샷(배타적)부터 증분 모드 시작. 빈 테이블이면 모든 미래 append 스냅샷이 발견돼요. INCREMENTAL_FROM_EARLIEST_SNAPSHOT: 가장 이른 스냅샷(포함)부터 증분 모드 시작. 빈 테이블이면 모든 미래 append 스냅샷이 발견돼요. INCREMENTAL_FROM_SNAPSHOT_ID: 특정 id의 스냅샷(포함)부터 증분 모드 시작. INCREMENTAL_FROM_SNAPSHOT_TIMESTAMP: 특정 타임스탬프의 스냅샷(포함)부터 증분 모드 시작. 타임스탬프가 두 스냅샷 사이에 있으면 타임스탬프 이후의 스냅샷부터 시작해요. FIP27 Source 전용.
start-snapshot-timestamp N/A N/A null 주어진 시간(밀리초)을 기준으로 가장 최근 스냅샷의 데이터를 읽기 시작해요.
start-snapshot-id N/A N/A null 지정된 snapshot-id의 데이터를 읽기 시작해요.
end-snapshot-id N/A N/A 최신 스냅샷 id 끝 스냅샷을 지정해요.
branch N/A N/A main 배치 모드에서 읽을 브랜치를 지정해요
tag N/A N/A null 배치 모드에서 읽을 태그를 지정해요
start-tag N/A N/A null 증분 읽기의 시작 태그를 지정해요
end-tag N/A N/A null 증분 읽기의 끝 태그를 지정해요
split-size connector.iceberg.split-size read.split.target-size 128 MB 입력 스플릿을 결합할 때의 목표 크기.
split-lookback connector.iceberg.split-file-open-cost read.split.planning-lookback 10 입력 스플릿을 결합할 때 고려할 빈의 개수.
split-file-open-cost connector.iceberg.split-file-open-cost read.split.open-file-cost 4MB 파일을 여는 데 드는 예상 비용. 스플릿 결합 시 최소 가중치로 사용돼요.
streaming connector.iceberg.streaming N/A false 현재 작업이 스트리밍 모드로 실행되는지 배치 모드로 실행되는지 설정해요.
monitor-interval connector.iceberg.monitor-interval N/A 60s 새 스냅샷에서 스플릿을 발견하는 모니터링 간격. 스트리밍 읽기에만 적용돼요.
include-column-stats connector.iceberg.include-column-stats N/A false 각 데이터 파일과 함께 컬럼 통계를 로드하는 새 스캔을 만들어요. 컬럼 통계에는 값 개수, null 값 개수, 하한값, 상한값이 포함돼요.
max-planning-snapshot-count connector.iceberg.max-planning-snapshot-count N/A Integer.MAX_VALUE 스플릿 열거당 제한되는 최대 스냅샷 수. 스트리밍 읽기에만 적용돼요.
limit connector.iceberg.limit N/A -1 출력 행 수 제한.
max-allowed-planning-failures connector.iceberg.max-allowed-planning-failures N/A 3 작업을 실패시키기 전에 스캔 계획에 허용되는 최대 연속 실패 횟수. -1로 설정하면 스캔 계획 실패로 작업이 실패하지 않아요.
watermark-column connector.iceberg.watermark-column N/A null 워터마크 생성을 위해 사용할 워터마크 컬럼을 지정해요. 이 옵션이 있으면 splitAssignerFactory가 OrderedSplitAssignerFactory로 재정의돼요.
watermark-column-time-unit connector.iceberg.watermark-column-time-unit N/A TimeUnit.MICROSECONDS 워터마크 생성에 사용할 워터마크 시간 단위를 지정해요. 가능한 값은 DAYS, HOURS, MINUTES, SECONDS, MILLISECONDS, MICROSECONDS, NANOSECONDS.

쓰기 옵션 (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') */
...
Flink 옵션 기본값 설명
write-format Table write.format.default 이 쓰기 연산에 사용할 파일 포맷; parquet, avro 또는 orc
target-file-size-bytes 테이블 속성 기준 이 테이블의 write.target-file-size-bytes를 재정의해요
upsert-enabled Table write.upsert.enabled 이 테이블의 write.upsert.enabled를 재정의해요
overwrite-enabled false 테이블 데이터를 덮어써요. UPSERT 데이터 스트림을 사용하도록 구성할 때는 overwrite 모드를 활성화하면 안 돼요.
distribution-mode Table write.distribution-mode 이 테이블의 write.distribution-mode를 재정의해요. RANGE 분포는 실험 상태예요.
range-distribution-statistics-type Auto 범위 분포 데이터 통계 수집 타입: Map, Sketch, Auto. 자세한 내용은 여기 참조.
range-distribution-sort-key-base-weight 0.0 (double) 작성자 작업별 목표 트래픽 가중치에 대한 각 정렬 키의 기본 가중치. 자세한 내용은 여기 참조.
compression-codec Table write.(fileformat).compression-codec 이 쓰기에 대한 테이블의 압축 코덱을 재정의해요
compression-level Table write.(fileformat).compression-level Parquet 및 Avro 테이블의 이 쓰기에 대한 테이블 압축 수준을 재정의해요
compression-strategy Table write.orc.compression-strategy ORC 테이블의 이 쓰기에 대한 테이블 압축 전략을 재정의해요
write-parallelism 업스트림 연산자 병렬도 작성자 병렬도를 재정의해요
uid-suffix 테이블 속성 기준 이 테이블의 기본 IcebergSink에서 사용되는 uid 접미사를 재정의해요

범위 분포 통계 유형 (Range distribution statistics type)

구성 값은 Map, Sketch, Auto라는 enum 타입이에요.

  • Map: 모든 개별 키에 대한 정확한 샘플링 개수를 수집해요. 낮은 카디널리티(수백~수천) 시나리오에 사용해야 해요. Sketch: reservoir sampling을 통해 균일한 무작위 샘플링을 구성해요. 메모리 사용량이 낮게 유지되므로 높은 카디널리티(수백만) 시나리오에 잘 맞아요. Auto: Map 통계로 시작해요. 하지만 카디널리티가 임계값(현재 10,000)보다 높게 감지되면 통계가 자동으로 Sketch로 전환돼요.

범위 분포 정렬 키 기본 가중치 (Range distribution sort key base weight)

range-distribution-sort-key-base-weight: 0.0.

정렬 순서에 파티션 컬럼이 포함되면 각 정렬 키는 하나의 파티션과 데이터 파일에 매핑돼요. 이 상대 가중치는 트래픽이 낮은 정렬 키에 너무 많은 작은 파일을 배치하는 것을 피할 수 있어요. 이것은 각 정렬 키의 최소 가중치를 정의하는 double 값이에요. 0.02는 각 키가 작성자 작업별 목표 트래픽 가중치의 2%를 기본 가중치로 가진다는 뜻이에요.

예를 들어 싱크 아이스버그 테이블이 이벤트 시간으로 매일 파티셔닝된다고 해볼게요. 데이터 스트림에 지금부터 180일 전까지의 이벤트가 포함된다고 가정해요. 이벤트 시간으로 볼 때 서로 다른 날짜들에 걸친 트래픽 가중치 분포는 전형적으로 긴 꼬리(long tail) 패턴을 가져요. 현재 날짜가 가장 많은 트래픽을 포함하고, 오래된 날짜들(긴 꼬리)은 점점 더 적은 트래픽을 포함해요. 작성자 병렬도가 10이라고 가정해요. 180일 전체의 총 가중치는 10,000이에요. 작성자 작업별 목표 트래픽 가중치는 1,000이 돼요. 가장 오래된 150일의 가중치 합이 1,000이라고 가정해요. 일반적으로 범위 파티셔너는 가장 오래된 150일을 하나의 작성자 작업에 넣어요. 그 작성자 작업은 150개의 작은 파일(하루에 하나)을 쓰게 돼요. 150개의 열린 파일을 유지하면 많은 메모리를 소비할 수 있어요. 체크포인트 시점에 150개의 파일(아무리 작아도)을 플러시하고 업로드하는 것도 느릴 수 있어요. 이 구성이 0.02로 설정되면, 모든 쓰기 작업에 대해 각 정렬 키가 목표 가중치 1,000의 2%인 기본 가중치를 가진다는 뜻이에요. 이렇게 하면 아무리 작아도 하나의 작성자 작업에 50개 이상의 데이터 파일(하루에 하나)을 배치하지 않게 돼요.

이것은 저카디널리티 시나리오의 StatisticsType.Map에만 적용돼요. StatisticsType.Sketch 고카디널리티 정렬 컬럼의 경우 보통 파티션 컬럼으로 사용되지 않아요. 그렇지 않으면 쓰기 중 너무 많은 파티션과 작은 파일이 생성될 수 있기 때문이에요. Sketch 범위 파티셔너는 간단히 고카디널리티 키를 정렬된 범위로 나눠요.

더 알아보기 (Learn more)