Hybrid Source
Hybrid Source (하이브리드 소스)
HybridSource는 구체적인 소스(sources) 목록을 포함하는 소스입니다. 이종(heterogeneous) 소스들의 입력을 순차적으로 읽어 단일 입력 스트림을 만들어내는 문제를 해결합니다. 예를 들어 부트스트랩(bootstrap) 사용 사례에서는 S3에서 며칠 분량의 유한(bounded) 입력을 먼저 읽은 뒤 Kafka에서 최신의 무한(unbounded) 입력을 계속 이어서 읽어야 할 수 있습니다.
출처: 문서
본문
HybridSource는 구체적인 소스 목록을 포함하는 소스입니다. 이종 소스로부터의 입력을 순차적으로 읽어 단일 입력 스트림을 만들어내는 문제를 해결합니다.
예를 들어 부트스트랩 사용 사례에서는 S3에서 며칠 분량의 유한 입력을 먼저 읽은 뒤 Kafka의 최신 무한 입력을 이어서 읽어야 할 수 있습니다. HybridSource는 유한 파일 입력이 끝나면 애플리케이션을 중단하지 않고 FileSource에서 KafkaSource로 전환합니다.
HybridSource 이전에는 여러 소스를 가진 토폴로지를 만들고 사용자 영역에서 전환 메커니즘을 정의해야 했으며, 이는 운영 복잡성과 비효율을 초래했습니다.
HybridSource를 사용하면 여러 소스가 Flink 작업 그래프와 DataStream API 관점에서 단일 소스로 나타납니다.
자세한 배경은 FLIP-150을 참조하세요.
커넥터를 사용하려면 프로젝트에 flink-connector-base 의존성을 추가하세요.
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-base</artifactId>
<version>2.3.0</version>
</dependency>
(보통 구체적인 소스의 전이적 의존성(transitive dependency)으로 포함됩니다.)
다음 소스의 시작 위치 (Start position for next source)
여러 소스를 HybridSource에 배치하려면 마지막 소스를 제외한 모든 소스가 유한(bounded)이어야 합니다. 따라서 소스들은 보통 시작 위치와 끝 위치를 할당받아야 합니다. 마지막 소스는 유한일 수 있으며, 이 경우 HybridSource는 유한하고 그렇지 않으면 무한합니다. 세부 사항은 특정 소스와 외부 저장 시스템에 따라 달라집니다.
여기서는 File/Kafka 예시를 따라 가장 기본적인 시나리오와 더 복잡한 시나리오를 다룹니다.
그래프 구성 시점의 고정 시작 위치 (Fixed start position at graph construction time)
예시: 미리 정해진 전환 시점까지 파일에서 읽고, 그 다음 Kafka에서 이어서 읽습니다. 각 소스는 사전에 알려진 범위를 다루므로 포함된 소스들을 직접 사용하는 것처럼 미리 만들 수 있습니다.
Java
long switchTimestamp = ...; // derive from file input paths
FileSource<String> fileSource =
FileSource.forRecordStreamFormat(new TextLineInputFormat(), Path.fromLocalFile(testDir)).build();
KafkaSource<String> kafkaSource =
KafkaSource.<String>builder()
.setStartingOffsets(OffsetsInitializer.timestamp(switchTimestamp + 1))
.build();
HybridSource<String> hybridSource =
HybridSource.builder(fileSource)
.addSource(kafkaSource)
.build();
Python
switch_timestamp = ... # derive from file input paths
file_source = FileSource \
.for_record_stream_format(StreamFormat.text_line_format(), test_dir) \
.build()
kafka_source = KafkaSource \
.builder() \
.set_bootstrap_servers('localhost:9092') \
.set_group_id('MY_GROUP') \
.set_topics('quickstart-events') \
.set_value_only_deserializer(SimpleStringSchema()) \
.set_starting_offsets(KafkaOffsetsInitializer.timestamp(switch_timestamp)) \
.build()
hybrid_source = HybridSource.builder(file_source).add_source(kafka_source).build()
전환 시점의 동적 시작 위치 (Dynamic start position at switch time)
예시: 파일 소스가 매우 큰 백로그를 읽으며, 잠재적으로 다음 소스에 사용 가능한 보존 기간(retention)보다 오래 걸릴 수 있습니다. 전환은 "현재 시간 - X"에 발생해야 합니다. 이는 다음 소스의 시작 시간이 전환 시점에 설정되어야 함을 의미합니다. 여기서는 SourceFactory를 구현해 이전 파일 enumerator의 끝 위치를 전달받아 KafkaSource를 지연 구성(deferred construction)해야 합니다.
enumerator가 끝 타임스탬프를 가져오는 것을 지원해야 한다는 점을 주의하세요. 현재는 소스 사용자 정의가 필요할 수 있습니다. FileSource에 동적 끝 위치 지원을 추가하는 것은 FLINK-23633에서 추적됩니다.
Java
FileSource<String> fileSource = CustomFileSource.readTillOneDayFromLatest();
HybridSource<String> hybridSource =
HybridSource.<String, CustomFileSplitEnumerator>builder(fileSource)
.addSource(
switchContext -> {
CustomFileSplitEnumerator previousEnumerator =
switchContext.getPreviousEnumerator();
// how to get timestamp depends on specific enumerator
long switchTimestamp = previousEnumerator.getEndTimestamp();
KafkaSource<String> kafkaSource =
KafkaSource.<String>builder()
.setStartingOffsets(OffsetsInitializer.timestamp(switchTimestamp + 1))
.build();
return kafkaSource;
},
Boundedness.CONTINUOUS_UNBOUNDED)
.build();
Python
Still not supported in Python API.