Hive 읽기·쓰기
Hive 읽기·쓰기 (Hive Read & Write)
HiveCatalog를 사용하면 Apache Flink를 Apache Hive 테이블의 통합 BATCH·STREAM 처리를 위해 사용할 수 있어요. 이는 Flink를 Hive 배치 엔진의 더 우수한 성능 대안으로 쓰거나, Hive 테이블에 데이터를 지속적으로 읽고 쓰는 실시간 데이터 웨어하우징 애플리케이션을 구동하는 데 사용할 수 있음을 의미해요.
출처: 문서
본문
읽기 (Reading)
Flink는 BATCH와 STREAMING 모드 모두에서 Hive 데이터 읽기를 지원해요. BATCH 애플리케이션으로 실행하면 Flink는 질의가 실행되는 시점의 테이블 상태에 대해 질의를 실행해요. STREAMING 읽기는 테이블을 지속적으로 모니터링하고 새 데이터가 생기면 점진적으로 가져와요. Flink는 기본적으로 테이블을 유계(bounded)로 읽어요.
STREAMING 읽기는 파티션 테이블과 비파티션 테이블 모두 소비를 지원해요. 파티션 테이블의 경우 Flink가 새 파티션 생성을 모니터링하고 가능해지면 점진적으로 읽어요. 비파티션 테이블의 경우 Flink가 폴더의 새 파일 생성을 모니터링하고 새 파일을 점진적으로 읽어요.
스트리밍 소스 옵션:
| 키 (Key) | 기본값 (Default) | 타입 (Type) | 설명 (Description) |
|---|---|---|---|
| streaming-source.enable | false | Boolean | 스트리밍 소스를 활성화할지 여부. 참고: 각 파티션/파일이 원자적으로 쓰여졌는지 확인해야 해요. 그렇지 않으면 리더가 불완전한 데이터를 받을 수 있어요. |
| streaming-source.partition.include | all | String | 읽을 파티션을 설정하는 옵션. 지원값은 all과 latest이며, all은 모든 파티션 읽기, latest는 'streaming-source.partition.order' 순서의 최신 파티션 읽기를 뜻해요. latest는 스트리밍 hive 소스 테이블을 temporal table로 사용할 때만 동작해요. 기본값은 all이에요. 'streaming-source.enable'을 켜고 'streaming-source.partition.include'를 'latest'로 설정하면 Flink가 최신 hive 파티션의 temporal join을 지원해요. 동시에 파티션 비교 순서와 데이터 갱신 간격을 파티션 관련 옵션으로 지정할 수 있어요. |
| streaming-source.monitor-interval | None | Duration | 파티션/파일을 연속 모니터링하는 시간 간격. 참고: hive 스트리밍 읽기의 기본 간격은 '1 min'이고, hive 스트리밍 temporal join의 기본 간격은 '60 min'이에요. 이는 현재 hive 스트리밍 temporal join 구현에서 모든 TM이 Hive metastore를 방문해 metastore에 압력을 줄 수 있는 프레임워크 제약 때문이며, 향후 개선될 예정이에요. |
| streaming-source.partition-order | partition-name | String | 스트리밍 소스의 파티션 순서. create-time, partition-time, partition-name을 지원해요. create-time은 파티션/파일 생성 시간을 비교해요(이는 Hive metastore의 파티션 생성 시간이 아니라 파일 시스템의 폴더/파일 수정 시간이며, 폴더에 새 파일 추가 같은 파티션 폴더 갱신은 데이터 소비 방식에 영향을 줄 수 있어요). partition-time은 파티션 이름에서 추출한 시간을 비교해요. partition-name은 파티션 이름의 알파벳순을 비교해요. 비파티션 테이블의 경우 이 값은 항상 'create-time'이어야 해요. 기본값은 partition-name이에요. 이 옵션은 deprecated 옵션 'streaming-source.consume-order'와 동등해요. |
| streaming-source.consume-start-offset | None | String | 스트리밍 소비의 시작 오프셋. 오프셋을 해석하고 비교하는 방법은 순서에 따라 달라져요. create-time과 partition-time의 경우 타임스탬프 문자열(yyyy-[m]m-[d]d [hh:mm:ss])이어야 해요. partition-time의 경우 파티션 시간 추출기가 파티션에서 시간을 추출해요. partition-name의 경우 파티션 이름 문자열입니다(예: pt_year=2020/pt_mon=10/pt_day=01). |
SQL Hints를 사용해 Hive metastore의 정의를 바꾸지 않고 Hive 테이블에 구성을 적용할 수 있어요.
SELECT *
FROM hive_table
/*+ OPTIONS('streaming-source.enable'='true', 'streaming-source.consume-start-offset'='2020-05-20') */;
참고:
- 모니터 전략은 현재 location 경로의 모든 디렉터리/파일을 스캔하는 것이에요. 파티션이 많으면 성능이 저하될 수 있어요.
- 비파티션 테이블의 스트리밍 읽기는 각 파일이 대상 디렉터리에 원자적으로 쓰여져야 해요.
- 파티션 테이블의 스트리밍 읽기는 각 파티션이 Hive metastore 관점에서 원자적으로 추가되어야 해요. 그렇지 않으면 기존 파티션에 추가된 새 데이터가 소비돼요.
- 스트리밍 읽기는 Flink DDL의 워터마크 문법을 지원하지 않아요. 이 테이블들은 윈도우 연산자에 사용할 수 없어요.
Hive 뷰 읽기 (Reading Hive Views)
Flink는 Hive 정의 뷰를 읽을 수 있지만 일부 제약이 적용돼요:
- 뷰를 조회하기 전에 Hive 카탈로그를 현재 카탈로그로 설정해야 해요. 이는 Table API에서
tableEnv.useCatalog(...)또는 SQL Client에서USE CATALOG ...로 할 수 있어요. - Hive와 Flink SQL은 다른 구문(예: 다른 예약 키워드와 리터럴)을 가져요. 뷰의 질의가 Flink 문법과 호환되는지 확인해요.
읽기 시 벡터화 최적화 (Vectorized Optimization upon Read)
다음 조건이 충족되면 Flink가 자동으로 Hive 테이블의 벡터화된 읽기를 사용해요:
- 포맷: ORC 또는 Parquet
- List, Map, Struct, Union 같은 hive 타입의 복합 데이터 타입이 없는 컬럼
이 기능은 기본적으로 활성화돼요. 다음 구성으로 비활성화할 수 있어요.
table.exec.hive.fallback-mapred-reader=true
소스 병렬도 추론 (Source Parallelism Inference)
기본적으로 Flink는 파일 수와 각 파일의 블록 수를 기반으로 Hive 리더의 최적 병렬도를 추론해요. Flink는 병렬도 추론 정책을 유연하게 구성할 수 있게 해줘요. TableConfig에서 다음 파라미터를 구성할 수 있어요(이 파라미터들은 작업의 모든 소스에 영향을 줌):
| 키 (Key) | 기본값 (Default) | 타입 (Type) | 설명 (Description) |
|---|---|---|---|
| table.exec.hive.infer-source-parallelism.mode | dynamic | InferMode | splits 수에 따라 병렬도를 추론하는 hive 소스 병렬도 추론 모드를 선택하는 옵션. 'static'은 정적 추론으로 작업 생성 단계에서 소스 병렬도를 추론해요. 'dynamic'은 동적 추론으로 작업 실행 단계에서 병렬도를 추론하며 소스 병렬도를 더 정확히 추론할 수 있어요. 'none'은 병렬도 추론을 비활성화하는 것. deprecated 옵션 'table.exec.hive.infer-source-parallelism'의 영향을 여전히 받으며, 병렬도 추론을 활성화하려면 그 값이 true여야 해요. |
| table.exec.hive.infer-source-parallelism.max | 1000 | Integer | 소스 연산자의 최대 추론 병렬도를 설정해요. 기본값은 정적 병렬도 추론 모드에서만 유효해요. |
Hive 테이블 읽기 시 split 크기 튜닝 (Tuning Split Size While Reading Hive Table)
Hive 테이블을 읽을 때 데이터 파일이 splits로 열거되며, 각 split은 소스가 소비하는 데이터의 일부예요. Splits는 소스가 작업을 분배하고 데이터 읽기를 병렬화하는 단위예요. 사용자는 다음 구성으로 split 크기를 조정해 성능 튜닝을 할 수 있어요.
| 키 (Key) | 기본값 (Default) | 타입 (Type) | 설명 (Description) |
|---|---|---|---|
| table.exec.hive.split-max-size | 128mb | MemorySize | Hive 테이블을 읽을 때 하나의 split에 담을 최대 바이트 수(기본 128MB). |
| table.exec.hive.file-open-cost | 4mb | MemorySize | 파일을 여는 데 드는 추정 비용(기본 4MB). Hive 파일을 splits로 열거하는 데 사용돼요. 값을 과대평가하면 Flink가 Hive 데이터를 더 적은 splits로 묶는 경향이 있어, Hive 테이블에 작은 파일이 많을 때 유용해요. 값을 과소평가하면 Flink가 Hive 데이터를 더 많은 splits로 묶는 경향이 있어 병렬도 개선에 도움이 돼요. |
참고:
- split 크기를 튜닝하려면 Flink가 먼저 모든 파티션의 모든 파일 크기를 얻어야 해요. 파티션이 너무 많으면 시간이 걸릴 수 있는데, 작업 구성
table.exec.hive.calculate-partition-size.thread-num(기본 3)을 더 큰 값으로 설정해 더 많은 스레드로 과정을 가속화할 수 있어요. - 현재 이 split 크기 튜닝 구성은 ORC 포맷으로 저장된 Hive 테이블에만 동작해요.
테이블 통계 읽기 (Read Table Statistics)
Hive metastore에서 테이블 통계를 쓸 수 없으면 Flink는 테이블을 스캔해 통계를 얻어 더 나은 실행 계획을 만들려고 해요. 통계를 얻는 데 시간이 걸릴 수 있어요. 더 빠르게 얻으려면 table.exec.hive.read-statistics.thread-num을 사용해 테이블 스캔에 사용할 스레드 수를 구성할 수 있어요. 기본값은 현재 시스템의 사용 가능한 프로세서 수이며 구성 값은 0보다 커야 해요.
파티션 splits 로드 (Load Partition Splits)
Hive 파티션을 나누는 데 다중 스레드가 사용돼요. table.exec.hive.load-partition-splits.thread-num으로 스레드 수를 구성할 수 있어요. 기본값은 3이고 구성 값은 0보다 커야 해요.
하위 디렉터리가 있는 파티션 읽기 (Read Partition With Subdirectory)
어떤 경우 다른 테이블을 참조하는 외부 테이블을 만들 수 있는데, 파티션 컬럼이 참조된 테이블의 부분집합일 수 있어요. 예를 들어 day와 hour 파티션이 있는 파티션 테이블 fact_tz가 있다고 해봐요:
CREATE TABLE fact_tz(x int) PARTITIONED BY (day STRING, hour STRING);
그리고 day 파티션으로 fact_tz 테이블을 참조하는 외부 테이블 fact_daily가 있다고 해봐요:
CREATE EXTERNAL TABLE fact_daily(x int) PARTITIONED BY (ds STRING) LOCATION '/path/to/fact_tz';
그러면 외부 테이블 fact_daily를 읽을 때 테이블의 파티션 디렉터리에 하위 디렉터리(hour=1부터 hour=24까지)가 있을 거예요. 기본적으로 외부 테이블에 하위 디렉터리가 있는 파티션을 추가할 수 있어요. Flink SQL은 모든 하위 디렉터리를 재귀적으로 스캔해 모든 하위 디렉터리의 모든 데이터를 가져올 수 있어요.
ALTER TABLE fact_daily ADD PARTITION (ds='2022-07-07') location '/path/to/fact_tz/ds=2022-07-07';
작업 구성 table.exec.hive.read-partition-with-subdirectory.enabled(기본 true)를 false로 설정해 Flink가 하위 디렉터리를 읽지 못하게 할 수 있어요. 구성이 false이고 디렉터리에 파일이 없고 하위 디렉터리로만 구성되어 있으면 Flink는 java.io.IOException: Not a file: /path/to/data/* 예외로 실패해요.
Temporal Table Join
Hive 테이블을 temporal table로 사용할 수 있고, 스트림은 temporal join으로 Hive 테이블을 연관(correlate)시킬 수 있어요. temporal join에 대한 자세한 내용은 temporal join을 참고해요.
Flink는 processing-time temporal join Hive Table을 지원하며, processing-time temporal join은 항상 temporal table의 최신 버전과 조인해요. Flink는 파티션 테이블과 Hive 비파티션 테이블 모두의 temporal join을 지원하며, 파티션 테이블의 경우 Flink가 Hive 테이블의 최신 파티션을 자동으로 추적하는 것을 지원해요.
참고: Flink는 아직 event-time temporal join Hive table을 지원하지 않아요.
최신 파티션과의 Temporal Join (Temporal Join The Latest Partition)
시간이 지나며 변하는 파티션 테이블의 경우 유계 없는 스트림으로 읽을 수 있어요. 각 파티션이 한 버전의 완전한 데이터를 포함하면 파티션을 temporal table의 버전으로 간주할 수 있고, temporal table의 버전은 파티션의 데이터를 유지해요. Flink는 processing time temporal join에서 temporal table의 최신 파티션(버전)을 자동 추적하는 것을 지원하며, 최신 파티션(버전)은 'streaming-source.partition-order' 옵션으로 정의돼요. 이는 Flink 스트림 애플리케이션 작업에서 Hive 테이블을 디멘전 테이블로 사용하는 가장 흔한 사용 사례예요.
참고: 이 기능은 Flink STREAMING 모드에서만 지원돼요.
다음 데모는 전형적인 비즈니스 파이프라인을 보여줘요. 디멘전 테이블은 Hive에서 오고 배치 파이프라인 작업이나 Flink 작업이 하루에 한 번 갱신하며, Kafka 스트림은 실시간 온라인 비즈니스 데이터 또는 로그에서 오며 디멘전 테이블과 조인해 스트림을 보강해야 해요.
-- Assume the data in hive table is updated per day, every day contains the latest and complete dimension data
SET table.sql-dialect=hive;
CREATE TABLE dimension_table (
product_id STRING,
product_name STRING,
unit_price DECIMAL(10, 4),
pv_count BIGINT,
like_count BIGINT,
comment_count BIGINT,
update_time TIMESTAMP(3),
update_user STRING,
...
) PARTITIONED BY (pt_year STRING, pt_month STRING, pt_day STRING) TBLPROPERTIES (
-- using default partition-name order to load the latest partition every 12h (the most recommended and convenient way)
'streaming-source.enable' = 'true',
'streaming-source.partition.include' = 'latest',
'streaming-source.monitor-interval' = '12 h',
'streaming-source.partition-order' = 'partition-name', -- option with default value, can be ignored.
-- using partition file create-time order to load the latest partition every 12h
'streaming-source.enable' = 'true',
'streaming-source.partition.include' = 'latest',
'streaming-source.partition-order' = 'create-time',
'streaming-source.monitor-interval' = '12 h'
-- using partition-time order to load the latest partition every 12h
'streaming-source.enable' = 'true',
'streaming-source.partition.include' = 'latest',
'streaming-source.monitor-interval' = '12 h',
'streaming-source.partition-order' = 'partition-time',
'partition.time-extractor.kind' = 'default',
'partition.time-extractor.timestamp-pattern' = '$pt_year-$pt_month-$pt_day 00:00:00'
);
SET table.sql-dialect=default;
CREATE TABLE orders_table (
order_id STRING,
order_amount DOUBLE,
product_id STRING,
log_ts TIMESTAMP(3),
proctime as PROCTIME()
) WITH (...);
-- streaming sql, kafka temporal join a hive dimension table. Flink will automatically reload data from the
-- configured latest partition in the interval of 'streaming-source.monitor-interval'.
SELECT * FROM orders_table AS o
JOIN dimension_table FOR SYSTEM_TIME AS OF o.proctime AS dim
ON o.product_id = dim.product_id;
최신 테이블과의 Temporal Join (Temporal Join The Latest Table)
Hive 테이블은 유계 스트림으로 읽을 수 있어요. 이 경우 Hive 테이블은 질의 시점의 최신 버전만 추적할 수 있어요. 테이블의 최신 버전은 Hive 테이블의 모든 데이터를 유지해요. 최신 Hive 테이블과 temporal join을 수행할 때 Hive 테이블은 Slot 메모리에 캐시되고 스트림의 각 레코드는 키로 테이블과 조인되어 일치가 있는지 결정돼요. 최신 Hive 테이블을 temporal table로 사용하는 것은 추가 구성이 필요 없어요. 선택적으로 다음 프로퍼티로 Hive 테이블 캐시의 TTL을 구성할 수 있어요. 캐시가 만료된 후 Hive 테이블이 다시 스캔되어 최신 데이터를 로드해요.
| 키 (Key) | 기본값 (Default) | 타입 (Type) | 설명 (Description) |
|---|---|---|---|
| lookup.join.cache.ttl | 60 min | Duration | lookup join에서 빌드 테이블의 캐시 TTL(예: 10min). 기본 TTL은 60분이에요. 참고: 이 옵션은 유계 hive 테이블 소스를 사용할 때만 동작해요. 스트리밍 hive 소스를 temporal table로 사용한다면 'streaming-source.monitor-interval'로 데이터 갱신 간격을 구성해요. |
다음 데모는 Hive 테이블의 모든 데이터를 temporal table로 로드하는 것을 보여줘요.
-- Assume the data in hive table is overwrite by batch pipeline.
SET table.sql-dialect=hive;
CREATE TABLE dimension_table (
product_id STRING,
product_name STRING,
unit_price DECIMAL(10, 4),
pv_count BIGINT,
like_count BIGINT,
comment_count BIGINT,
update_time TIMESTAMP(3),
update_user STRING,
...
) TBLPROPERTIES (
'streaming-source.enable' = 'false', -- option with default value, can be ignored.
'streaming-source.partition.include' = 'all', -- option with default value, can be ignored.
'lookup.join.cache.ttl' = '12 h'
);
SET table.sql-dialect=default;
CREATE TABLE orders_table (
order_id STRING,
order_amount DOUBLE,
product_id STRING,
log_ts TIMESTAMP(3),
proctime as PROCTIME()
) WITH (...);
-- streaming sql, kafka join a hive dimension table. Flink will reload all data from dimension_table after cache ttl is expired.
SELECT * FROM orders_table AS o
JOIN dimension_table FOR SYSTEM_TIME AS OF o.proctime AS dim
ON o.product_id = dim.product_id;
참고:
- 각 조인 하위 작업은 Hive 테이블의 자체 캐시를 유지해야 해요. Hive 테이블이 TM 작업 슬롯의 메모리에 들어갈 수 있는지 확인해요.
streaming-source.monitor-interval(최신 파티션을 temporal table로) 또는lookup.join.cache.ttl(모든 파티션을 temporal table로)에 비교적 큰 값을 설정하는 것을 권장해요. 그렇지 않으면 테이블을 너무 자주 갱신·재로드해야 하므로 작업이 성능 문제에 취약해져요.- 현재는 캐시를 새로고침해야 할 때마다 전체 Hive 테이블을 로드해요. 새 데이터와 이전 데이터를 구분할 방법이 없어요.
쓰기 (Writing)
Flink는 BATCH와 STREAMING 모드 모두에서 Hive 데이터 쓰기를 지원해요. BATCH 애플리케이션으로 실행하면 Flink는 작업이 끝날 때만 레코드를 보이게 만들어 Hive 테이블에 써요. BATCH 쓰기는 기존 테이블에 대한 추가와 덮어쓰기 모두를 지원해요.
# ------ INSERT INTO will append to the table or partition, keeping the existing data intact ------
Flink SQL> INSERT INTO mytable SELECT 'Tom', 25;
# ------ INSERT OVERWRITE will overwrite any existing data in the table or partition ------
Flink SQL> INSERT OVERWRITE mytable SELECT 'Tom', 25;
특정 파티션에도 데이터를 삽입할 수 있어요.
# ------ Insert with static partition ------
Flink SQL> INSERT OVERWRITE myparttable PARTITION (my_type='type_1', my_date='2019-08-08') SELECT 'Tom', 25;
# ------ Insert with dynamic partition ------
Flink SQL> INSERT OVERWRITE myparttable SELECT 'Tom', 25, 'type_1', '2019-08-08';
# ------ Insert with static(my_type) and dynamic(my_date) partition ------
Flink SQL> INSERT OVERWRITE myparttable PARTITION (my_type='type_1') SELECT 'Tom', 25, '2019-08-08';
STREAMING 쓰기는 Hive에 새 데이터를 지속적으로 추가하며 레코드를 점진적으로 커밋해 보이게 만들어요. 사용자는 여러 프로퍼티로 커밋을 언제/어떻게 트리거할지 제어해요. 스트리밍 쓰기에는 insert overwrite가 지원되지 않아요.
아래 예제는 스트리밍 싱크를 사용해 Kafka의 데이터를 파티션-커밋으로 Hive 테이블에 쓰는 스트리밍 질의를 작성하고, 배치 질의로 그 데이터를 다시 읽는 방법을 보여줘요. 사용 가능한 구성 전체 목록은 streaming sink를 참고해요.
SET table.sql-dialect=hive;
CREATE TABLE hive_table (
user_id STRING,
order_amount DOUBLE
) PARTITIONED BY (dt STRING, hr STRING) STORED AS parquet TBLPROPERTIES (
'partition.time-extractor.timestamp-pattern'='$dt $hr:00:00',
'sink.partition-commit.trigger'='partition-time',
'sink.partition-commit.delay'='1 h',
'sink.partition-commit.policy.kind'='metastore,success-file'
);
SET table.sql-dialect=default;
CREATE TABLE kafka_table (
user_id STRING,
order_amount DOUBLE,
log_ts TIMESTAMP(3),
WATERMARK FOR log_ts AS log_ts - INTERVAL '5' SECOND -- Define watermark on TIMESTAMP column
) WITH (...);
-- streaming sql, insert into hive table
INSERT INTO TABLE hive_table
SELECT user_id, order_amount, DATE_FORMAT(log_ts, 'yyyy-MM-dd'), DATE_FORMAT(log_ts, 'HH')
FROM kafka_table;
-- batch sql, select with partition pruning
SELECT * FROM hive_table WHERE dt='2020-05-20' and hr='12';
워터마크가 TIMESTAMP_LTZ 컬럼에 정의되고 partition-time으로 커밋한다면, sink.partition-commit.watermark-time-zone을 세션 시간대로 설정해야 해요. 그렇지 않으면 파티션 커밋이 몇 시간 늦게 일어날 수 있어요.
SET table.sql-dialect=hive;
CREATE TABLE hive_table (
user_id STRING,
order_amount DOUBLE
) PARTITIONED BY (dt STRING, hr STRING) STORED AS parquet TBLPROPERTIES (
'partition.time-extractor.timestamp-pattern'='$dt $hr:00:00',
'sink.partition-commit.trigger'='partition-time',
'sink.partition-commit.delay'='1 h',
'sink.partition-commit.watermark-time-zone'='Asia/Shanghai', -- Assume user configured time zone is 'Asia/Shanghai'
'sink.partition-commit.policy.kind'='metastore,success-file'
);
SET table.sql-dialect=default;
CREATE TABLE kafka_table (
user_id STRING,
order_amount DOUBLE,
ts BIGINT, -- time in epoch milliseconds
ts_ltz AS TO_TIMESTAMP_LTZ(ts, 3),
WATERMARK FOR ts_ltz AS ts_ltz - INTERVAL '5' SECOND -- Define watermark on TIMESTAMP_LTZ column
) WITH (...);
-- streaming sql, insert into hive table
INSERT INTO TABLE hive_table
SELECT user_id, order_amount, DATE_FORMAT(ts_ltz, 'yyyy-MM-dd'), DATE_FORMAT(ts_ltz, 'HH')
FROM kafka_table;
-- batch sql, select with partition pruning
SELECT * FROM hive_table WHERE dt='2020-05-20' and hr='12';
기본적으로 스트리밍 쓰기의 경우 Flink는 이름 바꾸기 커미터(renaming committers)만 지원하므로 S3 파일 시스템은 exactly-once 스트리밍 쓰기를 지원할 수 없어요. S3에 대한 exactly-once 쓰기는 다음 파라미터를 false로 구성해 달성할 수 있어요. 이는 싱크가 Flink 네이티브 writer를 사용하도록 지시하지만 parquet와 orc 파일 타입에서만 동작해요. 이 구성은 TableConfig에서 설정되며 작업의 모든 싱크에 영향을 줘요.
| 키 (Key) | 기본값 (Default) | 타입 (Type) | 설명 (Description) |
|---|---|---|---|
| table.exec.hive.fallback-mapred-writer | true | Boolean | false이면 flink 네이티브 writer로 parquet·orc 파일을 쓰고, true이면 hadoop mapred record writer로 parquet·orc 파일을 써요. |
동적 파티션 쓰기 (Dynamic Partition Writing)
파티션 컬럼 값을 사용자가 지정해야 하는 정적 파티션 쓰기와 달리, 동적 파티션 쓰기는 파티션 컬럼 값을 지정하지 않아도 돼요. 예를 들어 다음과 같은 파티션 테이블의 경우:
CREATE TABLE fact_tz(x int) PARTITIONED BY (day STRING, hour STRING);
다음 SQL 문으로 파티션 테이블 fact_tz에 데이터를 쓸 수 있어요:
INSERT INTO TABLE fact_tz PARTITION (day, hour) select 1, '2022-8-8', '14';
SQL 문에서 파티션 컬럼 값을 지정하지 않으므로 동적 파티션 쓰기의 전형적인 경우예요. 기본적으로 동적 파티션 쓰기의 경우 Flink가 싱크 테이블에 쓰기 전에 동적 파티션 컬럼으로 데이터를 추가로 정렬해요. 이는 싱크가 한 파티션의 모든 요소를 받고 나서 다른 파티션의 모든 요소를 받는다는 뜻이에요. 서로 다른 파티션의 요소는 섞이지 않아요. 이는 Hive 싱크가 한 번에 한 파티션을 씀으로써 파티션 writer 수를 줄이고 쓰기 성능을 개선하는 데 도움이 돼요. 그렇지 않으면 너무 많은 파티션 writer가 OutOfMemory 예외를 유발할 수 있어요.
추가 정렬을 피하려면 작업 구성 table.exec.hive.sink.sort-by-dynamic-partition.enable(기본 true)을 false로 설정할 수 있어요. 하지만 앞서 말했듯이 그러한 구성으로는 같은 싱크 노드에 너무 많은 파티션이 들어오면 OutOfMemory 예외가 발생할 수 있어요.
너무 많은 파티션 writer 문제를 완화하려면, 데이터가 치우치지 않았다면 SQL 문에 DISTRIBUTED BY <partition_field>를 추가해 같은 파티션의 데이터를 같은 노드로 셔플할 수 있어요. 또한 SQL 문에 SORTED BY <partition_field>를 수동으로 추가해 table.exec.hive.sink.sort-by-dynamic-partition.enable=true와 같은 목적을 달성할 수 있어요.
참고:
- 구성
table.exec.hive.sink.sort-by-dynamic-partition.enable은 Flink BATCH 모드에서만 동작해요. - 현재 DISTRIBUTED BY와 SORTED BY는 Flink BATCH 모드에서 Hive dialect를 사용할 때만 지원돼요.
통계 자동 수집 (Auto Gather Statistic)
기본적으로 Flink는 Hive 테이블을 쓸 때 통계를 자동으로 수집해 Hive metastore에 커밋해요. 하지만 통계 수집이 시간이 걸릴 수 있어 어떤 경우에는 비활성화하고 싶을 수 있어요. 작업 구성 table.exec.hive.sink.statistic-auto-gather.enable(기본 true)을 false로 설정해 비활성화할 수 있어요.
Hive 테이블이 Parquet 또는 ORC 포맷으로 저장되어 있으면 numFiles/totalSize/numRows/rawDataSize를 수집할 수 있어요. 그렇지 않으면 numFiles/totalSize만 수집할 수 있어요. Parquet와 ORC 포맷의 numRows/rawDataSize 통계를 수집하기 위해 Flink는 파일의 footer만 읽어 빠른 수집을 해요. 하지만 파일이 너무 많으면 여전히 시간이 걸릴 수 있는데, 작업 구성 table.exec.hive.sink.statistic-auto-gather.thread-num(기본 3)으로 더 많은 스레드를 사용해 수집을 가속화할 수 있어요.
참고:
- BATCH 모드만 통계 자동 수집을 지원하며, STREAMING 모드는 아직 지원하지 않아요.
파일 압축 (File Compaction)
Hive 싱크는 파일 압축도 지원하며, 이를 통해 애플리케이션이 Hive에 쓸 때 생성되는 파일 수를 줄일 수 있어요.
스트림 모드 (Stream Mode)
스트림 모드에서는 FileSystem 싱크와 동작이 같아요. 자세한 내용은 File Compaction을 참고해요.
배치 모드 (Batch Mode)
배치 모드이고 자동 압축이 활성화되어 있으면 파일 쓰기를 마친 후 Flink가 각 파티션의 쓰여진 파일의 평균 크기를 계산해요. 평균 크기가 구성된 임계값보다 작으면 Flink는 이 파일들을 대상 크기의 파일로 압축하려고 해요. 다음은 파일 압축의 테이블 옵션이에요.
| 옵션 (Option) | 필수 (Required) | 전달 (Forwarded) | 기본값 (Default) | 타입 (Type) | 설명 (Description) |
|---|---|---|---|---|---|
| auto-compaction | optional | no | false | Boolean | Hive 싱크에서 자동 압축을 활성화할지 여부. 데이터는 먼저 임시 파일에 쓰여져요. 임시 파일은 압축 전까지 보이지 않아요. |
| compaction.small-files.avg-size | optional | yes | 16MB | MemorySize | 파일 압축 임계값. 파일의 평균 크기가 이 값보다 작으면 Flink가 이 파일들을 압축해요. 기본값은 16MB예요. |
| compaction.file-size | optional | yes | (none) | MemorySize | 압축 대상 파일 크기. 기본값은 롤링 파일 크기예요. |
| compaction.parallelism | optional | no | (none) | Integer | 파일 압축 병렬도. 설정하지 않으면 싱크 병렬도를 사용해요. 적응형 배치 스케줄러를 사용할 때 스케줄러가 추론한 compact 연산자의 병렬도가 작아 압축을 끝내는 데 많은 시간이 걸릴 수 있어요. 그런 경우 이 옵션을 더 큰 값으로 수동 설정해주세요. |
포맷 (Formats)
Flink의 Hive 통합은 다음 파일 포맷에 대해 테스트되었어요:
- Text
- CSV
- SequenceFile
- ORC
- Parquet