Dynamic Kafka SQL 커넥터

Dynamic Kafka SQL 커넥터

Dynamic Kafka 커넥터를 사용하면 잡을 재시작하지 않고 클러스터를 이동할 수 있는 Kafka 토픽에서 데이터를 읽을 수 있어요. 스트림은 Kafka 메타데이터 서비스를 통해 해석됩니다.

출처: Dynamic Kafka SQL Connector

본문

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를 참고하세요.

더 알아보기 (Learn more)