테이블 & SQL 커넥터
테이블 & SQL 커넥터 (Table & SQL Connectors)
Flink의 Table API & SQL 프로그램은 배치 및 스트리밍 테이블을 읽고 쓰기 위해 다른 외부 시스템에 연결될 수 있습니다. 테이블 소스(table source)는 데이터베이스, 키-값 저장소, 메시지 큐, 파일 시스템 같은 외부 시스템에 저장된 데이터에 대한 접근을 제공합니다. 테이블 싱크(table sink)는 테이블을 외부 저장 시스템으로 내보냅니다. 소스와 싱크의 유형에 따라 CSV, Avro, Parquet, ORC 같은 다양한 포맷을 지원합니다.
이 페이지는 네이티브로 지원되는 커넥터를 사용해 Flink에서 테이블 소스와 테이블 싱크를 등록하는 방법을 설명합니다. 소스나 싱크가 등록되면 Table API & SQL 문으로 접근할 수 있습니다.
자체 커스텀 테이블 소스나 싱크를 구현하려면 사용자 정의 소스 & 싱크 페이지를 살펴보세요.
출처: 문서
본문
지원되는 커넥터 (Supported Connectors)
Flink는 다양한 커넥터를 네이티브로 지원합니다. 다음 표는 사용 가능한 모든 커넥터를 나열합니다.
| 이름 | 버전 | 소스 | 싱크 |
|---|---|---|---|
| Filesystem | 유한 및 무한 스캔 | 스트리밍 싱크, 배치 싱크 | |
| Elasticsearch | 6.x & 7.x | 지원 안 함 | 스트리밍 싱크, 배치 싱크 |
| Opensearch | 1.x & 2.x | 지원 안 함 | 스트리밍 싱크, 배치 싱크 |
| Apache Kafka | 0.10+ | 무한 스캔 | 스트리밍 싱크, 배치 싱크 |
| Amazon DynamoDB | 지원 안 함 | 스트리밍 싱크, 배치 싱크 | |
| Amazon Kinesis Data Streams | 무한 스캔 | 스트리밍 싱크 | |
| Amazon Kinesis Data Firehose | 지원 안 함 | 스트리밍 싱크 | |
| JDBC | 유한 스캔, Lookup | 스트리밍 싱크, 배치 싱크 | |
| Apache HBase | 1.4.x & 2.2.x | 유한 스캔, Lookup | 스트리밍 싱크, 배치 싱크 |
| Apache Hive | 지원 버전 | 무한 스캔, 유한 스캔, Lookup | 스트리밍 싱크, 배치 싱크 |
| MongoDB | 3.6.x & 4.x & 5.x & 6.x & 7.0.x | 유한 스캔, Lookup | 스트리밍 싱크, 배치 싱크 |
커넥터를 의존성으로 추가하는 방법은 구성(Configuration) 섹션을 참고하세요.
커넥터 사용법 (How to use connectors)
Flink는 SQL CREATE TABLE 문으로 테이블을 등록하는 것을 지원합니다. 테이블 이름, 테이블 스키마, 외부 시스템 연결을 위한 테이블 옵션을 정의할 수 있습니다. 테이블 생성에 대한 자세한 내용은 SQL 섹션을 참고하세요.
다음 코드는 Kafka에 연결해 JSON 레코드를 읽고 쓰는 전체 예시입니다.
SQL
CREATE TABLE MyUserTable (
-- declare the schema of the table
`user` BIGINT,
`message` STRING,
`rowtime` TIMESTAMP(3) METADATA FROM 'timestamp', -- use a metadata column to access Kafka's record timestamp
`proctime` AS PROCTIME(), -- use a computed column to define a proctime attribute
WATERMARK FOR `rowtime` AS `rowtime` - INTERVAL '5' SECOND -- use a WATERMARK statement to define a rowtime attribute
) WITH (
-- declare the external system to connect to
'connector' = 'kafka',
'topic' = 'topic_name',
'scan.startup.mode' = 'earliest-offset',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json' -- declare a format for this system
)
원하는 연결 속성은 문자열 기반의 키-값 쌍으로 변환됩니다. 팩토리(factory)는 팩토리 식별자(이 예시에서는 kafka와 json)를 기반으로 키-값 쌍에서 구성된 테이블 소스, 테이블 싱크, 해당 포맷을 만듭니다. 각 컴포넌트에 대해 정확히 하나의 일치하는 팩토리를 검색할 때 Java의 Service Provider Interface(SPI)로 찾을 수 있는 모든 팩토리가 고려됩니다.
주어진 속성에 대해 팩토리를 찾을 수 없거나 여러 팩토리가 일치하면 고려된 팩토리와 지원되는 속성에 대한 추가 정보와 함께 예외가 던져집니다.
테이블 커넥터/포맷 리소스 변환 (Transform table connector/format resources)
Flink는 Java의 Service Provider Interface(SPI)를 사용해 식별자로 테이블 커넥터/포맷 팩토리를 로드합니다. 모든 테이블 커넥터/포맷에 대한 org.apache.flink.table.factories.Factory라는 SPI 리소스 파일이 같은 디렉터리 META-INF/services 아래에 있으므로, 둘 이상의 테이블 커넥터/포맷을 사용하는 프로젝트의 uber-jar를 빌드할 때 이러한 리소스 파일이 서로 덮어쓰게 되어 Flink가 테이블 커넥터/포맷 팩토리를 로드하지 못할 수 있습니다.
이 상황에서 권장하는 방법은 maven shade 플러그인의 ServicesResourceTransformer로 META-INF/services 디렉터리 아래의 리소스 파일을 변환하는 것입니다. flink-sql-connector-hive-3.1.3 커넥터와 flink-parquet 포맷을 포함한 프로젝트 예시의 pom.xml 내용은 다음과 같습니다.
<modelVersion>4.0.0</modelVersion>
<groupId>org.example</groupId>
<artifactId>myProject</artifactId>
<version>1.0-SNAPSHOT</version>
<dependencies>
<!-- other project dependencies ...-->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-sql-connector-hive-3.1.3_2.12</artifactId>
<version>2.3.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-parquet_2.12</artifactId>
<version>2.3.0</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<executions>
<execution>
<id>shade</id>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
<configuration>
<transformers combine.children="append">
<!-- The service transformer is needed to merge META-INF/services files -->
<transformer implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"/>
<!-- ... -->
</transformers>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
ServicesResourceTransformer를 구성하면 위 프로젝트의 uber-jar를 빌드할 때 META-INF/services 디렉터리 아래의 테이블 커넥터/포맷 리소스 파일이 서로 덮어쓰지 않고 병합됩니다.
스키마 매핑 (Schema Mapping)
SQL CREATE TABLE 문의 본문 절은 물리적 컬럼의 이름과 타입, 제약 조건, 워터마크를 정의합니다. Flink는 데이터를 보유하지 않으므로 스키마 정의는 외부 시스템의 물리적 컬럼을 Flink의 표현으로 매핑하는 방법만 선언합니다. 매핑은 이름으로 되지 않을 수도 있으며 포맷과 커넥터의 구현에 따라 달라집니다. 예를 들어 MySQL 데이터베이스 테이블은 필드 이름(대소문자 구분 안 함)으로 매핑되고, CSV 파일 시스템은 필드 순서(필드 이름은 임의일 수 있음)로 매핑됩니다. 이는 각 커넥터에서 설명됩니다.
다음 예시는 시간 속성이 없고 입출력과 테이블 컬럼의 일대일 필드 매핑이 있는 단순한 스키마입니다.
SQL
CREATE TABLE MyTable (
MyField1 INT,
MyField2 STRING,
MyField3 BOOLEAN
) WITH (
...
)
메타데이터 (Metadata)
일부 커넥터와 포맷은 물리적 페이로드 컬럼 옆의 메타데이터 컬럼에서 접근할 수 있는 추가 메타데이터 필드를 노출합니다. 메타데이터 컬럼에 대한 자세한 내용은 CREATE TABLE 섹션을 참고하세요.
기본 키 (Primary Key)
기본 키 제약 조건은 테이블의 한 컬럼 또는 컬럼 집합이 고유하고 null을 포함하지 않음을 나타냅니다. 기본 키는 테이블의 한 행을 고유하게 식별합니다.
소스 테이블의 기본 키는 최적화를 위한 메타데이터 정보입니다. 싱크 테이블의 기본 키는 보통 싱크 구현이 upsert하는 데 사용합니다.
SQL 표준은 제약 조건이 ENFORCED 또는 NOT ENFORCED일 수 있음을 지정합니다. 이는 들어오고 나가는 데이터에 대해 제약 조건 검사가 수행되는지 여부를 제어합니다. Flink는 데이터를 소유하지 않으므로 지원하려는 유일한 모드는 NOT ENFORCED입니다. 쿼리가 키 무결성을 강제하도록 하는 것은 사용자의 몫입니다.
CREATE TABLE MyTable (
MyField1 INT,
MyField2 STRING,
MyField3 BOOLEAN,
PRIMARY KEY (MyField1, MyField2) NOT ENFORCED -- defines a primary key on columns
) WITH (
...
)
시간 속성 (Time Attributes)
시간 속성은 무한 스트리밍 테이블을 다룰 때 필수적입니다. 따라서 proctime과 rowtime 속성 모두 스키마의 일부로 정의할 수 있습니다. Flink에서의 시간 처리, 특히 이벤트 시간에 대한 자세한 내용은 일반 이벤트-시간 섹션을 권장합니다.
Proctime 속성 (Proctime Attributes)
스키마에서 proctime 속성을 선언하려면 Computed Column 문법을 사용해 PROCTIME() 내장 함수에서 생성되는 계산 컬럼을 선언할 수 있습니다. 계산 컬럼은 물리적 데이터에 저장되지 않는 가상(virtual) 컬럼입니다.
CREATE TABLE MyTable (
MyField1 INT,
MyField2 STRING,
MyField3 BOOLEAN,
MyField4 AS PROCTIME() -- declares a proctime attribute
) WITH (
...
)
Rowtime 속성 (Rowtime Attributes)
테이블의 이벤트-시간 동작을 제어하기 위해 Flink는 사전 정의된 타임스탬프 추출기와 워터마크 전략을 제공합니다. DDL에서 시간 속성을 정의하는 방법에 대한 자세한 내용은 CREATE TABLE 문을 참고하세요.
다음 타임스탬프 추출기가 지원됩니다:
-- use the existing TIMESTAMP(3) field in schema as the rowtime attribute
CREATE TABLE MyTable (
ts_field TIMESTAMP(3),
WATERMARK FOR ts_field AS ...
) WITH (
...
)
-- use system functions or UDFs or expressions to extract the expected TIMESTAMP(3) rowtime field
CREATE TABLE MyTable (
log_ts STRING,
ts_field AS TO_TIMESTAMP(log_ts),
WATERMARK FOR ts_field AS ...
) WITH (
...
)
다음 워터마크 전략이 지원됩니다:
-- Sets a watermark strategy for strictly ascending rowtime attributes. Emits a watermark of the
-- maximum observed timestamp so far. Rows that have a timestamp bigger to the max timestamp
-- are not late.
CREATE TABLE MyTable (
ts_field TIMESTAMP(3),
WATERMARK FOR ts_field AS ts_field
) WITH (
...
)
-- Sets a watermark strategy for ascending rowtime attributes. Emits a watermark of the maximum
-- observed timestamp so far minus 1. Rows that have a timestamp bigger or equal to the max timestamp
-- are not late.
CREATE TABLE MyTable (
ts_field TIMESTAMP(3),
WATERMARK FOR ts_field AS ts_field - INTERVAL '0.001' SECOND
) WITH (
...
)
-- Sets a watermark strategy for rowtime attributes which are out-of-order by a bounded time interval.
-- Emits watermarks which are the maximum observed timestamp minus the specified delay, e.g. 2 seconds.
CREATE TABLE MyTable (
ts_field TIMESTAMP(3),
WATERMARK FOR ts_field AS ts_field - INTERVAL '2' SECOND
) WITH (
...
)
항상 타임스탬프와 워터마크를 함께 선언해야 합니다. 워터마크는 시간 기반 연산을 트리거하는 데 필요합니다.
SQL 타입 (SQL Types)
SQL에서 타입을 선언하는 방법은 데이터 타입(Data Types) 페이지를 참고하세요.