동적 Kafka 소스 `Experimental`
동적 Kafka 소스 Experimental (Dynamic Kafka Source)
Flink는 하나 이상의 Kafka 클러스터의 Kafka 토픽에서 데이터를 읽기 위한 Apache Kafka 커넥터를 제공합니다. Dynamic Kafka 커넥터는 Kafka 메타데이터 서비스를 사용해 클러스터와 토픽을 발견하며, 작업을 재시작할 필요 없이 토픽 및/또는 클러스터의 변경을 수용하며 동적으로 읽을 수 있습니다. 이는 새 Kafka 클러스터/토픽을 읽어야 하거나 기존 Kafka 클러스터/토픽 읽기를 중지해야 할 때(클러스터 마이그레이션/페이오버/기타 인프라 변경)와 Hybrid Source와의 직접 통합이 필요할 때 특히 유용합니다. 이 솔루션은 이러한 연산을 자동화해 Kafka 컨슈머에게 투명하게 만듭니다.
출처: 문서
본문
의존성 (Dependency)
Kafka 호환성에 대한 자세한 내용은 공식 Kafka 문서를 참고하세요.
현재 Flink 2.3 버전용 커넥터는 아직 없습니다.
Flink의 스트리밍 커넥터는 바이너리 배포의 일부가 아닙니다. 클러스터 실행을 위해 연결하는 방법은 여기를 참고하세요.
Dynamic Kafka 소스 (Dynamic Kafka Source)
이 부분은 새 데이터 소스 API에 기반한 Dynamic Kafka Source를 설명합니다.
사용법 (Usage)
Dynamic Kafka Source는 DynamicKafkaSource를 초기화하기 위한 빌더 클래스를 제공합니다. 아래 코드 스니펫은 "MyKafkaMetadataService"로 "input-stream"에 해당하는 클러스터와 토픽을 해석하면서, "input-stream" 스트림의 가장 이른 오프셋에서 메시지를 소비하고 ConsumerRecord의 값만 문자열로 역직렬화하는 DynamicKafkaSource를 빌드하는 방법을 보여줍니다.
Java
DynamicKafkaSource<String> source = DynamicKafkaSource.<String>builder()
.setKafkaMetadataService(new MyKafkaMetadataService())
.setStreamIds(Collections.singleton("input-stream"))
.setEnumeratorMode(DynamicKafkaSourceOptions.EnumeratorMode.PER_CLUSTER)
.setStartingOffsets(OffsetsInitializer.earliest())
.setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class))
.setProperties(properties)
.build();
env.fromSource(source, WatermarkStrategy.noWatermarks(), "Dynamic Kafka Source");
Python
metadata_service = SingleClusterTopicMetadataService(
"cluster-a",
{"bootstrap.servers": "localhost:9092"})
source = DynamicKafkaSource.builder() \
.set_kafka_metadata_service(metadata_service) \
.set_stream_ids({"input-stream"}) \
.set_starting_offsets(KafkaOffsetsInitializer.earliest()) \
.set_value_only_deserializer(SimpleStringSchema()) \
.set_properties(properties) \
.build()
env.from_source(source, WatermarkStrategy.no_watermarks(), "Dynamic Kafka Source")
DynamicKafkaSource를 빌드하는 데 다음 속성이 필요합니다:
setKafkaMetadataService(KafkaMetadataService)로 구성하는 Kafka 메타데이터 서비스- 구독할 스트림 id(자세한 내용은 아래 Kafka 스트림 구독 섹션 참고)
- Kafka 메시지를 파싱하는 역직렬화기(자세한 내용은 Kafka Source 문서 참고)
오프셋 초기화 (Offsets Initialization)
시작 및 중지 오프셋을 빌더를 통해 전역으로 구성할 수 있습니다. 시작 오프셋은 유한(bounded) 및 무한(unbounded) 소스 모두에 적용되는 반면, 중지 오프셋은 소스가 유한 모드로 실행될 때만 적용됩니다. 클러스터 메타데이터에는 선택적으로 클러스터별 시작 또는 중지 오프셋 초기화기가 포함될 수 있습니다. 존재하면 해당 클러스터의 전역 기본값을 재정의합니다.
예시: 메타데이터를 통해 특정 클러스터의 오프셋을 재정의합니다.
Java
Properties cluster0Props = new Properties();
cluster0Props.setProperty(
CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, "cluster0:9092");
Properties cluster1Props = new Properties();
cluster1Props.setProperty(
CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, "cluster1:9092");
KafkaStream stream =
new KafkaStream(
"input-stream",
Map.of(
"cluster0",
new ClusterMetadata(
Set.of("topic-a"),
cluster0Props,
OffsetsInitializer.earliest(),
OffsetsInitializer.latest()),
"cluster1",
new ClusterMetadata(
Set.of("topic-b"),
cluster1Props,
OffsetsInitializer.latest(),
null)));
DynamicKafkaSource<String> source =
DynamicKafkaSource.<String>builder()
.setStreamIds(Set.of(stream.getStreamId()))
.setKafkaMetadataService(new MockKafkaMetadataService(Set.of(stream)))
.setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class))
// Overridden by per-cluster starting offsets in metadata when present.
.setStartingOffsets(OffsetsInitializer.earliest())
.setBounded(OffsetsInitializer.latest())
.build();
스플릿 할당 모드 (Split Assignment Mode)
Dynamic Kafka Source는 두 가지 스플릿 할당 모드를 지원합니다:
per_cluster(기본값): 각 Kafka 클러스터 내에서 독립적으로 스플릿을 할당합니다.global: 하나의 전역 밸런싱 전략으로 발견된 모든 클러스터에 걸쳐 스플릿을 할당합니다.
모드는 빌더 API 또는 소스 속성으로 구성할 수 있습니다.
새로 발견된 스플릿의 다음 소유자(owner) 선택 방식:
per_cluster:KafkaSourceEnumerator와 같은 소유자 로직을 사용합니다. 토픽 파티션(topic, partition)에 대해numReaders = P일 때:startIndex = ((topic.hashCode() * 31) & 0x7FFFFFFF) % P,owner = (startIndex + partition) % P.global: 모든 클러스터에 걸쳐 하나의 전역 소유자 커서를 사용합니다:owner = knownActiveSplitIds.size() % numReaders, 그런 다음 스플릿 id를knownActiveSplitIds에 추가합니다. (addSplitsBack으로 스플릿이 반환되면, 유효할 때 이전 소유자를 선호해 재사용합니다.)
global 모드에서 밸런싱은 미래지향적입니다. 새로 발견된 스플릿은 미래 분포가 균형을 유지하도록 할당되며, 이미 할당된 활성 스플릿은 재밸런싱을 위해 선제적으로 마이그레이션되지 않습니다.
축소/제거로 global 할당이 치우치게 되어도 열거자(enumerator)는 이미 활성화된 스플릿을 스스로 재밸런싱하지 않습니다. 기존 소유권을 재밸런싱하려면 병렬도 변경과 함께 복원을 사용하세요(예: rescale 후 savepoint/checkpoint 복원). 그러면 Flink 런타임이 소스 리더 연산자 상태를 재파티셔닝합니다.
Java
DynamicKafkaSource<String> source =
DynamicKafkaSource.<String>builder()
.setKafkaMetadataService(new MyKafkaMetadataService())
.setStreamIds(Set.of("input-stream"))
.setEnumeratorMode(DynamicKafkaSourceOptions.EnumeratorMode.GLOBAL)
.setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class))
.build();
워터마크 정렬 (Watermark Alignment)
워터마크 정렬이 활성화되고 동적 소스 리더가 여러 스플릿을 소유하면, Flink는 개별 스플릿의 워터마크를 구성된 드리프트 내로 유지하기 위해 개별 스플릿을 일시 중지하거나 재개할 수 있습니다. Dynamic Kafka Source는 이러한 스플릿 일시 중지·재개 요청을 기본 Kafka 소스 리더로 전달하므로, 스플릿 레벨 워터마크 정렬이 일반 Kafka Source와 동일한 방식으로 동작합니다.
Kafka 스트림 구독 (Kafka Stream Subscription)
Dynamic Kafka Source는 Kafka 스트림 구독 방법 2가지를 제공합니다.
- Kafka 스트림 id 집합. 예를 들어:
Java
DynamicKafkaSource.builder().setStreamIds(Set.of("stream-a", "stream-b"));
Python
DynamicKafkaSource.builder().set_stream_ids({"stream-a", "stream-b"})
- 주어진 regex와 일치하는 모든 Kafka 스트림 id를 구독하는 regex 패턴. 예를 들어:
Java
DynamicKafkaSource.builder().setStreamPattern(Pattern.of("stream.*"));
Python
DynamicKafkaSource.builder().set_stream_pattern("stream.*")
Kafka 메타데이터 서비스 (Kafka Metadata Service)
논리적 Kafka 스트림을 해당 물리적 토픽과 클러스터로 해석하기 위한 인터페이스가 제공됩니다. 일반적으로 이러한 구현은 내부 Kafka 인프라와 잘 맞는 서비스에 기반합니다. 이를 사용할 수 없다면 인메모리 구현도 동작합니다. 인메모리 구현의 예는 테스트에서 찾을 수 있습니다.
이 소스는 주기적으로 이 Kafka 메타데이터 서비스를 폴링해 Kafka 스트림의 변경 사항을 감지하고, 서비스가 반환한 새 Kafka 메타데이터를 구독하도록 리더 태스크를 조정(reconcile)함으로써 동적 특성을 달성합니다. 예를 들어 Kafka 마이그레이션의 경우, 서비스가 Kafka 스트림 메타데이터에서 그 변경을 만들면 소스는 한 클러스터에서 새 클러스터로 전환합니다.
클러스터 메타데이터는 선택적으로 클러스터별 시작 및 중지 오프셋 초기화기를 운반할 수 있습니다. 이들은 영향을 받는 클러스터의 전역 빌더 구성을 재정의합니다.
추가 속성 (Additional Properties)
빌더를 통해 properties에 구성할 수 있는 DynamicKafkaSourceOptions의 구성 옵션이 있습니다:
| 옵션 | 필수 | 기본값 | 타입 | 설명 |
|---|---|---|---|---|
| stream-metadata-discovery-interval-ms | required | -1 | Long | 소스가 스트림 메타데이터의 변경을 발견하는 간격(밀리초). 양수가 아닌 값은 스트림 메타데이터 발견을 비활성화합니다. |
| stream-metadata-discovery-failure-threshold | required | 1 | Integer | Kafka 메타데이터 서비스 발견의 예외가 jobmanager 실패와 전역 페일오버를 유발하기 전의 연속 실패 횟수. 기본값은 1로 최소한 시작 실패를 잡아냅니다. |
| stream-enumerator-mode | required | per_cluster | String | 동적 Kafka 스플릿 할당을 위한 열거자 구현. 지원 값은 per_cluster(클러스터 로컬 할당)와 global(클러스터 간 전역 균형 할당)입니다. |
이 목록 외에 적용 가능한 속성 목록은 일반 Kafka 커넥터를 참고하세요.
메트릭 (Metrics)
| 범위 | 메트릭 | 사용자 변수 | 설명 | 타입 |
|---|---|---|---|---|
| Operator | currentEmitEventTimeLag | n/a | 레코드 이벤트 타임스탬프부터 소스 커넥터가 레코드를 내보낸 시간까지의 시간 간격: currentEmitEventTimeLag = EmitTime - EventTime. |
Gauge |
| watermarkLag | n/a | 워터마크가 벽시계 시간보다 뒤처진 시간 간격: watermarkLag = CurrentTime - Watermark |
Gauge | |
| sourceIdleTime | n/a | 소스가 어떤 레코드도 처리하지 않은 시간 간격: sourceIdleTime = CurrentTime - LastRecordProcessTime |
Gauge | |
| pendingRecords | n/a | 소스가 아직 가져오지 않은 레코드 수. 예: Kafka 파티션에서 컨슈머 오프셋 이후의 사용 가능 레코드. | Gauge | |
| kafkaClustersCount | n/a | 이 리더가 읽는 Kafka 클러스터의 총 수. | Gauge |
이 목록 외에도 보고되는 KafkaSourceReader 메트릭은 일반 Kafka 커넥터를 참고하세요.
추가 세부 사항 (Additional Details)
역직렬화, 이벤트 시간 및 워터마크, 유휴 상태, 컨슈머 오프셋 커밋, 보안 등에 대한 추가 세부 사항은 Kafka Source 문서를 참고할 수 있습니다. Dynamic Kafka Source가 Kafka Source의 컴포넌트를 활용하므로 가능하며, 구현은 다음 섹션에서 논의됩니다.
뒷면의 동작 (Behind the Scene)
새 데이터 소스 API 설계에서 Kafka 소스가 어떻게 동작하는지 관심이 있다면 이 부분을 참조로 읽는 것이 좋습니다. 새 데이터 소스 API에 대한 자세한 내용은 데이터 소스 문서와 FLIP-27이 더 설명적인 논의를 제공합니다.
새 데이터 소스 API 추상화 아래에서 Dynamic Kafka Source는 다음 컴포넌트로 구성됩니다:
소스 스플릿 (Source Split)
Dynamic Kafka Source의 소스 스플릿은 클러스터 정보가 있는 Kafka 토픽의 파티션을 나타냅니다. 다음으로 구성됩니다:
- Kafka 메타데이터 서비스로 해석할 수 있는 Kafka 클러스터 id.
- Kafka Source Split(TopicPartition, starting offset, stopping offset).
자세한 내용은 DynamicKafkaSourceSplit 클래스를 확인할 수 있습니다.
스플릿 열거자 (Split Enumerator)
이 열거자는 하나 이상의 클러스터에서 스플릿을 발견하고 할당하는 역할을 합니다. 시작 시 열거자는 Kafka 스트림 id에 속하는 메타데이터를 발견합니다. 메타데이터를 사용해 리더에게 스플릿을 할당하는 기능을 처리하는 KafkaSourceEnumerators를 초기화할 수 있습니다. 또한 소스 이벤트가 소스 리더로 전송되어 메타데이터를 조정합니다. 이 열거자는 스트림 발견을 위해 KafkaMetadataService를 주기적으로 폴링할 수 있습니다. 또한 클러스터가 제거될 수 있으므로 그 메트릭도 제거되어야 하므로, 메타데이터가 변경될 때 열거자를 재시작하는 것은 오래된 메트릭을 지우는 것을 수반합니다.
소스 리더 (Source Reader)
이 리더는 하나 이상의 클러스터에서 읽는 역할을 하며 KafkaSourceReader를 사용해 메타데이터에 기반해 토픽과 클러스터에서 레코드를 가져옵니다. 열거자가 새 메타데이터를 발견하면 리더는 메타데이터 변경을 조정하여 새 토픽과 클러스터 집합에서 읽도록 KafkaSourceReader를 재시작할 수도 있습니다.
Kafka 메타데이터 서비스 (Kafka Metadata Service)
이 인터페이스는 구성된 Kafka 스트림 id에 대한 현재 메타데이터의 진실 원천(source of truth)을 나타냅니다. 폴링 사이에 제거된 메타데이터는 비활성으로 간주됩니다(예: 반환 값에서 클러스터를 제거하면 그 클러스터는 비활성이며 읽지 않아야 함을 뜻함). 클러스터 메타데이터는 불변의 Kafka 클러스터 id, 토픽 집합, Kafka 클러스터에 연결하는 데 필요한 속성을 포함합니다.
FLIP 246
더 많은 뒷면 동작을 이해하려면 자세한 내용과 논의를 위해 FLIP-246을 읽으세요.