Dynamic Kafka SQL 커넥터
Dynamic Kafka SQL 커넥터
Dynamic Kafka 커넥터를 사용하면 잡을 재시작하지 않고 클러스터를 이동할 수 있는 Kafka 토픽에서 데이터를 읽을 수 있어요. 스트림은 Kafka 메타데이터 서비스를 통해 해석됩니다.
본문
Dynamic Kafka 커넥터는 잡을 재시작하지 않고 클러스터를 이동할 수 있는 Kafka 토픽에서 데이터를 읽는 것을 허용해요. 스트림은 Kafka 메타데이터 서비스를 통해 해석돼요. 이는 특히 클러스터 마이그레이션과 동적 토픽/클러스터 변경에 유용해요.
Dependencies
Flink 버전 2.3용 커넥터는 아직 제공되지 않아요.
Kafka 커넥터는 바이너리 배포판의 일부가 아니에요. 클러스터 실행을 위해 연결하는 방법은 여기를 참고하세요.
How to create a Dynamic Kafka table
아래 예시는 내장 단일-클러스터 메타데이터 서비스를 사용해 Dynamic Kafka 테이블을 만드는 방법을 보여줘요. 이 서비스로 스트림 id는 단일 Kafka 클러스터의 토픽으로 해석돼요.
CREATE TABLE DynamicKafkaTable (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING,
`kafka_cluster` STRING METADATA FROM 'kafka_cluster',
`ts` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp'
) WITH (
'connector' = 'dynamic-kafka',
'stream-ids' = 'user_behavior;user_behavior_v2',
'metadata-service' = 'single-cluster',
'metadata-service.cluster-id' = 'cluster-0',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'scan.startup.mode' = 'earliest-offset',
'format' = 'csv'
)
커넥터는 metadata-service 옵션을 통한 커스텀 메타데이터 서비스도 지원해요. 서비스 클래스는 KafkaMetadataService를 구현해야 하며, 공개 no-arg 생성자 또는 Properties를 받는 생성자를 가져야 해요. 커넥터는 가능할 때 컨스트럭터에 Kafka 속성(모든 properties.* 옵션)을 전달해요.
Available Metadata
Dynamic Kafka 커넥터는 Kafka 커넥터의 모든 메타데이터 컬럼을 노출하고 동적 소스에 특화된 메타데이터 컬럼 하나를 추가해요:
kafka_cluster(STRING NOT NULL, read-only): 레코드에 대해 메타데이터 서비스가 해석한 클러스터 id.
공유 메타데이터 컬럼은 Kafka SQL Connector를 참고하세요.
예시:
CREATE TABLE DynamicKafkaTable (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING,
`kafka_cluster` STRING METADATA FROM 'kafka_cluster' VIRTUAL
) WITH (
'connector' = 'dynamic-kafka',
'stream-ids' = 'user_behavior;user_behavior_v2',
'metadata-service' = 'single-cluster',
'metadata-service.cluster-id' = 'cluster-0',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'scan.startup.mode' = 'earliest-offset',
'format' = 'csv'
);
Connector Options
| Option | Required | Forwarded | Default | Type | Description |
|---|---|---|---|---|---|
| connector | required | no | (none) | String | 사용할 커넥터 지정. Dynamic Kafka는 'dynamic-kafka' 사용. |
| stream-ids | optional | no | (none) | String | 구독할 세미콜론으로 구분된 스트림 id. stream-ids와 stream-pattern 중 하나만 설정 가능. |
| stream-pattern | optional | no | (none) | String | 구독할 스트림 id의 정규식 패턴. stream-ids와 stream-pattern 중 하나만 설정 가능. |
| metadata-service | required | no | (none) | String | 메타데이터 서비스 식별자. 'single-cluster' 또는 KafkaMetadataService를 구현하는 완전 한정 클래스 이름 사용. |
| metadata-service.cluster-id | optional | no | (none) | String | single-cluster 메타데이터 서비스에 필요한 클러스터 id. |
| stream-metadata-discovery-interval-ms | optional | no | -1 | Long | 스트림 메타데이터 변경을 발견하는 간격(밀리초). 양수가 아닌 값은 발견을 비활성화. |
| stream-metadata-discovery-failure-threshold | optional | no | 1 | Integer | 잡을 실패시키기 전의 연속 발견 실패 횟수. |
커넥터는 Kafka 커넥터와 동일한 포맷 옵션과 Kafka 클라이언트 속성도 지원해요. 전체 포맷 옵션과 Kafka 속성(모든 properties.* 옵션) 목록은 Kafka SQL Connector를 참고하세요.