플링크 시작하기

플링크 시작하기 (Getting Started)

아파치 아이스버그는 아파치 플링크의 DataStream API와 Table API를 모두 지원해요. 이 문서에서는 Flink에서 아이스버그를 시작하는 방법을 알려드릴게요. Flink SQL 클라이언트 준비, Python API 사용, 카탈로그 추가, 테이블 생성, 쓰기, 읽기, 그리고 Flink·Iceberg 타입 변환까지 실제 명령과 코드 예시로 차근차근 살펴볼게요.

출처: 문서

본문

팁 (Tip)

Flink에서 아이스버그를 사용하는 것에 대한 개요는 Flink Quickstart 문서를 참고해주세요.

아파치 아이스버그는 아파치 플링크의 DataStream API와 Table API를 모두 지원해요. 아파치 플링크의 통합에 대한 내용은 다중 엔진 지원(Multi-Engine Support) 페이지를 참고해주세요.

기능 지원 Flink 참고
SQL create catalog ✔️
SQL create database ✔️
SQL create table ✔️
SQL create table like ✔️
SQL alter table ✔️ 테이블 속성 변경만 지원, 컬럼과 파티션 변경은 지원하지 않음
SQL drop_table ✔️
SQL select ✔️ 스트리밍과 배치 모드 모두 지원
SQL insert into ✔️ 스트리밍과 배치 모드 모두 지원
SQL insert overwrite ✔️
DataStream read ✔️
DataStream append ✔️
DataStream overwrite ✔️
Metadata tables ✔️
Rewrite files action ✔️

Flink에서 아이스버그 테이블을 만들 때는 Flink SQL Client를 사용하는 것이 좋아요. 개념을 이해하기 더 쉽기 때문이에요.

Apache 다운로드 페이지에서 Flink를 다운로드해요. 아이스버그는 Apache iceberg-flink-runtime jar를 컴파일할 때 Scala 2.12을 사용하므로 Scala 2.12이 번들된 Flink 2.3을 사용하는 것을 권장해요.

FLINK_VERSION=2.3.0
SCALA_VERSION=2.12
APACHE_FLINK_URL=https://archive.apache.org/dist/flink/
wget ${APACHE_FLINK_URL}/flink-${FLINK_VERSION}/flink-${FLINK_VERSION}-bin-scala_${SCALA_VERSION}.tgz
tar xzvf flink-${FLINK_VERSION}-bin-scala_${SCALA_VERSION}.tgz

Hadoop 환경 안에서 독립 실행형(standalone) Flink 클러스터를 시작해요.

# HADOOP_HOME is your hadoop root directory after unpack the binary package.
APACHE_HADOOP_URL=https://archive.apache.org/dist/hadoop/
HADOOP_VERSION=2.8.5
wget ${APACHE_HADOOP_URL}/common/hadoop-${HADOOP_VERSION}/hadoop-${HADOOP_VERSION}.tar.gz
tar xzvf hadoop-${HADOOP_VERSION}.tar.gz
HADOOP_HOME=`pwd`/hadoop-${HADOOP_VERSION}

export HADOOP_CLASSPATH=`$HADOOP_HOME/bin/hadoop classpath`

# Start the flink standalone cluster
cd flink-${FLINK_VERSION}/
./bin/start-cluster.sh

Flink SQL 클라이언트를 시작해요. 아이스버그 프로젝트에는 bundled jar를 생성하는 별도의 flink-runtime 모듈이 있는데, 이 jar는 Flink SQL 클라이언트가 직접 로드할 수 있어요. flink-runtime bundled jar를 수동으로 빌드하려면 아이스버그 프로젝트를 빌드하면 /flink-runtime/build/libs 아래에 jar가 생성돼요. 또는 Apache 저장소에서 flink-runtime jar를 다운로드해요.

# HADOOP_HOME is your hadoop root directory after unpack the binary package.
export HADOOP_CLASSPATH=`$HADOOP_HOME/bin/hadoop classpath`   

# Below works for Flink 1.15 or earlier
./bin/sql-client.sh embedded -j <flink-runtime-directory>/iceberg-flink-runtime-2.3-1.11.0.jar shell

# Flink 1.16+ has a regression in loading external jars via -j. See FLINK-30035 for details.
# put iceberg-flink-runtime-2.3-1.11.0.jar in flink/lib dir
./bin/sql-client.sh embedded shell

기본적으로 아이스버그는 Hadoop 카탈로그용 Hadoop jar를 포함해요. Hive 카탈로그를 사용하려면 Flink SQL 클라이언트를 열 때 Hive jar를 로드해요. 다행히 Flink는 SQL 클라이언트용 bundled hive jar를 제공해요. 의존성을 다운로드하고 시작하는 예시:

# HADOOP_HOME is your hadoop root directory after unpack the binary package.
export HADOOP_CLASSPATH=`$HADOOP_HOME/bin/hadoop classpath`

ICEBERG_VERSION=1.11.0
MAVEN_URL=https://repo1.maven.org/maven2
ICEBERG_MAVEN_URL=${MAVEN_URL}/org/apache/iceberg
ICEBERG_PACKAGE=iceberg-flink-runtime
FLINK_VERSION_MAJOR=2.3
wget ${ICEBERG_MAVEN_URL}/${ICEBERG_PACKAGE}-${FLINK_VERSION_MAJOR}/${ICEBERG_VERSION}/${ICEBERG_PACKAGE}-${FLINK_VERSION_MAJOR}-${ICEBERG_VERSION}.jar -P lib/

HIVE_VERSION=2.3.9
SCALA_VERSION=2.12
FLINK_VERSION=2.3.0
FLINK_CONNECTOR_URL=${MAVEN_URL}/org/apache/flink
FLINK_CONNECTOR_PACKAGE=flink-sql-connector-hive
wget ${FLINK_CONNECTOR_URL}/${FLINK_CONNECTOR_PACKAGE}-${HIVE_VERSION}_${SCALA_VERSION}/${FLINK_VERSION}/${FLINK_CONNECTOR_PACKAGE}-${HIVE_VERSION}_${SCALA_VERSION}-${FLINK_VERSION}.jar

./bin/sql-client.sh embedded shell

정보 (Info)

PyFlink 1.6.1은 Apple Silicon이 있는 macOS에서 알려진 문제가 있어요. FLINK-28786 참조.

pip로 Apache Flink 의존성을 설치해요.

pip install apache-flink==2.3.0

iceberg-flink-runtime jar에 file:// 경로를 제공해요. 이 jar는 프로젝트를 빌드하고 /flink-runtime/build/libs를 보거나 Apache 공식 저장소에서 다운로드해서 얻을 수 있어요. 서드파티 jar는 다음과 같이 pyflink에 추가할 수 있어요.

이것은 공식 문서에도 언급돼 있어요. 아래 예시는 env.add_jars(..)를 사용해요.

import os

from pyflink.datastream import StreamExecutionEnvironment

env = StreamExecutionEnvironment.get_execution_environment()
iceberg_flink_runtime_jar = os.path.join(os.getcwd(), "iceberg-flink-runtime-2.3-1.11.0.jar")

env.add_jars("file://{}".format(iceberg_flink_runtime_jar))

다음으로 StreamTableEnvironment를 만들고 Flink SQL 문을 실행해요. 아래 예시는 Python Table API로 커스텀 카탈로그를 만드는 방법을 보여줘요.

from pyflink.table import StreamTableEnvironment
table_env = StreamTableEnvironment.create(env)
table_env.execute_sql("""
CREATE CATALOG my_catalog WITH (
    'type'='iceberg', 
    'catalog-impl'='com.my.custom.CatalogImpl',
    'my-additional-catalog-config'='my-value'
)
""")

쿼리를 실행해요.

(table_env
    .sql_query("SELECT PULocationID, DOLocationID, passenger_count FROM my_catalog.nyc.taxis LIMIT 5")
    .execute()
    .print()) 
+----+----------------------+----------------------+--------------------------------+
| op |         PULocationID |         DOLocationID |                passenger_count |
+----+----------------------+----------------------+--------------------------------+
| +I |                  249 |                   48 |                            1.0 |
| +I |                  132 |                  233 |                            1.0 |
| +I |                  164 |                  107 |                            1.0 |
| +I |                   90 |                  229 |                            1.0 |
| +I |                  137 |                  249 |                            1.0 |
+----+----------------------+----------------------+--------------------------------+
5 rows in set

자세한 내용은 Python Table API를 참고해주세요.

카탈로그 추가 (Adding catalogs)

Flink는 Flink SQL을 사용한 카탈로그 생성을 지원해요.

카탈로그 구성 (Catalog Configuration)

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

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

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

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

Hive 카탈로그 (Hive catalog)

'catalog-type'='hive'로 구성할 수 있고 Hive metastore에서 테이블을 로드하는 hive_catalog라는 아이스버그 카탈로그를 만들어요.

CREATE CATALOG hive_catalog WITH (
  'type'='iceberg',
  'catalog-type'='hive',
  'uri'='thrift://localhost:9083',
  'clients'='5',
  'property-version'='1',
  'warehouse'='hdfs://nn:8020/warehouse/path'
);

REST 카탈로그 (REST catalog)

'catalog-type'='rest'로 구성할 수 있고 REST 카탈로그에서 테이블을 로드하는 rest_catalog라는 아이스버그 카탈로그를 만들어요.

CREATE CATALOG rest_catalog WITH (
  'type'='iceberg',
  'catalog-type'='rest',
  'uri'='https://localhost/'
);

테이블 만들기 (Creating a table)

CREATE TABLE `hive_catalog`.`default`.`sample` (
    id BIGINT COMMENT 'unique id',
    data STRING
);

쓰기 (Writing)

테이블에 새 데이터를 추가하는 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를 사용해요(Flink 스트리밍 작업은 INSERT OVERWRITE를 지원하지 않아요). Overwrite는 아이스버그 테이블에서 원자적 연산이에요.

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

INSERT OVERWRITE `hive_catalog`.`default`.`sample` VALUES (1, 'a');

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

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

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");

브랜치 쓰기 (Branch Writes)

아이스버그 테이블의 브랜치에 쓰는 것도 FlinkSink의 toBranch API로 지원돼요.

브랜치에 대한 자세한 내용은 branches 문서를 참고해주세요.

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

읽기 (Reading)

다음 문장으로 Flink 배치 작업을 제출해요.

-- Execute the flink job in batch mode for current session context
SET execution.runtime-mode = batch;
SELECT * FROM `hive_catalog`.`default`.`sample`;

아이스버그는 과거 스냅샷 ID에서 시작하는 Flink 스트리밍 작업에서 증분 데이터를 처리하는 것을 지원해요.

-- Submit the flink job in streaming mode for current session.
SET execution.runtime-mode = streaming;

-- Enable this switch because streaming read SQL will provide few job options in flink SQL hint options.
SET table.dynamic-table-options.enabled=true;

-- Read all the records from the iceberg current snapshot, and then read incremental data starting from that snapshot.
SELECT * FROM `hive_catalog`.`default`.`sample` /*+ OPTIONS('streaming'='true', 'monitor-interval'='1s')*/ ;

-- Read all incremental data starting from the snapshot-id '3821550127947089987' (records from this snapshot will be excluded).
SELECT * FROM `hive_catalog`.`default`.`sample` /*+ OPTIONS('streaming'='true', 'monitor-interval'='1s', 'start-snapshot-id'='3821550127947089987')*/ ;

SQL은 테이블을 검사하는 데도 권장되는 방법이에요. 테이블의 모든 스냅샷을 보려면 snapshots 메타데이터 테이블을 사용해요.

SELECT * FROM `hive_catalog`.`default`.`sample$snapshots`;

아이스버그는 자바 API에서 스트리밍 또는 배치 읽기를 지원해요.

DataStream<RowData> batch = FlinkSource.forRowData()
     .env(env)
     .tableLoader(tableLoader)
     .streaming(false)
     .build();

타입 변환 (Type conversion)

아이스버그의 Flink 통합은 Flink 타입과 아이스버그 타입을 자동으로 변환해요. Flink가 지원하지 않는 타입(예: UUID)을 가진 테이블에 쓸 때, 아이스버그는 Flink 타입의 값을 받아 변환해요.

Flink 타입은 다음 표에 따라 아이스버그 타입으로 변환돼요.

Flink Iceberg 참고
boolean boolean
tinyint integer
smallint integer
integer integer
bigint long
float float
double double
char string
varchar string
string string
binary binary
varbinary fixed
decimal decimal
date date
time time
timestamp timestamp without timezone
timestamp_ltz timestamp with timezone
array list
map map
multiset map
row struct
raw 지원되지 않음
interval 지원되지 않음
structured 지원되지 않음
timestamp with zone 지원되지 않음
distinct 지원되지 않음
null 지원되지 않음
symbol 지원되지 않음
logical 지원되지 않음

아이스버그 타입은 다음 표에 따라 Flink 타입으로 변환돼요.

Iceberg Flink 참고
boolean boolean
struct row
list array
map map
integer integer
long bigint
float float
double double
date date
time time
timestamp without timezone timestamp(6)
timestamp with timezone timestamp_ltz(6)
string varchar(2147483647)
uuid binary(16)
fixed(N) binary(N)
binary varbinary(2147483647)
decimal(P, S) decimal(P, S)
nanosecond timestamp timestamp(9)
nanosecond timestamp with timezone timestamp_ltz(9)
unknown null
variant 지원되지 않음
geometry 지원되지 않음
geography 지원되지 않음

향후 개선 사항 (Future improvements)

현재 Flink 아이스버그 통합에서 아직 지원되지 않는 기능이 몇 가지 있어요.

  • 숨겨진 파티셔닝으로 아이스버그 테이블 만들기. flink 메일 리스트에서 논의 중.
  • 계산된 컬럼(computed column)으로 아이스버그 테이블 만들기.

더 알아보기 (Learn more)