데이터 소스
데이터 소스 (Data Sources)
이 페이지는 Flink의 Data Source API와 그 배후의 개념 및 아키텍처를 설명합니다. Flink에서 데이터 소스가 어떻게 동작하는지 관심이 있거나, 새 Data Source를 구현하려면 이것을 읽으세요. 미리 정의된 소스 커넥터를 찾고 있다면 Connector Docs를 확인하세요.
출처: 문서
본문
이 페이지는 Flink의 Data Source API와 그 배후의 개념과 아키텍처를 설명합니다. 미리 정의된 소스 커넥터를 찾고 있다면 Connector Docs를 확인하세요.
Data Source 개념 (Data Source Concepts)
핵심 컴포넌트
Data Source에는 Splits, SplitEnumerator, SourceReader 세 가지 핵심 컴포넌트가 있습니다.
- Split은 소스가 소비하는 데이터의 일부입니다. 예를 들어 파일이나 로그 파티션입니다. Splits는 소스가 작업을 분배하고 데이터 읽기를 병렬화하는 단위입니다.
- SourceReader는 Splits를 요청하고 처리합니다. 예를 들어 Split이 나타내는 파일이나 로그 파티션을 읽습니다. SourceReader는 Task Managers에서
SourceOperators안에 병렬로 실행되며 병렬 이벤트/레코드 스트림을 생성합니다. - SplitEnumerator는 Splits를 생성하고 SourceReaders에게 할당합니다. Job Manager에서 단일 인스턴스로 실행되며, 대기 중인 Splits의 백로그를 유지하고 균형 있게 리더에게 할당하는 책임이 있습니다.
Source 클래스는 위의 세 컴포넌트를 함께 묶는 API 진입점입니다.
SplitEnumerator와 SourceReader 상호작용 예시
스트리밍과 배치에 걸친 통합 — Data Source API는 유한(bounded) 배치 소스와 무한(unbounded) 스트리밍 소스를 통합된 방식으로 지원합니다.
두 경우의 차이는 최소입니다. 유한/배치 경우에는 enumerator가 고정된 split 집합을 생성하며, 각 split은 반드시 유한합니다. 무한 스트리밍 경우에는 이 중 하나가 성립하지 않습니다(split이 유한하지 않거나, enumerator가 계속 새 split을 생성하거나).
복구 시 Split 재할당 (Split Reassignment On Recovery) — 정상적인 상황에서는 SplitEnumerator가 SourceReaders에게 Splits를 할당하면 이 splits가 다른 리더에게 다시 할당되지 않습니다. 소스가 실패에서 복구될 때, 저장된 상태의 splits는 리더들에게 즉시 다시 추가됩니다.
소스가 SupportsSplitReassignmentOnRecovery 인터페이스를 구현하면 복구 과정이 다르게 동작합니다. 복구 시 splits를 같은 SourceReaders에게 즉시 재할당하는 대신, 모든 splits를 수집해 SplitEnumerator에게 다시 추가합니다. 그러면 SplitEnumerator가 사용 가능한 SourceReaders 사이에 이 splits를 균형 있게 재분배하는 책임을 집니다. 이 메커니즘은 중앙의 SplitEnumerator가 split 분배에 대해 정보에 기반한 결정을 내릴 수 있게 하여 더 유연하고 효율적인 복구를 가능하게 합니다.
예시 (Examples)
다음은 스트리밍과 배치 경우에서 데이터 소스 컴포넌트가 어떻게 상호작용하는지 설명하는 몇 가지 단순화된 개념적 예시입니다.
이것은 Kafka와 File 소스 구현이 실제로 어떻게 동작하는지 정확히 설명하지 않으며, 일부는 설명을 위해 단순화되었습니다.
유한 파일 소스 (Bounded File Source) — 소스는 읽을 디렉터리의 URI/Path와 파일을 파싱하는 방법을 정의하는 Format을 가집니다.
- Split은 파일 또는 파일의 영역(데이터 포맷이 파일 분할을 지원하는 경우)입니다.
- SplitEnumerator는 주어진 디렉터리 경로 아래의 모든 파일을 나열합니다. Split을 요청하는 다음 리더에게 Splits를 할당합니다. 모든 Splits가 할당되면 요청에 NoMoreSplits로 응답합니다.
- SourceReader는 Split을 요청하고 할당된 Split(파일 또는 파일 영역)을 읽고 주어진 Format으로 파싱합니다. 다른 Split이 아닌 NoMoreSplits 메시지를 받으면 종료됩니다.
무한 스트리밍 파일 소스 (Unbounded Streaming File Source) — 이 소스는 위와 같은 방식으로 동작하지만, SplitEnumerator가 NoMoreSplits로 응답하지 않고 주어진 URI/Path 아래의 내용을 주기적으로 나열해 새 파일을 확인한다는 점이 다릅니다. 새 파일을 찾으면 새 Splits를 생성하고 사용 가능한 SourceReaders에게 할당할 수 있습니다.
무한 스트리밍 Kafka 소스 (Unbounded Streaming Kafka Source) — 소스는 Kafka Topic(또는 Topic 목록 또는 Topic 정규식)과 레코드를 파싱하는 Deserializer를 가집니다.
- Split은 Kafka Topic Partition입니다.
- SplitEnumerator는 브로커에 연결해 구독된 토픽에 포함된 모든 토픽 파티션을 나열합니다. enumerator는 선택적으로 이 작업을 반복해 새로 추가된 토픽/파티션을 발견할 수 있습니다.
- SourceReader는 KafkaConsumer로 할당된 Splits(Topic Partitions)를 읽고 제공된 Deserializer로 레코드를 역직렬화합니다. splits(Topic Partitions)에는 끝이 없으므로 리더는 데이터의 끝에 도달하지 않습니다.
유한 Kafka 소스 (Bounded Kafka Source) — 각 Split(Topic Partition)에 정의된 끝 오프셋이 있다는 점을 제외하면 위와 동일합니다. SourceReader가 Split의 끝 오프셋에 도달하면 해당 Split을 종료합니다. 모든 할당된 Splits가 종료되면 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 작업을 실행할 수 있습니다.
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는 SplitEnumerator가 생성되거나 복원될 때 Source에 제공되는 SplitEnumeratorContext의 도움으로 위 작업을 수행할 수 있습니다. SplitEnumeratorContext는 SplitEnumerator가 리더의 필요한 정보를 검색하고 조정 작업을 수행할 수 있게 합니다. Source 구현은 SplitEnumeratorContext를 SplitEnumerator 인스턴스에 전달해야 합니다.
SplitEnumerator 구현은 메서드가 호출될 때만 조정 작업을 수행하는 반응적 방식으로 잘 동작할 수 있지만, 일부 SplitEnumerator 구현은 적극적으로 작업을 수행하고 싶어할 수 있습니다. 예를 들어 SplitEnumerator가 주기적으로 split 발견을 실행하고 새 split을 SourceReaders에 할당하고 싶을 수 있습니다. 이러한 구현은 SplitEnumeratorContext의 callAsync() 메서드가 유용함을 알게 될 것입니다. 아래 코드 스니펫은 SplitEnumerator 구현이 자신의 스레드를 유지하지 않고 이를 달성하는 방법을 보여줍니다.
Java
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
Still not supported in Python API.
SourceReader
SourceReader는 Splits에서 레코드를 소비하기 위해 Task Managers에서 실행되는 컴포넌트입니다.
SourceReader는 풀 기반(pull-based) 소비 인터페이스를 노출합니다. Flink 작업은 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(ReaderOutput) 호출에서 여러 레코드를 내보내는 것을 피해야 합니다. SourceReader에서 폴링하는 작업 스레드는 이벤트 루프에서 작동하며 차단될 수 없기 때문입니다.
SourceReader의 모든 상태는 snapshotState() 호출에서 반환되는 SourceSplit 내부에 유지되어야 합니다. 이렇게 하면 필요할 때 SourceSplit을 다른 SourceReader에 재할당할 수 있습니다.
SourceReader 생성 시 Source에 SourceReaderContext가 제공됩니다. Source가 컨텍스트를 SourceReader 인스턴스에 전달할 것으로 기대됩니다. SourceReader는 SourceReaderContext를 통해 SplitEnumerator에게 SourceEvent를 보낼 수 있습니다. Source의 일반적인 설계 패턴은 SourceReader가 로컬 정보를 전역 뷰를 가진 SplitEnumerator에게 보고해 결정을 내리도록 하는 것입니다.
SourceReader API는 사용자가 split을 수동으로 다루고 레코드를 가져와 넘겨주는 자신의 스레딩 모델을 가질 수 있게 하는 저수준 API입니다. SourceReader 구현을 돕기 위해 Flink는 SourceReader 작성에 필요한 작업량을 크게 줄여주는 SourceReaderBase 클래스를 제공합니다. 커넥터 개발자는 SourceReader를 처음부터 작성하는 대신 SourceReaderBase를 활용할 것을 강력히 권장합니다. 자세한 내용은 Split Reader API 섹션을 확인하세요.
Source 사용 (Use the Source)
Source에서 DataStream을 만들려면 Source를 StreamExecutionEnvironment에 전달해야 합니다. 예:
Java
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
Source mySource = new MySource(...);
DataStream<Integer> stream = env.fromSource(
mySource,
WatermarkStrategy.noWatermarks(),
"MySourceName");
...
Python
env = StreamExecutionEnvironment.get_execution_environment()
my_source = ...
env.from_source(
my_source,
WatermarkStrategy.no_watermarks(),
"my_source_name")
Split Reader API
핵심 SourceReader API는 완전히 비동기이며 구현이 split 읽기를 비동기적으로 수동 관리하도록 요구합니다. 하지만 실제로 대부분의 소스는 클라이언트의 차단 poll() 호출(예: KafkaConsumer)이나 분산 파일시스템(HDFS, S3 등)의 차단 I/O 연산처럼 차단 연산을 수행합니다. 이를 비동기 Source API와 호환되게 만들려면 이러한 차단(동기) 연산이 별도의 스레드에서 발생해야 하며, 데이터를 리더의 비동기 부분에 넘겨주어야 합니다.
SplitReader는 파일 읽기, Kafka 등과 같은 단순한 동기 읽기/폴링 기반 소스 구현을 위한 고수준 API입니다.
핵심은 SplitReader를 받아 SplitReader를 실행하는 fetcher 스레드를 만들고 다양한 소비 스레딩 모델을 지원하는 SourceReaderBase 클래스입니다.
SplitReader
SplitReader API에는 세 가지 메서드만 있습니다.
- RecordsWithSplitIds를 반환하는 차단 fetch 메서드
- split 변경을 처리하는 비차단 메서드
- 차단 fetch 연산을 깨우는 비차단 wake up 메서드
SplitReader는 외부 시스템에서 레코드를 읽는 데만 집중하므로 SourceReader보다 훨씬 간단합니다. 자세한 내용은 클래스의 Java doc을 확인하세요.
SourceReaderBase
SourceReader 구현이 다음을 수행하는 것은 매우 흔합니다.
- 외부 시스템의 split에서 차단 방식으로 가져오는 스레드 풀을 가집니다.
- 내부 fetching 스레드와
pollNext(ReaderOutput)같은 다른 메서드 호출 사이의 동기화를 처리합니다. - 워터마크 정렬을 위해 split별 워터마크를 유지합니다.
- 체크포인트를 위해 각 split의 상태를 유지합니다.
- 레코드 내보내기 속도를 제한합니다.
새 SourceReader를 작성하는 작업을 줄이기 위해 Flink는 SourceReader의 기본 구현 역할을 하는 SourceReaderBase 클래스를 제공합니다. SourceReaderBase는 위의 모든 작업을 기본으로 수행합니다. 새 SourceReader를 작성하려면 SourceReader 구현이 SourceReaderBase를 상속하고 몇 가지 메서드를 채우고 고수준 SplitReader를 구현하기만 하면 됩니다.
속도 제한 (Rate Limiting)
사용자는 RateLimiterStrategy를 SourceReaderBase의 생성자에 전달해 레코드 내보내기 속도를 제한할 수 있습니다. 기본적으로 속도 제한은 SourceReaderBase에 적용되지 않습니다.
SplitFetcherManager
SourceReaderBase는 함께 동작하는 SplitFetcherManager의 동작에 따라 몇 가지 스레딩 모델을 기본으로 지원합니다. SplitFetcherManager는 각각 SplitReader로 fetching하는 SplitFetcher 풀을 만들고 유지하는 데 도움을 줍니다. 또한 각 split fetcher에 split을 어떻게 할당할지도 결정합니다.
예를 들어 아래와 같이 SplitFetcherManager는 각각 SourceReader에 할당된 일부 split에서 fetching하는 고정된 수의 스레드를 가질 수 있습니다.
다음 코드 스니펫은 이 스레딩 모델을 구현합니다.
Java
/**
* 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
Still not supported in Python API.
그리고 이 스레딩 모델을 사용하는 SourceReader는 다음과 같이 만들 수 있습니다.
Java
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
Still not supported in Python API.
SourceReader 구현은 SplitFetcherManager와 SourceReaderBase 위에서 자신의 스레딩 모델을 쉽게 구현할 수도 있습니다.
이벤트 시간과 워터마크 (Event Time and Watermarks)
Event Time 할당과 Watermark Generation은 데이터 소스의 일부로 발생합니다. Source Readers에서 나가는 이벤트 스트림은 이벤트 타임스탬프를 가지며 (스트리밍 실행 동안) 워터마크를 포함합니다. Event Time과 Watermark에 대한 소개는 Timely Stream Processing을 참조하세요.
API
WatermarkStrategy는 DataStream API에서 생성 중에 Source에 전달되며 TimestampAssigner와 WatermarkGenerator를 모두 만듭니다.
Java
environment.fromSource(
Source<OUT, ?, ?> source,
WatermarkStrategy<OUT> timestampsAndWatermarks,
String sourceName);
Python
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)를 호출해 source record timestamp를 이벤트에 첨부할 수 있습니다. 이는 Kafka, Kinesis, Pulsar, Pravega처럼 레코드 기반이고 타임스탬프가 있는 데이터 소스에만 관련됩니다. 타임스탬프가 있는 레코드에 기반하지 않은 소스(파일처럼)는 source record timestamp를 갖지 않습니다. 이 단계는 소스 커넥터 구현의 일부이며 애플리케이션이 파라미터화하지 않습니다. - 애플리케이션이 구성한
TimestampAssigner가 최종 타임스탬프를 할당합니다.TimestampAssigner는 원래 source record timestamp와 이벤트를 봅니다. assigner는 source record timestamp를 사용하거나 이벤트의 필드에 접근해 최종 이벤트 타임스탬프를 얻을 수 있습니다.
이 두 단계 접근 방식은 사용자가 소스 시스템의 타임스탬프와 이벤트 데이터의 타임스탬프를 모두 이벤트 타임스탬프로 참조할 수 있게 합니다.
참고: source record timestamps가 없는 데이터 소스(파일처럼)를 사용하고 source record timestamp를 최종 이벤트 타임스탬프로 선택하면, 이벤트는 LONG_MIN *(=-9,223,372,036,854,775,808)*과 같은 기본 타임스탬프를 얻습니다.
워터마크 생성 (Watermark Generation)
Watermark Generator는 스트리밍 실행 동안에만 활성화됩니다. 배치 실행은 Watermark Generator를 비활성화하며, 아래에 설명된 모든 관련 연산은 효과적으로 no-op이 됩니다.
Data source API는 워터마크 생성기를 split별로 개별 실행하는 것을 지원합니다. 이는 Flink가 split별로 이벤트 시간 진행을 개별적으로 관찰할 수 있게 하며, 이는 event time skew를 적절히 처리하고 idle partitions이 전체 애플리케이션의 이벤트 시간 진행을 붙잡는 것을 방지하는 데 중요합니다.
두 개의 Splits가 있는 Source에서의 Watermark 생성.
Split Reader API를 사용해 소스 커넥터를 구현하면 이는 자동으로 처리됩니다. Split Reader API에 기반한 모든 구현은 기본적으로 split 인지(split-aware) 워터마크를 가집니다.
저수준 SourceReader API 구현이 split 인지 워터마크 생성을 사용하려면 다른 split의 이벤트를 다른 출력, 즉 Split-local SourceOutputs에 출력해야 합니다. Split-local 출력은 메인 ReaderOutput에서 createOutputForSplit(splitId)와 releaseOutputForSplit(splitId) 메서드를 통해 생성되고 해제될 수 있습니다. 자세한 내용은 클래스와 메서드의 JavaDocs를 참조하세요.
Split 레벨 워터마크 정렬 (Split Level Watermark Alignment)
소스 오퍼레이터 워터마크 정렬은 Flink 런타임이 처리하지만, split 레벨 워터마크 정렬을 달성하기 위해 소스는 추가로 SourceReader#pauseOrResumeSplits와 SplitReader#pauseOrResumeSplits를 구현해야 합니다. Split 레벨 워터마크 정렬은 소스 리더에 여러 split이 할당될 때 유용합니다. 기본적으로 이 구현들은 UnsupportedOperationException을 던지며, 하나 이상의 split이 할당되고 split이 WatermarkStrategy로 구성된 워터마크 정렬 임계값을 초과하면 pipeline.watermark-alignment.allow-unaligned-source-splits가 false로 설정됩니다. SourceReaderBase는 SourceReader#pauseOrResumeSplits의 구현을 포함하므로 상속하는 소스는 SplitReader#pauseOrResumeSplits만 구현하면 됩니다. 구현 힌트는 javadocs를 참조하세요.