사용자 정의 소스 & 싱크

사용자 정의 소스 & 싱크 (User-defined Sources & Sinks)

*동적 테이블(Dynamic tables)*은 유계(bounded)와 무계(unbounded) 데이터를 통합된 방식으로 처리하기 위한 Flink Table & SQL API의 핵심 개념이에요. 동적 테이블은 논리적 개념일 뿐이므로 Flink가 데이터 자체를 소유하지 않아요. 대신 동적 테이블의 내용은 외부 시스템(데이터베이스, 키-값 저장소, 메시지 큐 등)이나 파일에 저장돼요. 동적 소스동적 싱크는 외부 시스템에서/으로 데이터를 읽고 쓰는 데 사용되며, 문서에서 소스와 싱크는 종종 커넥터라는 용어로 요약돼요.

출처: 문서

본문

Flink는 Kafka, Hive, 다양한 파일 시스템용 사전 정의 커넥터를 제공해요. 내장 테이블 소스·싱크에 대한 자세한 내용은 커넥터 섹션을 참고해요. 이 페이지는 커스텀 사용자 정의 커넥터를 개발하는 방법에 초점을 맞춰요.

Flink v1.16부터 TableEnvironment는 테이블 프로그램, SQL Client, SQL Gateway에서 일관된 클래스 로딩 동작을 위해 사용자 클래스 로더를 도입했어요. 사용자 클래스로더는 ADD JAR이나 CREATE FUNCTION .. USING JAR .. 문으로 추가된 jar 같은 모든 사용자 jar를 관리해요. 사용자 정의 커넥터는 Thread.currentThread().getContextClassLoader()를 사용자 클래스 로더로 대체해 클래스를 로드해야 해요. 그렇지 않으면 ClassNotFoundException이 발생할 수 있어요. 사용자 클래스 로더는 DynamicTableFactory.Context로 접근할 수 있어요.

개요 (Overview)

많은 경우 구현자는 처음부터 새 커넥터를 만들 필요 없이 기존 커넥터를 약간 수정하거나 기존 스택에 훅을 걸기를 원해요. 다른 경우 구현자는 특화된 커넥터를 만들고 싶어해요. 이 섹션은 두 종류의 사용 사례를 모두 돕습니다. API에서의 순수 선언부터 클러스터에서 실행될 런타임 코드까지 테이블 커넥터의 일반 아키텍처를 설명해요. 채워진 화살표는 번역 과정에서 한 단계에서 다음 단계로 객체가 어떻게 변환되는지 보여줘요.

메타데이터 (Metadata)

Table API와 SQL은 모두 선언적 API이며 테이블의 선언을 포함해요. 따라서 CREATE TABLE 문을 실행하면 대상 카탈로그의 메타데이터가 갱신돼요. 대부분의 카탈로그 구현에서 그러한 연산으로 외부 시스템의 물리 데이터는 수정되지 않아요. 커넥터 특정 의존성은 아직 클래스패스에 있을 필요가 없어요. WITH 절에 선언된 옵션은 검증되거나 다른 방식으로 해석되지 않아요.

동적 테이블(DDL로 생성되거나 카탈로그가 제공)의 메타데이터는 CatalogTable 인스턴스로 표현돼요. 테이블 이름은 필요할 때 내부적으로 CatalogTable로 해석돼요.

계획 (Planning)

테이블 프로그램의 계획과 최적화와 관련해서 CatalogTable은 DynamicTableSource(선택 질의에서 읽기용)와 DynamicTableSink(INSERT INTO 문에서 쓰기용)로 해석되어야 해요.

DynamicTableSourceFactory와 DynamicTableSinkFactory는 CatalogTable의 메타데이터를 DynamicTableSource와 DynamicTableSink 인스턴스로 번역하는 커넥터 특정 로직을 제공해요. 대부분의 경우 팩토리의 목적은 옵션을 검증하고(예제의 'port' = '5022'), 인코딩/디코딩 포맷을 구성하고(필요한 경우), 테이블 커넥터의 파라미터화된 인스턴스를 만드는 것이에요.

기본적으로 DynamicTableSourceFactory와 DynamicTableSinkFactory의 인스턴스는 Java의 Service Provider Interfaces (SPI)를 사용해 발견돼요. 커넥터 옵션(예제의 'connector' = 'custom')은 유효한 팩토리 식별자와 일치해야 해요.

클래스 이름에 나타나지 않을 수 있지만 DynamicTableSource와 DynamicTableSink는 실제 데이터를 읽고 쓰기 위한 구체 런타임 구현을 결국 만들게 하는 상태 있는 팩토리로도 볼 수 있어요. 플래너는 최적의 논리 계획을 찾을 때까지 소스·싱크 인스턴스를 사용해 커넥터 특정 양방향 통신을 수행해요. 선택적으로 선언된 능력 인터페이스(예: SupportsProjectionPushDown 또는 SupportsOverwrite)에 따라 플래너는 인스턴스에 변경을 적용해 생성된 런타임 구현을 변경할 수 있어요.

런타임 (Runtime)

논리 계획이 완료되면 플래너는 테이블 커넥터에서 런타임 구현을 얻어요. 런타임 로직은 InputFormat이나 SourceFunction 같은 Flink의 핵심 커넥터 인터페이스로 구현돼요. 이 인터페이스들은 ScanRuntimeProvider, LookupRuntimeProvider, SinkRuntimeProvider의 하위 클래스로 또 다른 추상화 수준으로 그룹화돼요.

예를 들어 OutputFormatProvider(org.apache.flink.api.common.io.OutputFormat 제공)와 SinkFunctionProvider(org.apache.flink.streaming.api.functions.sink.SinkFunction 제공)는 모두 플래너가 처리할 수 있는 SinkRuntimeProvider의 구체 인스턴스예요.

프로젝트 구성 (Project Configuration)

커스텀 커넥터나 커스텀 포맷을 구현하려면 일반적으로 다음 의존성이면 충분해요.

Maven — 프로젝트 디렉터리의 pom.xml 파일을 열고 dependencies 블록에 다음을 추가해요:

<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-common</artifactId>
<version>2.3.0</version>
<scope>provided</scope>
</dependency>

자세한 내용은 Project configuration을 참고해요.

Gradle — 프로젝트 디렉터리의 build.gradle 파일을 열고 dependencies 블록에 다음을 추가해요:

runtime "org.apache.flink:flink-table-common:2.3.0"

참고: 우리의 Gradle 빌드 스크립트나 quickstart 스크립트로 프로젝트를 만들었다고 가정해요. 자세한 내용은 Project configuration을 참고해요.

DataStream API와 브리지해야 하는 커넥터를 개발하려면(즉 DataStream 커넥터를 Table API에 맞게 적응시키려면) 이 의존성을 추가해야 해요.

Maven:

<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-java-bridge</artifactId>
<version>2.3.0</version>
<scope>provided</scope>
</dependency>

Gradle:

runtime "org.apache.flink:flink-table-api-java-bridge:2.3.0"

커넥터/포맷을 개발할 때 얇은(thhin) JAR과 uber JAR을 모두 배포하는 것을 제안해요. 그래야 사용자가 SQL 클라이언트나 Flink 배포에서 uber JAR을 쉽게 로드해 사용할 수 있어요. uber JAR은 위에 나열된 테이블 의존성을 제외한 커넥터의 모든 제3자 의존성을 포함해야 해요.

프로덕션 코드에서 flink-table-planner_2.12에 의존해서는 안 돼요. Flink 1.15에 도입된 새 모듈 flink-table-planner-loader로 애플리케이션의 클래스패스는 더 이상 org.apache.flink.table.planner 클래스에 직접 접근할 수 없어요. org.apache.flink.table.planner 패키지와 하위 패키지 내부에서만 사용 가능한 기능이 필요하면 이슈를 여세요. 자세한 내용은 Anatomy of Table Dependencies를 참고해요.

확장 지점 (Extension Points)

이 섹션은 Flink의 테이블 커넥터 확장에 사용할 수 있는 인터페이스를 설명해요.

동적 테이블 팩토리 (Dynamic Table Factories)

동적 테이블 팩토리는 카탈로그와 세션 정보에서 외부 저장 시스템용 동적 테이블 커넥터를 구성하는 데 사용돼요. org.apache.flink.table.factories.DynamicTableSourceFactory를 구현해 DynamicTableSource를 만들 수 있고, org.apache.flink.table.factories.DynamicTableSinkFactory를 구현해 DynamicTableSink를 만들 수 있어요.

기본적으로 팩토리는 Java의 Service Provider Interface로 커넥터 옵션 값을 팩토리 식별자로 사용해 발견돼요. JAR 파일에서 새 구현에 대한 참조를 서비스 파일에 추가할 수 있어요:

META-INF/services/org.apache.flink.table.factories.Factory

프레임워크는 팩토리 식별자와 요청된 기본 클래스(예: DynamicTableSourceFactory)로 고유하게 식별되는 단일 일치 팩토리를 확인해요. 필요하면 카탈로그 구현이 팩토리 발견 과정을 우회할 수 있어요. 이를 위해 카탈로그는 org.apache.flink.table.catalog.Catalog#getFactory에서 요청된 기본 클래스를 구현하는 인스턴스를 반환해야 해요.

동적 테이블 소스 (Dynamic Table Source)

정의상 동적 테이블은 시간에 따라 바뀔 수 있어요. 동적 테이블을 읽을 때 내용은 다음 중 하나로 간주될 수 있어요:

  • 변경 로그(changelog, 유한 또는 무한)로 모든 변경이 로그가 다할 때까지 지속적으로 소비됨. 이는 ScanTableSource 인터페이스로 표현돼요.
  • 지속적으로 변하거나 매우 큰 외부 테이블로 내용이 보통 전체적으로 읽히지 않고 필요할 때 개별 값으로 조회됨. 이는 LookupTableSource 인터페이스로 표현돼요.
  • 벡터로 검색을 지원하는 테이블. 이는 VectorSearchTableSource 인터페이스로 표현돼요.

클래스는 동시에 이 모든 인터페이스를 구현할 수 있어요. 플래너는 지정된 질의에 따라 사용을 결정해요.

Scan Table Source

ScanTableSource는 런타임에 외부 저장 시스템의 모든 행을 스캔해요. 스캔된 행은 삽입만 포함할 필요가 없고 갱신과 삭제도 포함할 수 있어요. 따라서 테이블 소스는 (유한 또는 무한) changelog를 읽는 데 사용될 수 있어요. 반환된 changelog mode는 런타임 중에 플래너가 기대할 수 있는 변경 집합을 나타내요.

  • 일반 배치 시나리오의 경우 소스는 삽입 전용 행의 유계 스트림을 방출할 수 있어요.
  • 일반 스트리밍 시나리오의 경우 소스는 삽입 전용 행의 무계 스트림을 방출할 수 있어요.
  • CDC(change data capture) 시나리오의 경우 소스는 삽입, 갱신, 삭제 행이 있는 유계 또는 무계 스트림을 방출할 수 있어요.

테이블 소스는 SupportsProjectionPushDown 같은 추가 능력 인터페이스를 구현할 수 있으며, 이는 계획 중 인스턴스를 변경할 수 있어요. 모든 능력은 org.apache.flink.table.connector.source.abilities 패키지에서 찾을 수 있고 source abilities table에 나열돼요.

반환된 scan runtime provider는 데이터 읽기용 런타임 구현을 제공해요. 런타임 구현에는 여러 인터페이스가 있으며 그중 SourceProvider는 권장 핵심 인터페이스예요. 제공자 인터페이스와 무관하게 소스 런타임 구현은 내부 데이터 구조를 생성해야 해요. 따라서 레코드는 org.apache.flink.table.data.RowData로 방출되어야 해요. 프레임워크는 소스가 여전히 공통 데이터 구조로 동작하고 마지막에 변환을 수행할 수 있도록 런타임 변환기를 제공해요.

병렬도 설정을 지원하려면 동적 테이블 팩토리가 org.apache.flink.table.factories.FactoryUtil에 정의된 선택적 scan.parallelism 옵션을 지원하고 그 값을 ParallelismProvider 인터페이스도 구현하는 제공자에 전달해야 해요.

Lookup Table Source

LookupTableSource는 런타임에 하나 이상의 키로 외부 저장 시스템의 행을 조회해요. ScanTableSource와 비교해 소스는 전체 테이블을 읽을 필요가 없고 필요할 때 (지속적으로 변할 수 있는) 외부 테이블에서 개별 값을 느리게 가져올 수 있어요. ScanTableSource와 달리 LookupTableSource는 현재 삽입 전용 변경만 방출하는 것을 지원해요. 추가 능력은 지원되지 않아요. 자세한 내용은 org.apache.flink.table.connector.source.LookupTableSource 문서를 참고해요.

LookupTableSource의 런타임 구현은 TableFunction 또는 AsyncTableFunction이에요. 런타임 중 주어진 조회 키 값으로 함수가 호출돼요.

Vector Search Table Source

VectorSearchTableSource는 런타임에 입력 벡터를 사용해 외부 저장 시스템을 검색하고 가장 유사한 top-K 행을 반환해요. 사용자는 입력 데이터와 외부 시스템에 저장된 데이터 사이의 유사도를 계산하는 데 사용할 알고리즘을 결정할 수 있어요. 일반적으로 대부분의 벡터 데이터베이스는 유사도 계산에 유클리드 거리(Euclidean distance)나 코사인 거리(Cosine distance)를 사용하는 것을 지원해요.

ScanTableSource와 비교해 소스는 전체 테이블을 읽을 필요가 없고 필요할 때 (지속적으로 변할 수 있는) 외부 테이블에서 개별 값을 느리게 가져올 수 있어요. ScanTableSource와 달리 VectorSearchTableSource는 현재 삽입 전용 변경만 방출하는 것을 지원해요. LookupTableSource와 달리 VectorSearchTableSource는 행 일치 여부를 결정할 때 동등성(equality)을 사용하지 않아요. 추가 능력은 지원되지 않아요. 자세한 내용은 org.apache.flink.table.connector.source.VectorSearchTableSource 문서를 참고해요.

VectorSearchTableSource의 런타임 구현은 TableFunction 또는 AsyncTableFunction이에요. 런타임 중 주어진 벡터 값으로 함수가 호출돼요.

소스 능력 (Source Abilities)

인터페이스 (Interface) 설명 (Description)
SupportsFilterPushDown DynamicTableSource로 필터를 밀어 내리는(push down) 것을 활성화. 효율성을 위해 소스는 실제 데이터 생성에 가깝도록 필터를 더 아래로 밀어낼 수 있어요.
SupportsLimitPushDown DynamicTableSource로 limit(생성될 예상 최대 레코드 수)을 밀어 내리는 것을 활성화.
SupportsPartitionPushDown 사용 가능한 파티션을 플래너에 전달하고 DynamicTableSource로 파티션을 밀어 내리는 것을 활성화. 런타임 중 소스는 효율성을 위해 전달된 파티션 목록에서만 데이터를 읽어요.
SupportsProjectionPushDown DynamicTableSource로 (중첩될 수 있는) 프로젝션을 밀어 내리는 것을 활성화. 효율성을 위해 소스는 프로젝션을 더 아래로 밀어 실제 데이터 생성에 가깝게 할 수 있어요. 소스가 SupportsReadingMetadata도 구현하면 소스는 필요한 메타데이터만 읽어요.
SupportsReadingMetadata DynamicTableSource에서 메타데이터 컬럼을 읽는 것을 활성화. 소스는 생성된 행의 끝에 필요한 메타데이터를 추가할 책임이 있어요. 포함된 포맷에서 메타데이터 컬럼을 잠재적으로 전달하는 것을 포함해요.
SupportsWatermarkPushDown DynamicTableSource로 워터마크 전략을 밀어 내리는 것을 활성화. 워터마크 전략은 타임스탬프 추출과 워터마크 생성용 빌더/팩토리예요. 런타임 중 워터마크 생성기는 소스 내부에 있으며 파티션별 워터마크를 생성할 수 있어요.
SupportsSourceWatermark ScanTableSource 자체가 제공하는 워터마크 전략에 완전히 의존하는 것을 활성화. 따라서 CREATE TABLE DDL은 플래너가 감지하고 사용 가능하면 이 인터페이스로의 호출로 번역하는 내장 마커 함수인 SOURCE_WATERMARK()를 사용할 수 있어요.
SupportsRowLevelModificationScan 행 수준 수정을 지원하는 스캔.

주의: 위 인터페이스들은 현재 ScanTableSource에만 사용할 수 있으며 LookupTableSource나 VectorSearchTableSource에는 사용할 수 없어요.

동적 테이블 싱크 (Dynamic Table Sink)

정의상 동적 테이블은 시간에 따라 바뀔 수 있어요. 동적 테이블을 쓸 때 내용은 항상 changelog(유한 또는 무한)로 간주될 수 있으며, 모든 변경이 로그가 다할 때까지 지속적으로 쓰여져요. 반환된 changelog mode는 런타임 중 싱크가 받아들이는 변경 집합을 나타내요.

  • 일반 배치 시나리오의 경우 싱크는 삽입 전용 행만 받아들이고 유계 스트림을 쓸 수 있어요.
  • 일반 스트리밍 시나리오의 경우 싱크는 삽입 전용 행만 받아들이고 무계 스트림을 쓸 수 있어요.
  • CDC 시나리오의 경우 싱크는 삽입, 갱신, 삭제 행이 있는 유계 또는 무계 스트림을 쓸 수 있어요.

테이블 싱크는 SupportsOverwrite 같은 추가 능력 인터페이스를 구현할 수 있으며, 이는 계획 중 인스턴스를 변경할 수 있어요. 모든 능력은 org.apache.flink.table.connector.sink.abilities 패키지에서 찾을 수 있고 sink abilities table에 나열돼요.

반환된 sink runtime provider는 데이터 쓰기용 런타임 구현을 제공해요. 런타임 구현에는 여러 인터페이스가 있으며 그중 SinkV2Provider는 권장 핵심 인터페이스예요. 제공자 인터페이스와 무관하게 싱크 런타임 구현은 내부 데이터 구조를 소비해야 해요. 따라서 레코드는 org.apache.flink.table.data.RowData로 받아들여야 해요. 프레임워크는 싱크가 여전히 공통 데이터 구조로 동작하고 시작 부분에서 변환을 수행할 수 있도록 런타임 변환기를 제공해요.

병렬도 설정을 지원하려면 동적 테이블 팩토리가 org.apache.flink.table.factories.FactoryUtil에 정의된 선택적 sink.parallelism 옵션을 지원하고 그 값을 ParallelismProvider 인터페이스도 구현하는 제공자에 전달해야 해요.

싱크 능력 (Sink Abilities)

인터페이스 (Interface) 설명 (Description)
SupportsOverwrite DynamicTableSink에서 기존 데이터를 덮어쓰는 것을 활성화. 기본적으로 이 인터페이스가 구현되지 않으면 SQL INSERT OVERWRITE 절 같은 것으로 기존 테이블이나 파티션을 덮어쓸 수 없어요.
SupportsPartitioning DynamicTableSink에서 파티션 데이터를 쓰는 것을 활성화.
SupportsBucketing DynamicTableSink에 버케팅(bucketing)을 활성화.
SupportsWritingMetadata DynamicTableSink로 메타데이터 컬럼을 쓰는 것을 활성화. 테이블 싱크는 소비된 행의 끝에 요청된 메타데이터 컬럼을 받아 영속화할 책임이 있어요. 포함된 포맷으로 메타데이터 컬럼을 전달하는 것을 포함해요.
SupportsDeletePushDown DELETE 문의 WHERE 절에서 분해된 필터를 DynamicTableSink로 밀어 내리는 것을 활성화. 테이블 싱크는 필터에 따라 기존 데이터를 직접 삭제할 수 있어요.
SupportsRowLevelDelete DynamicTableSink에서 행 수준 변경에 따라 기존 데이터를 삭제하는 것을 활성화. 테이블 싱크는 플래너에게 행 변경을 어떻게 생성할지 알려주고 이를 소비해 행 삭제 목적을 달성할 책임이 있어요.
SupportsRowLevelUpdate DynamicTableSink에서 행 수준 변경에 따라 기존 데이터를 갱신하는 것을 활성화. 테이블 싱크는 플래너에게 행 변경을 어떻게 생성할지 알려주고 이를 소비해 행 갱신 목적을 달성할 책임이 있어요.
SupportsStaging DynamicTableSink에서 CTAS(CREATE TABLE AS SELECT) 또는 RTAS([CREATE OR] REPLACE TABLE AS SELECT)의 원자적 의미론을 지원하는 것을 활성화. 테이블 싱크는 원자적 의미론을 제공하는 StagedTable 객체를 반환할 책임이 있어요.

인코딩/디코딩 포맷 (Encoding / Decoding Formats)

일부 테이블 커넥터는 키 및/또는 값을 인코딩·디코딩하는 서로 다른 포맷을 받아들여요. 포맷은 DynamicTableSourceFactory → DynamicTableSource → ScanRuntimeProvider 패턴과 유사하게 동작하며, 팩토리가 옵션 번역을 담당하고 소스가 런타임 로직 생성을 담당해요.

포맷이 다른 모듈에 있을 수 있으므로 테이블 팩토리와 유사하게 Java의 Service Provider Interface로 발견돼요. 포맷 팩토리를 발견하기 위해 동적 테이블 팩토리는 팩토리 식별자와 커넥터 특정 기본 클래스에 해당하는 팩토리를 검색해요.

예를 들어 Kafka 테이블 소스는 디코딩 포맷용 런타임 인터페이스로 DeserializationSchema를 요구해요. 따라서 Kafka 테이블 소스 팩토리는 value.format 옵션 값을 사용해 DeserializationFormatFactory를 발견해요. 현재 다음 포맷 팩토리가 지원돼요:

org.apache.flink.table.factories.DeserializationFormatFactory
org.apache.flink.table.factories.SerializationFormatFactory

포맷 팩토리는 옵션을 EncodingFormat 또는 DecodingFormat으로 번역해요. 이 인터페이스들은 주어진 데이터 타입에 대해 특화된 포맷 런타임 로직을 생성하는 또 다른 종류의 팩토리예요. 예를 들어 Kafka 테이블 소스 팩토리의 경우 DeserializationFormatFactory가 Kafka 테이블 소스에 전달될 수 있는 EncodingFormat<DeserializationSchema>을 반환해요.

전체 스택 예제 (Full Stack Example)

이 섹션은 changelog 의미론을 지원하는 디코딩 포맷이 있는 scan table source를 구현하는 방법을 개략적으로 설명해요. 예제는 언급된 모든 구성 요소가 어떻게 함께 작동하는지 보여줍니다. 참조 구현 역할을 할 수 있어요. 특히 다음을 어떻게 하는지 보여줘요:

  • 옵션을 파싱·검증하는 팩토리 생성
  • 테이블 커넥터 구현
  • 커스텀 포맷 구현·발견
  • 데이터 구조 변환기와 FactoryUtil 같은 제공 유틸리티 사용

테이블 소스는 들어오는 바이트를 수신하는 소켓을 여는 간단한 단일 스레드 SourceFunction을 사용해요. 원시 바이트는 플러그인 가능한 포맷으로 행으로 디코딩돼요. 포맷은 첫 번째 컬럼으로 changelog 플래그를 기대해요.

다음 DDL을 활성화하기 위해 위에서 언급한 대부분의 인터페이스를 사용할게요:

CREATE TABLE UserScores (name STRING, score INT)
WITH (
'connector' = 'socket',
'hostname' = 'localhost',
'port' = '9999',
'byte-delimiter' = '10',
'format' = 'changelog-csv',
'changelog-csv.column-delimiter' = '|'
);

포맷이 changelog 의미론을 지원하므로 런타임 중 갱신을 수집하고 변화하는 데이터를 지속적으로 평가할 수 있는 갱신 뷰(updating view)를 만들 수 있어요:

SELECT name, SUM(score) FROM UserScores GROUP BY name;

터미널에서 데이터를 수집하려면 다음 명령을 사용해요:

> nc -lk 9999
INSERT|Alice|12
INSERT|Bob|5
DELETE|Alice|12
INSERT|Alice|18

팩토리 (Factories)

이 섹션은 카탈로그에서 오는 메타데이터를 구체 커넥터 인스턴스로 번역하는 방법을 설명해요. 두 팩토리 모두 META-INF/services 디렉터리에 추가됐어요.

SocketDynamicTableFactory

SocketDynamicTableFactory는 카탈로그 테이블을 테이블 소스로 번역해요. 테이블 소스에 디코딩 포맷이 필요하므로 편의상 제공된 FactoryUtil을 사용해 포맷을 발견해요.

import org.apache.flink.api.common.serialization.DeserializationSchema;
import org.apache.flink.configuration.ConfigOption;
import org.apache.flink.configuration.ConfigOptions;
import org.apache.flink.configuration.ReadableConfig;
import org.apache.flink.table.connector.format.DecodingFormat;
import org.apache.flink.table.connector.source.DynamicTableSource;
import org.apache.flink.table.data.RowData;
import org.apache.flink.table.factories.DeserializationFormatFactory;
import org.apache.flink.table.factories.DynamicTableSourceFactory;
import org.apache.flink.table.factories.FactoryUtil;
import org.apache.flink.table.types.DataType;

import java.util.HashSet;
import java.util.Set;

public class SocketDynamicTableFactory implements DynamicTableSourceFactory {

// define all options statically
public static final ConfigOption<String> HOSTNAME = ConfigOptions.key("hostname")
.stringType()
.noDefaultValue();

public static final ConfigOption<Integer> PORT = ConfigOptions.key("port")
.intType()
.noDefaultValue();

public static final ConfigOption<Integer> BYTE_DELIMITER = ConfigOptions.key("byte-delimiter")
.intType()
.defaultValue(10); // corresponds to '\n'

@Override
public String factoryIdentifier() {
return "socket"; // used for matching to `connector = '...'`
}

@Override
public Set<ConfigOption<?>> requiredOptions() {
final Set<ConfigOption<?>> options = new HashSet<>();
options.add(HOSTNAME);
options.add(PORT);
options.add(FactoryUtil.FORMAT); // use pre-defined option for format
return options;
}

@Override
public Set<ConfigOption<?>> optionalOptions() {
final Set<ConfigOption<?>> options = new HashSet<>();
options.add(BYTE_DELIMITER);
return options;
}

@Override
public DynamicTableSource createDynamicTableSource(Context context) {
// either implement your custom validation logic here ...
// or use the provided helper utility
final FactoryUtil.TableFactoryHelper helper = FactoryUtil.createTableFactoryHelper(this, context);

// discover a suitable decoding format
final DecodingFormat<DeserializationSchema<RowData>> decodingFormat = helper.discoverDecodingFormat(
DeserializationFormatFactory.class,
FactoryUtil.FORMAT);

// validate all options
helper.validate();

// get the validated options
final ReadableConfig options = helper.getOptions();
final String hostname = options.get(HOSTNAME);
final int port = options.get(PORT);
final byte byteDelimiter = (byte) (int) options.get(BYTE_DELIMITER);

// derive the produced data type (excluding computed columns) from the catalog table
final DataType producedDataType =
context.getCatalogTable().getResolvedSchema().toPhysicalRowDataType();

// create and return dynamic table source
return new SocketDynamicTableSource(hostname, port, byteDelimiter, decodingFormat, producedDataType);
}
}

ChangelogCsvFormatFactory

ChangelogCsvFormatFactory는 포맷 특정 옵션을 포맷으로 번역해요. SocketDynamicTableFactory의 FactoryUtil이 옵션 키를 그에 맞게 적응시키고 changelog-csv.column-delimiter 같은 접두사를 처리해요. 이 팩토리는 DeserializationFormatFactory를 구현하므로 Kafka 커넥터 같은 역직렬화 포맷을 지원하는 다른 커넥터에도 사용될 수 있어요.

import org.apache.flink.api.common.serialization.DeserializationSchema;
import org.apache.flink.configuration.ConfigOption;
import org.apache.flink.configuration.ConfigOptions;
import org.apache.flink.configuration.ReadableConfig;
import org.apache.flink.table.connector.format.DecodingFormat;
import org.apache.flink.table.data.RowData;
import org.apache.flink.table.factories.FactoryUtil;
import org.apache.flink.table.factories.DeserializationFormatFactory;
import org.apache.flink.table.factories.DynamicTableFactory;

import java.util.Collections;
import java.util.HashSet;
import java.util.Set;

public class ChangelogCsvFormatFactory implements DeserializationFormatFactory {

// define all options statically
public static final ConfigOption<String> COLUMN_DELIMITER = ConfigOptions.key("column-delimiter")
.stringType()
.defaultValue("|");

@Override
public String factoryIdentifier() {
return "changelog-csv";
}

@Override
public Set<ConfigOption<?>> requiredOptions() {
return Collections.emptySet();
}

@Override
public Set<ConfigOption<?>> optionalOptions() {
final Set<ConfigOption<?>> options = new HashSet<>();
options.add(COLUMN_DELIMITER);
return options;
}

@Override
public DecodingFormat<DeserializationSchema<RowData>> createDecodingFormat(
DynamicTableFactory.Context context,
ReadableConfig formatOptions) {
// either implement your custom validation logic here ...
// or use the provided helper method
FactoryUtil.validateFactoryOptions(this, formatOptions);

// get the validated options
final String columnDelimiter = formatOptions.get(COLUMN_DELIMITER);

// create and return the format
return new ChangelogCsvFormat(columnDelimiter);
}
}

테이블 소스와 디코딩 포맷 (Table Source and Decoding Format)

이 섹션은 계획 계층의 인스턴스를 클러스터로 전달되는 런타임 인스턴스로 번역하는 방법을 설명해요.

SocketDynamicTableSource

SocketDynamicTableSource는 계획 중에 사용돼요. 이 예제에서는 사용 가능한 능력 인터페이스를 구현하지 않아요. 따라서 주요 로직은 런타임용으로 필요한 SourceFunction과 그 DeserializationSchema를 인스턴스화하는 getScanRuntimeProvider(...)에 있어요. 두 인스턴스 모두 내부 데이터 구조(즉 RowData)를 반환하도록 파라미터화돼요.

import org.apache.flink.api.common.serialization.DeserializationSchema;
import org.apache.flink.legacy.table.connector.source.SourceFunctionProvider;
import org.apache.flink.streaming.api.functions.source.SourceFunction;
import org.apache.flink.table.connector.ChangelogMode;
import org.apache.flink.table.connector.format.DecodingFormat;
import org.apache.flink.table.connector.source.DynamicTableSource;
import org.apache.flink.table.connector.source.ScanTableSource;
import org.apache.flink.table.data.RowData;
import org.apache.flink.table.types.DataType;

public class SocketDynamicTableSource implements ScanTableSource {

private final String hostname;
private final int port;
private final byte byteDelimiter;
private final DecodingFormat<DeserializationSchema<RowData>> decodingFormat;
private final DataType producedDataType;

public SocketDynamicTableSource(
String hostname,
int port,
byte byteDelimiter,
DecodingFormat<DeserializationSchema<RowData>> decodingFormat,
DataType producedDataType) {
this.hostname = hostname;
this.port = port;
this.byteDelimiter = byteDelimiter;
this.decodingFormat = decodingFormat;
this.producedDataType = producedDataType;
}

@Override
public ChangelogMode getChangelogMode() {
// in our example the format decides about the changelog mode
// but it could also be the source itself
return decodingFormat.getChangelogMode();
}

@Override
public ScanRuntimeProvider getScanRuntimeProvider(ScanContext runtimeProviderContext) {

// create runtime classes that are shipped to the cluster

final DeserializationSchema<RowData> deserializer = decodingFormat.createRuntimeDecoder(
runtimeProviderContext,
producedDataType);

final SourceFunction<RowData> sourceFunction = new SocketSourceFunction(
hostname,
port,
byteDelimiter,
deserializer);

return SourceFunctionProvider.of(sourceFunction, false);
}

@Override
public DynamicTableSource copy() {
return new SocketDynamicTableSource(hostname, port, byteDelimiter, decodingFormat, producedDataType);
}

@Override
public String asSummaryString() {
return "Socket Table Source";
}
}

ChangelogCsvFormat

ChangelogCsvFormat은 런타임에 DeserializationSchema를 사용하는 디코딩 포맷이에요. INSERT와 DELETE 변경 방출을 지원해요.

import org.apache.flink.api.common.serialization.DeserializationSchema;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.table.connector.ChangelogMode;
import org.apache.flink.table.connector.format.DecodingFormat;
import org.apache.flink.table.connector.source.DynamicTableSource;
import org.apache.flink.table.connector.source.DynamicTableSource.DataStructureConverter;
import org.apache.flink.table.data.RowData;
import org.apache.flink.table.types.DataType;
import org.apache.flink.table.types.logical.LogicalType;
import org.apache.flink.types.RowKind;

import java.util.List;

public class ChangelogCsvFormat implements DecodingFormat<DeserializationSchema<RowData>> {

private final String columnDelimiter;

public ChangelogCsvFormat(String columnDelimiter) {
this.columnDelimiter = columnDelimiter;
}

@Override
@SuppressWarnings("unchecked")
public DeserializationSchema<RowData> createRuntimeDecoder(
DynamicTableSource.Context context,
DataType producedDataType) {
// create type information for the DeserializationSchema
final TypeInformation<RowData> producedTypeInfo = context.createTypeInformation(producedDataType);

// most of the code in DeserializationSchema will not work on internal data structures
// create a converter for conversion at the end
final DataStructureConverter converter = context.createDataStructureConverter(producedDataType);

// use logical types during runtime for parsing
final List<LogicalType> parsingTypes = producedDataType.getLogicalType().getChildren();

// create runtime class
return new ChangelogCsvDeserializer(parsingTypes, converter, producedTypeInfo, columnDelimiter);
}

@Override
public ChangelogMode getChangelogMode() {
// define that this format can produce INSERT and DELETE rows
return ChangelogMode.newBuilder()
.addContainedKind(RowKind.INSERT)
.addContainedKind(RowKind.DELETE)
.build();
}
}

런타임 (Runtime)

완전성을 위해 이 섹션은 SourceFunction과 DeserializationSchema의 런타임 로직을 설명해요.

ChangelogCsvDeserializer

ChangelogCsvDeserializer는 바이트를 행 종류(row kind)를 가진 Integer와 String의 Row로 변환하는 간단한 파싱 로직을 포함해요. 최종 변환 단계는 그것들을 내부 데이터 구조로 변환해요.

import org.apache.flink.api.common.serialization.DeserializationSchema;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.table.connector.RuntimeConverter.Context;
import org.apache.flink.table.connector.source.DynamicTableSource.DataStructureConverter;
import org.apache.flink.table.data.RowData;
import org.apache.flink.table.types.logical.LogicalType;
import org.apache.flink.table.types.logical.LogicalTypeRoot;
import org.apache.flink.types.Row;
import org.apache.flink.types.RowKind;

import java.util.List;
import java.util.regex.Pattern;

public class ChangelogCsvDeserializer implements DeserializationSchema<RowData> {

private final List<LogicalType> parsingTypes;
private final DataStructureConverter converter;
private final TypeInformation<RowData> producedTypeInfo;
private final String columnDelimiter;

public ChangelogCsvDeserializer(
List<LogicalType> parsingTypes,
DataStructureConverter converter,
TypeInformation<RowData> producedTypeInfo,
String columnDelimiter) {
this.parsingTypes = parsingTypes;
this.converter = converter;
this.producedTypeInfo = producedTypeInfo;
this.columnDelimiter = columnDelimiter;
}

@Override
public TypeInformation<RowData> getProducedType() {
// return the type information required by Flink's core interfaces
return producedTypeInfo;
}

@Override
public void open(InitializationContext context) {
// converters must be open
converter.open(Context.create(ChangelogCsvDeserializer.class.getClassLoader()));
}

@Override
public RowData deserialize(byte[] message) {
// parse the columns including a changelog flag
final String[] columns = new String(message).split(Pattern.quote(columnDelimiter));
final RowKind kind = RowKind.valueOf(columns[0]);
final Row row = new Row(kind, parsingTypes.size());
for (int i = 0; i < parsingTypes.size(); i++) {
row.setField(i, parse(parsingTypes.get(i).getTypeRoot(), columns[i + 1]));
}
// convert to internal data structure
return (RowData) converter.toInternal(row);
}

private static Object parse(LogicalTypeRoot root, String value) {
switch (root) {
case INTEGER:
return Integer.parseInt(value);
case VARCHAR:
return value;
default:
throw new IllegalArgumentException();
}
}

@Override
public boolean isEndOfStream(RowData nextElement) {
return false;
}
}

SocketSourceFunction

SocketSourceFunction은 소켓을 열고 바이트를 소비해요. 주어진 바이트 구분자(기본 \n)로 레코드를 나누고 디코딩을 플러그인 가능한 DeserializationSchema에 위임해요. 소스 함수는 병렬도 1로만 동작할 수 있어요.

import org.apache.flink.api.common.serialization.DeserializationSchema;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.typeutils.ResultTypeQueryable;
import org.apache.flink.streaming.api.functions.source.RichSourceFunction;
import org.apache.flink.table.data.RowData;

import java.io.ByteArrayOutputStream;
import java.io.InputStream;
import java.net.InetSocketAddress;
import java.net.Socket;

public class SocketSourceFunction extends RichSourceFunction<RowData> implements ResultTypeQueryable<RowData> {

private final String hostname;
private final int port;
private final byte byteDelimiter;
private final DeserializationSchema<RowData> deserializer;

private volatile boolean isRunning = true;
private Socket currentSocket;

public SocketSourceFunction(String hostname, int port, byte byteDelimiter, DeserializationSchema<RowData> deserializer) {
this.hostname = hostname;
this.port = port;
this.byteDelimiter = byteDelimiter;
this.deserializer = deserializer;
}

@Override
public TypeInformation<RowData> getProducedType() {
return deserializer.getProducedType();
}

@Override
public void run(SourceContext<RowData> ctx) throws Exception {
while (isRunning) {
// open and consume from socket
try (final Socket socket = new Socket()) {
currentSocket = socket;
socket.connect(new InetSocketAddress(hostname, port), 0);
try (InputStream stream = socket.getInputStream()) {
ByteArrayOutputStream buffer = new ByteArrayOutputStream();
int b;
while ((b = stream.read()) >= 0) {
// buffer until delimiter
if (b != byteDelimiter) {
buffer.write(b);
}
// decode and emit record
else {
ctx.collect(deserializer.deserialize(buffer.toByteArray()));
buffer.reset();
}
}
}
} catch (Throwable t) {
t.printStackTrace(); // print and continue
}
Thread.sleep(1000);
}
}

@Override
public void cancel() {
isRunning = false;
try {
currentSocket.close();
} catch (Throwable t) {
// ignore
}
}
}

더 알아보기 (Learn more)