데이터 소스
데이터 소스 (Data Sources)
이 페이지는 Flink의 Data Source API와 그 배경이 되는 개념·아키텍처를 설명해요. Flink에서 데이터 소스가 어떻게 동작하는지, 또는 새로운 Data Source를 어떻게 구현하는지 알고 싶다면 이 문서를 읽어보세요. 미리 정의된 소스 커넥터를 찾고 있다면 Connector Docs를 확인하세요.
출처: 문서
본문
Data Source 개념
핵심 구성 요소
Data Source에는 세 가지 핵심 구성 요소가 있어요. 바로 Splits, SplitEnumerator, SourceReader입니다.
- Split은 소스가 소비하는 데이터의 일부로, 파일이나 로그 파티션 같은 것을 말해요. Split은 소스가 작업을 분배하고 데이터 읽기를 병렬화하는 단위예요.
- SourceReader는 Split을 요청하고 처리해요. 예를 들어 Split이 나타내는 파일이나 로그 파티션을 읽는 방식이에요. SourceReader들은 Task Manager에서 SourceOperators 안에 병렬로 실행되며, 이벤트/레코드의 병렬 스트림을 만들어내요.
- SplitEnumerator는 Split을 생성하고 SourceReader에 할당해요. Job Manager에서 단일 인스턴스로 실행되며, 대기 중인 Split의 백로그를 유지하고 균형 있게 리더에 할당하는 역할을 해요.
Source 클래스는 위 세 가지 구성 요소를 하나로 묶는 API 진입점이에요.
스트리밍과 배치의 통합
Data Source API는 무한(unbounded) 스트리밍 소스와 유한(bounded) 배치 소스를 통합된 방식으로 지원해요.
두 경우의 차이는 아주 작아요. 유한/배치의 경우 열거자(enumerator)가 고정된 Split 집합을 생성하고, 각 Split은 반드시 유한해요. 무한 스트리밍의 경우에는 둘 중 하나가 성립하지 않아요 (Split이 유한하지 않거나, 열거자가 계속 새 Split을 생성하거나).
복구 시 Split 재할당 (Split Reassignment On Recovery)
정상적인 상황에서는 SplitEnumerator가 Splits를 SourceReaders에 할당하면, 이 splits가 다른 리더에 다시 할당되지 않아요. 소스가 실패에서 복구될 때는 저장된 상태의 splits가 즉시 리더에 다시 추가돼요.
소스가 SupportsSplitReassignmentOnRecovery 인터페이스를 구현하면 복구 과정이 다르게 동작해요.
- 복구 시
splits를 같은SourceReaders에 즉시 재할당하는 대신, 모든splits를 모아서SplitEnumerator에 다시 추가해요. - 그러면
SplitEnumerator가 이splits를 사용 가능한SourceReaders에 균형 있게 재분배하는 책임을 져요. - 이 메커니즘은 중앙의
SplitEnumerator가 split 분배에 대해 정보에 기반한 결정을 내릴 수 있게 해, 더 유연하고 효율적인 복구를 가능하게 해요.
예시
다음은 스트리밍과 배치 케이스에서 데이터 소스 구성 요소가 어떻게 상호작용하는지 보여주는, 단순화된 개념 예시예요. Kafka와 File 소스 구현이 실제로 이렇게 동작한다는 정확한 설명은 아니며, 설명을 위해 일부가 단순화되었어요.
유한 파일 소스 (Bounded File Source)
소스는 읽을 디렉터리의 URI/Path와, 파일을 파싱하는 방법을 정의하는 Format을 가져요.
- Split은 파일 또는 파일의 일부 영역이에요 (데이터 포맷이 파일 분할을 지원하는 경우).
- SplitEnumerator는 주어진 디렉터리 경로 아래의 모든 파일을 나열해요. Split을 요청하는 다음 리더에게 Split을 할당해요. 모든 Split이 할당되면 요청에
NoMoreSplits로 응답해요. - SourceReader는 Split을 요청하고 할당된 Split (파일 또는 파일 영역)을 읽은 뒤 주어진
Format을 사용해 파싱해요. 다른 Split이 아닌NoMoreSplits메시지를 받으면 종료해요.
무한 스트리밍 파일 소스
이 소스는 SplitEnumerator가 NoMoreSplits로 응답하지 않고, 새 파일이 있는지 주기적으로 주어진 URI/Path의 내용을 나열한다는 점을 제외하면 위와 동일하게 동작해요. 새 파일을 찾으면 그에 대한 새 Split을 생성하고 사용 가능한 SourceReader에 할당할 수 있어요.
무한 스트리밍 Kafka 소스
소스는 Kafka Topic(또는 Topic 목록, Topic 정규식)과 레코드를 파싱하기 위한 Deserializer를 가져요.
- Split은 Kafka Topic Partition이에요.
- SplitEnumerator는 브로커에 연결해 구독 중인 topic에 포함된 모든 topic 파티션을 나열해요. 열거자는 이 작업을 반복해 새로 추가된 topic/파티션을 발견할 수도 있어요.
- SourceReader는 KafkaConsumer를 이용해 할당된 Split (Topic Partition)을 읽고, 제공된
Deserializer로 레코드를 역직렬화해요. Split (Topic Partition)에는 끝이 없으므로 리더는 데이터의 끝에 도달하지 않아요.
유한 Kafka 소스
각 Split (Topic Partition)에 정의된 끝 offset이 있다는 점만 빼고 위와 동일해요. SourceReader가 Split의 끝 offset에 도달하면 그 Split을 종료하고, 할당된 모든 Split이 종료되면 SourceReader가 종료돼요.
Data Source API
이 절에서는 FLIP-27에서 도입된 새로운 Source API의 주요 인터페이스와, Source 개발자를 위한 팁을 설명해요.
Source
Source API는 다음 구성 요소를 생성하는 팩토리 스타일 인터페이스예요.
- Split Enumerator
- Source Reader
- Split Serializer
- Enumerator Checkpoint Serializer
이에 더해 Source는 소스의 boundedness 속성을 제공해서, Flink가 Flink job을 실행할 적절한 모드를 선택할 수 있게 해요.
Source 인스턴스는 런타임에 직렬화되어 Flink 클러스터에 업로드되므로, Source 구현은 직렬화 가능해야 해요.
SplitEnumerator
SplitEnumerator는 Source의 "두뇌"로 기대돼요. 일반적인 SplitEnumerator 구현은 다음을 수행해요.
SourceReader등록 처리SourceReader실패 처리:SourceReader가 실패하면addSplitsBack()메서드가 호출돼요. SplitEnumerator는 실패한SourceReader가 확인(acknowledge)하지 않은 split 할당을 되돌려받아야 해요.SourceEvent처리:SourceEvent는SplitEnumerator와SourceReader사이에서 전송되는 사용자 정의 이벤트예요. 구현은 이 메커니즘을 활용해 정교한 조정을 수행할 수 있어요.- Split 발견 및 할당:
SplitEnumerator는 새 split 발견, 새SourceReader등록,SourceReader실패 등 다양한 이벤트에 대응해SourceReader에 split을 할당할 수 있어요.
SplitEnumerator는 Source가 생성 또는 복원 시 제공하는 SplitEnumeratorContext의 도움으로 위 작업을 수행해요. SplitEnumeratorContext는 SplitEnumerator가 리더에 대한 필요한 정보를 검색하고 조정 작업을 수행할 수 있게 해요. Source 구현은 SplitEnumeratorContext를 SplitEnumerator 인스턴스에 전달해야 해요.
SplitEnumerator 구현이 메서드가 호출될 때만 조정 작업을 수행하는 반응형 방식으로 잘 동작할 수 있는 반면, 일부 구현은 능동적으로 작업을 수행하고 싶을 수도 있어요. 예를 들어 주기적으로 split 발견을 실행하고 새 split을 SourceReaders에 할당하고 싶을 수 있어요. 이런 구현은 SplitEnumeratorContext의 callAsync() 메서드가 유용해요. 아래 코드 스니펫은 SplitEnumerator 구현이 자체 스레드를 유지하지 않고도 이를 달성하는 방법을 보여줘요.
class MySplitEnumerator implements SplitEnumerator<MySplit, MyCheckpoint> {
private final long DISCOVER_INTERVAL = 60_000L;
private final SplitEnumeratorContext<MySplit> enumContext ;
/** The Source creates instances of SplitEnumerator and provides the context. */
MySplitEnumerator(SplitEnumeratorContext<MySplit> enumContext) {
this.enumContext = enumContext;
}
/**
* A method to discover the splits.
*/
private List<MySplit> discoverSplits() {...}
@Override
public void start() {
...
enumContext.callAsync(this::discoverSplits, (splits, thrown) -> {
Map<Integer, List<MySplit>> assignments = new HashMap<>();
int parallelism = enumContext.currentParallelism();
for (MySplit split : splits) {
int owner = split.splitId().hashCode() % parallelism;
assignments.computeIfAbsent(owner, s -> new ArrayList<>()).add(split);
}
enumContext.assignSplits(new SplitsAssignment<>(assignments));
}, 0L, DISCOVER_INTERVAL);
...
}
...
}
Python API에서는 아직 지원되지 않아요.
SourceReader
SourceReader는 Task Manager에서 실행되어 Splits의 레코드를 소비하는 구성 요소예요.
SourceReader는 풀(pull) 기반 소비 인터페이스를 노출해요. Flink task는 pollNext(ReaderOutput)를 루프로 계속 호출해 SourceReader에서 레코드를 가져와요. pollNext(ReaderOutput) 메서드의 반환값은 소스 리더의 상태를 나타내요.
MORE_AVAILABLE- SourceReader가 즉시 사용 가능한 더 많은 레코드를 가지고 있어요.NOTHING_AVAILABLE- SourceReader가 현재는 더 많은 레코드를 가지고 있지 않지만, 나중에는 가질 수도 있어요.END_OF_INPUT- SourceReader가 모든 레코드를 소진하고 데이터의 끝에 도달했어요. 즉 SourceReader를 닫을 수 있다는 뜻이에요.
성능을 위해 pollNext(ReaderOutput) 메서드에 ReaderOutput가 제공되므로, SourceReader는 필요한 경우 pollNext()를 한 번 호출해서 여러 레코드를 방출할 수 있어요. 예를 들어 외부 시스템이 블록 단위로 동작하는 경우가 있어요. 블록이 여러 레코드를 포함할 수 있지만 소스는 블록 경계에서만 체크포인트할 수 있어요. 이 경우 SourceReader는 한 블록의 모든 레코드를 ReaderOutput에 한 번에 방출할 수 있어요. 다만 SourceReader 구현은 단일 pollNext() 호출에서 여러 레코드를 방출하는 것을 피해야 해요.
Split Reader API
핵심 SourceReader API는 완전히 비동기적이며, 구현이 split 읽기를 수동으로 비동기적으로 관리해야 해요. 그러나 실제로 대부분의 소스는 클라이언트에 대한 차단 poll() 호출(예: KafkaConsumer)이나 분산 파일 시스템(HDFS, S3 등)에 대한 차단 I/O 같은 차단(blocking) 연산을 수행해요. 이를 비동기 Source API와 호환되게 하려면, 이러한 차단(동기) 연산은 별도의 스레드에서 수행되고, 데이터를 리더의 비동기 부분에 넘겨줘야 해요.
SplitReader는 파일 읽기, Kafka 등과 같은 단순한 동기 읽기/폴링 기반 소스 구현을 위한 고수준 API예요. 핵심은 SourceReaderBase 클래스로, SplitReader를 받아 fetch 스레드를 생성해 SplitReader를 실행하며 다양한 소비 스레딩 모델을 지원해요.
SplitReader
SplitReader API에는 세 가지 메서드만 있어요.
RecordsWithSplitIds를 반환하는 차단 fetch 메서드- split 변경을 처리하는 비차단 메서드
- 차단 fetch 연산을 깨우는 비차단 wake up 메서드
SplitReader는 외부 시스템에서 레코드를 읽는 것에만 집중하므로, SourceReader보다 훨씬 단순해요. 자세한 내용은 클래스의 Java doc을 확인하세요.
SourceReaderBase
SourceReader 구현이 흔히 다음을 수행한다는 점은 아주 보편적이에요.
- 외부 시스템의 split에서 차단 방식으로 가져오는 스레드 풀을 가진다
- 내부 fetch 스레드와
pollNext(ReaderOutput)같은 다른 메서드 호출 사이의 동기화를 처리한다 - 워터마크 정렬을 위해 split별 워터마크를 유지한다
- 체크포인트를 위해 각 split의 상태를 유지한다
- 레코드 방출 속도를 제한한다
새 SourceReader를 작성하는 작업을 줄이기 위해, Flink는 SourceReader의 기본 구현 역할을 하는 SourceReaderBase 클래스를 제공해요. SourceReaderBase는 위 작업을 모두 기본적으로 처리해요. 새 SourceReader를 작성하려면 SourceReader 구현이 SourceReaderBase를 상속하고, 몇 가지 메서드를 채우고, 고수준 SplitReader를 구현하면 돼요.
속도 제한 (Rate Limiting)
사용자는 RateLimiterStrategy를 SourceReaderBase 생성자에 전달해 레코드 방출 속도를 제한할 수 있어요. 기본적으로 SourceReaderBase에는 속도 제한이 적용되지 않아요.
SplitFetcherManager
SourceReaderBase는 함께 동작하는 SplitFetcherManager의 동작에 따라 몇 가지 스레딩 모델을 기본으로 지원해요. SplitFetcherManager는 각각 SplitReader로 fetch하는 SplitFetchers 풀을 생성하고 유지하는 데 도움을 줘요. 또한 각 split fetcher에 split을 어떻게 할당할지 결정해요.
예를 들어 아래에 설명된 대로 SplitFetcherManager는 고정된 수의 스레드를 가질 수 있고, 각 스레드는 SourceReader에 할당된 일부 split에서 fetch해요. 다음 코드 스니펫은 이 스레딩 모델을 구현해요.
/**
* A SplitFetcherManager that has a fixed size of split fetchers and assign splits
* to the split fetchers based on the hash code of split IDs.
*/
public class FixedSizeSplitFetcherManager<E, SplitT extends SourceSplit>
extends SplitFetcherManager<E, SplitT> {
private final int numFetchers;
public FixedSizeSplitFetcherManager(
int numFetchers,
Supplier<SplitReader<E, SplitT>> splitReaderSupplier,
Configuration config) {
super(splitReaderSupplier, config);
this.numFetchers = numFetchers;
// Create numFetchers split fetchers.
for (int i = 0; i < numFetchers; i++) {
startFetcher(createSplitFetcher());
}
}
@Override
public void addSplits(List<SplitT> splitsToAdd) {
// Group splits by their owner fetchers.
Map<Integer, List<SplitT>> splitsByFetcherIndex = new HashMap<>();
splitsToAdd.forEach(split -> {
int ownerFetcherIndex = split.hashCode() % numFetchers;
splitsByFetcherIndex
.computeIfAbsent(ownerFetcherIndex, s -> new ArrayList<>())
.add(split);
});
// Assign the splits to their owner fetcher.
splitsByFetcherIndex.forEach((fetcherIndex, splitsForFetcher) -> {
fetchers.get(fetcherIndex).addSplits(splitsForFetcher);
});
}
}
Python API에서는 아직 지원되지 않아요.
이 스레딩 모델을 사용하는 SourceReader는 다음과 같이 만들 수 있어요.
public class FixedFetcherSizeSourceReader<E, T, SplitT extends SourceSplit, SplitStateT>
extends SourceReaderBase<E, T, SplitT, SplitStateT> {
public FixedFetcherSizeSourceReader(
Supplier<SplitReader<E, SplitT>> splitFetcherSupplier,
RecordEmitter<E, T, SplitStateT> recordEmitter,
Configuration config,
SourceReaderContext context) {
super(
new FixedSizeSplitFetcherManager<>(
config.get(SourceConfig.NUM_FETCHERS),
splitFetcherSupplier,
config),
recordEmitter,
config,
context);
}
@Override
protected void onSplitFinished(Map<String, SplitStateT> finishedSplitIds) {
// Do something in the callback for the finished splits.
}
@Override
protected SplitStateT initializedState(SplitT split) {
...
}
@Override
protected SplitT toSplitType(String splitId, SplitStateT splitState) {
...
}
}
Python API에서는 아직 지원되지 않아요.
SourceReader 구현은 SplitFetcherManager와 SourceReaderBase 위에 자신만의 스레딩 모델을 쉽게 구현할 수도 있어요.
이벤트 시간과 워터마크 (Event Time and Watermarks)
이벤트 시간(Event Time) 할당과 워터마크 생성은 데이터 소스의 일부로 이루어져요. Source Reader를 떠나는 이벤트 스트림은 이벤트 타임스탬프를 가지며 (스트리밍 실행 중에는) 워터마크를 포함해요. Event Time과 Watermark에 대한 소개는 Timely Stream Processing을 참고하세요.
API
WatermarkStrategy는 DataStream API에서 Source 생성 시 전달되며, TimestampAssigner와 WatermarkGenerator를 모두 생성해요.
environment.fromSource(
Source<OUT, ?, ?> source,
WatermarkStrategy<OUT> timestampsAndWatermarks,
String sourceName);
environment.from_source(
source: Source,
watermark_strategy: WatermarkStrategy,
source_name: str,
type_info: TypeInformation = None)
TimestampAssigner와 WatermarkGenerator는 ReaderOutput(또는 SourceOutput)의 일부로 투명하게 실행되므로, 소스 구현자는 타임스탬프 추출 및 워터마크 생성 코드를 구현할 필요가 없어요.
이벤트 타임스탬프 (Event Timestamps)
이벤트 타임스탬프는 두 단계로 할당돼요.
SourceReader는SourceOutput.collect(event, timestamp)를 호출해 소스 레코드 타임스탬프를 이벤트에 붙일 수 있어요. 이는 Kafka, Kinesis, Pulsar, Pravega처럼 레코드 기반이고 타임스탬프가 있는 데이터 소스에만 관련돼요. 레코드에 타임스탬프가 없는 소스(파일 등)는 소스 레코드 타임스탬프가 없어요. 이 단계는 소스 커넥터 구현의 일부이며, 소스를 사용하는 애플리케이션이 매개변수화하지 않아요.- 애플리케이션이 구성한
TimestampAssigner가 최종 타임스탬프를 할당해요.TimestampAssigner는 원본 소스 레코드 타임스탬프와 이벤트를 봐요. 할당자는 소스 레코드 타임스탬프를 사용하거나 이벤트의 필드에 접근해 최종 이벤트 타임스탬프를 얻을 수 있어요.
이 두 단계 접근 방식은 사용자가 소스 시스템의 타임스탬프와 이벤트 데이터의 타임스탬프를 모두 이벤트 타임스탬프로 참조할 수 있게 해요.
참고: 소스 레코드 타임스탬프가 없는 데이터 소스(파일 등)를 사용하면서 소스 레코드 타임스탬프를 최종 이벤트 타임스탬프로 선택하면, 이벤트는 LONG_MIN(=-9,223,372,036,854,775,808)과 같은 기본 타임스탬프를 갖게 돼요.
워터마크 생성 (Watermark Generation)
워터마크 생성기는 스트리밍 실행 중에만 활성화돼요. 배치 실행은 워터마크 생성기를 비활성화하고, 아래에 설명된 모든 관련 연산은 사실상 no-op이 돼요.
데이터 소스 API는 split별로 워터마크 생성기를 개별적으로 실행하는 것을 지원해요. 이를 통해 Flink가 split별로 이벤트 시간 진행을 관찰할 수 있는데, 이는 이벤트 시간 스큐(skew)를 제대로 처리하고 유휴(idle) 파티션이 전체 애플리케이션의 이벤트 시간 진행을 지연시키는 것을 방지하는 데 중요해요.
Split Reader API를 사용해 소스 커넥터를 구현하면 이 처리가 자동으로 이루어져요. Split Reader API 기반의 모든 구현은 기본으로 split 인식(split-aware) 워터마크를 가져요.
Split 인식 워터마크 생성을 사용하기 위해 저수준 SourceReader API를 구현하려면, 구현은 서로 다른 split의 이벤트를 서로 다른 출력, 즉 Split-로컬 SourceOutputs로 출력해야 해요. Split-로컬 출력은 메인 ReaderOutput에서 createOutputForSplit(splitId)와 releaseOutputForSplit(splitId) 메서드를 통해 생성 및 해제할 수 있어요. 자세한 내용은 클래스와 메서드의 JavaDocs를 참고하세요.
Split 수준 워터마크 정렬 (Split Level Watermark Alignment)
소스 연산자 워터마크 정렬은 Flink 런타임이 처리하지만, split 수준 워터마크 정렬을 달성하려면 소스가 추가로 SourceReader#pauseOrResumeSplits와 SplitReader#pauseOrResumeSplits를 구현해야 해요. Split 수준 워터마크 정렬은 소스 리더에 여러 split이 할당된 경우에 유용해요. 기본적으로 이 구현들은 소스 리더에 둘 이상의 split이 할당되고, split이 WatermarkStrategy가 구성한 워터마크 정렬 임계값을 초과하면 UnsupportedOperationException을 던져요 (pipeline.watermark-alignment.allow-unaligned-source-splits가 false로 설정된 경우). SourceReaderBase는 SourceReader#pauseOrResumeSplits에 대한 구현을 포함하므로, 상속하는 소스는 SplitReader#pauseOrResumeSplits만 구현하면 돼요. 더 많은 구현 힌트는 javadocs를 참고하세요.