플링크 커넥터

아파치 플링크는 명시적인 Flink 카탈로그를 만들지 않고도 Flink SQL에서 아이스버그 테이블을 바로 만들 수 있게 지원해요. 이 문서에서는 Flink SQL의 CREATE TABLE 문에 'connector'='iceberg' 테이블 옵션을 지정해서 아이스버그 테이블에 연결하는 방법을 알려드릴게요. Hive 카탈로그, hadoop 카탈로그, REST 카탈로그, 커스텀 카탈로그를 각각 어떻게 지정하는지 살펴볼게요.

출처: 문서

본문

아파치 플링크는 명시적인 Flink 카탈로그를 만들지 않고 Flink SQL에서 아이스버그 테이블을 직접 만드는 것을 지원해요. 즉 Flink 공식 문서에서의 사용법과 비슷하게, Flink SQL에서 'connector'='iceberg' 테이블 옵션을 지정해서 아이스버그 테이블을 만들 수 있어요.

Flink에서 SQL CREATE TABLE test (..) WITH ('connector'='iceberg', ...)는 현재 Flink 카탈로그(기본적으로 GenericInMemoryCatalog)에 Flink 테이블을 만들어요. 이 테이블은 기본 아이스버그 테이블에 매핑(mapping)될 뿐, 현재 Flink 카탈로그에서 직접 아이스버그 테이블을 유지 관리하는 것은 아니에요.

SQL 구문 CREATE TABLE test (..) WITH ('connector'='iceberg', ...)으로 Flink SQL에서 테이블을 만들려면, Flink 아이스버그 커넥터는 테이블 속성을 통해 카탈로그 속성을 설정할 수 있게 해줘요. 유효한 속성 값은 플링크 구성(Flink Configuration) 페이지에 자세히 설명돼 있어요.

Hive 카탈로그에서 관리되는 테이블 (Table managed in Hive catalog)

다음 SQL을 실행하기 전에, 퀵스타트 문서에 따라 Flink SQL 클라이언트를 올바르게 구성했는지 확인해주세요.

다음 SQL은 현재 Flink 카탈로그에 Flink 테이블을 만들고, 이 테이블은 아이스버그 카탈로그에서 관리되는 아이스버그 테이블 default_database.flink_table에 매핑돼요.

CREATE TABLE flink_table (
    id   BIGINT,
    data STRING
) WITH (
    'connector'='iceberg',
    'catalog-name'='hive_prod',
    'uri'='thrift://localhost:9083',
    'warehouse'='hdfs://nn:8020/path/to/warehouse'
);

Hive 카탈로그에서 관리되는 서로 다른 아이스버그 테이블(예: Hive의 hive_db.hive_iceberg_table)에 매핑하는 Flink 테이블을 만들고 싶다면, 다음과 같이 Flink 테이블을 만들 수 있어요.

CREATE TABLE flink_table (
    id   BIGINT,
    data STRING
) WITH (
    'connector'='iceberg',
    'catalog-name'='hive_prod',
    'catalog-database'='hive_db',
    'catalog-table'='hive_iceberg_table',
    'uri'='thrift://localhost:9083',
    'warehouse'='hdfs://nn:8020/path/to/warehouse'
);

정보 (Info)

기본 카탈로그 데이터베이스(위 예시의 hive_db)는 Flink 테이블에 레코드를 쓸 때 존재하지 않으면 자동으로 생성돼요.

hadoop 카탈로그에서 관리되는 테이블 (Table managed in hadoop catalog)

다음 SQL은 현재 Flink 카탈로그에 Flink 테이블을 만들고, 이 테이블은 hadoop 카탈로그에서 관리되는 아이스버그 테이블 default_database.flink_table에 매핑돼요.

CREATE TABLE flink_table (
    id   BIGINT,
    data STRING
) WITH (
    'connector'='iceberg',
    'catalog-name'='hadoop_prod',
    'catalog-type'='hadoop',
    'warehouse'='hdfs://nn:8020/path/to/warehouse'
);

REST 카탈로그에서 관리되는 테이블 (Table managed in REST catalog)

다음 SQL은 현재 Flink 카탈로그에 Flink 테이블을 만들고, 이 테이블은 REST 카탈로그에서 관리되는 아이스버그 테이블 default_database.flink_table에 매핑돼요.

CREATE TABLE flink_table (
    id   BIGINT,
    data STRING
) WITH (
    'connector'='iceberg',
    'catalog-name'='rest_prod',
    'catalog-type'='rest',
    'uri'='https://localhost/'
    'credential'='xxxx' -- Optional
    'token'='xxxx' -- Optional
    'scope'='xxxx' -- Optional
     ...
);

커스텀 카탈로그에서 관리되는 테이블 (Table managed in custom catalog)

다음 SQL은 현재 Flink 카탈로그에 Flink 테이블을 만들고, 이 테이블은 com.my.custom.CatalogImpl 타입의 커스텀 카탈로그에서 관리되는 아이스버그 테이블 default_database.flink_table에 매핑돼요.

CREATE TABLE flink_table (
    id   BIGINT,
    data STRING
) WITH (
    'connector'='iceberg',
    'catalog-name'='custom_prod',
    'catalog-impl'='com.my.custom.CatalogImpl',
     -- More table properties for the customized catalog
    'my-additional-catalog-config'='my-value',
     ...
);

모든 커스텀 카탈로그는 Integrations 탭 아래의 섹션들을 확인해주세요.

완전한 예시 (A complete example)

Hive 카탈로그를 예로 들게요.

CREATE TABLE flink_table (
    id   BIGINT,
    data STRING
) WITH (
    'connector'='iceberg',
    'catalog-name'='hive_prod',
    'uri'='thrift://localhost:9083',
    'warehouse'='file:///path/to/warehouse'
);

INSERT INTO flink_table VALUES (1, 'AAA'), (2, 'BBB'), (3, 'CCC');

SET execution.result-mode=tableau;
SELECT * FROM flink_table;

+----+------+
| id | data |
+----+------+
|  1 |  AAA |
|  2 |  BBB |
|  3 |  CCC |
+----+------+
3 rows in set

자세한 내용은 Iceberg Flink 문서를 참고해주세요.

더 알아보기 (Learn more)