데이터 싱크
데이터 싱크 (Data Sinks)
이 페이지는 Flink의 Data Sink API와 그 배경이 되는 개념·아키텍처를 설명해요. Flink에서 데이터 싱크가 어떻게 동작하는지, 또는 새로운 Data Sink를 어떻게 구현하는지 알고 싶다면 이 문서를 읽어보세요. 미리 정의된 싱크 커넥터를 찾고 있다면 Connector Docs를 확인하세요.
출처: 문서
본문
Data Sink API
이 절에서는 FLIP-191과 FLIP-372에서 도입된 새로운 Sink API의 주요 인터페이스와, Sink 개발자를 위한 팁을 설명해요.
Sink
Sink API는 데이터를 쓰기 위한 SinkWriter를 생성하는 팩토리 스타일 인터페이스예요. Sink 인스턴스는 런타임에 직렬화되어 Flink 클러스터에 업로드되므로, Sink 구현은 직렬화 가능해야 해요.
Sink 사용하기
DataStream.sinkTo(Sink) 메서드를 호출해 DataStream에 Sink를 추가할 수 있어요. 예를 들어,
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
Source mySource = new MySource(...);
DataStream<Integer> stream = env.fromSource(
mySource,
WatermarkStrategy.noWatermarks(),
"MySourceName");
Sink mySink = new MySink(...);
stream.sinkTo(mySink);
...
SinkWriter
핵심 SinkWriter API는 데이터를 다운스트림 시스템에 쓰는 역할을 담당해요. SinkWriter API에는 세 가지 메서드만 있어요.
write(InputT element, Context context): 요소를 writer에 추가해요.flush(boolean endOfInput): 체크포인트 또는 입력의 끝에서 호출되며, 이 플래그를 설정하면 writer가 at-least-once를 위해 대기 중인 모든 데이터를 flush해요. exactly-once 의미론을 달성하려면 writer가SupportsCommitter인터페이스를 구현해야 해요.writeWatermark(Watermark watermark): watermark를 writer에 추가해요.
자세한 내용은 클래스의 Java doc을 확인하세요.
고급 Sink API (Advanced Sink API)
SupportsWriterState
SupportsWriterState 인터페이스는 싱크가 writer 상태를 지원함을 나타내는 데 사용돼요. 즉 싱크가 실패에서 복구될 수 있음을 뜻해요. SupportsWriterState 인터페이스는 SinkWriter가 StatefulSinkWriter 인터페이스를 구현하도록 요구해요.
SupportsCommitter
SupportsCommitter 인터페이스는 싱크가 2단계 커밋(two-phase commit) 프로토콜을 사용해 exactly-once 의미론을 지원함을 나타내는 데 사용돼요.
Sink는 사전 커밋(precommit)을 수행하는 CommittingSinkWriter와 실제로 데이터를 커밋하는 Committer로 구성돼요. 이 분리를 돕기 위해 CommittingSinkWriter는 체크포인트 또는 입력의 끝에서 committable을 생성하고 이를 Committer로 보내요.
Sink는 직렬화 가능해야 해요. 모든 구성은 즉시 검증되어야 해요. 각 싱크 writer와 committer는 일시적이며 TaskManager의 subtask에서만 생성돼요.
사용자 정의 싱크 토폴로지 (Custom sink topology)
고급 개발자를 위해, committable을 한 subtask로 모아 함께 처리하거나, Committer 이후 작은 파일 병합 같은 연산을 수행하는 등 자신만의 싱크 연산자 토폴로지(일련의 연산자로 구성된 구조)를 지정하고 싶을 수 있어요. Flink는 전문 사용자가 싱크 연산자 토폴로지를 사용자 정의할 수 있도록 다음 인터페이스를 제공해요.
SupportsPreWriteTopology
SupportsPreWriteTopology 인터페이스는 전문 사용자가 SinkWriter 이전에 사용자 정의 연산자 토폴로지를 구현할 수 있게 해요. 이를 이용해 입력 데이터를 처리하거나 재분배할 수 있어요. 예를 들어 같은 파티션의 데이터를 Kafka 또는 Iceberg의 같은 SinkWriter로 보내는 경우예요.
다음 그림은 SupportsPreWriteTopology를 사용한 연산자 토폴로지를 보여줘요. 위 그림에서 사용자는 SupportsPreWriteTopology 토폴로지에 PrePartition과 PostPartition 연산자를 추가하고, 입력 데이터를 SinkWriter로 재분배해요.
SupportsPreCommitTopology
SupportsPreCommitTopology 인터페이스는 전문 사용자가 SinkWriter 이후, Committer 이전에 사용자 정의 연산자 토폴로지를 구현할 수 있게 해요. 이를 이용해 커밋 메시지를 처리하거나 재분배할 수 있어요.
다음 그림은 SupportsPreCommitTopology를 사용한 연산자 토폴로지를 보여줘요. 위 그림에서 사용자는 SupportsPreCommitTopology 토폴로지에 CollectCommit 연산자를 추가하고, SinkWriter의 모든 커밋 메시지를 하나의 subtask로 모은 다음 Committer로 보내 중앙에서 처리하게 해요. 이렇게 하면 서버와의 상호작용 횟수를 줄일 수 있어요.
여기서는 표시 목적으로만 parallelism이 수정되었다는 점에 주의하세요. 실제로 parallelism은 사용자가 설정할 수 있어요.
SupportsPostCommitTopology
SupportsPostCommitTopology 인터페이스는 전문 사용자가 Committer 이후에 사용자 정의 연산자 토폴로지를 구현할 수 있게 해요.
다음 그림은 SupportsPostCommitTopology를 사용한 연산자 토폴로지를 보여줘요. 위 그림에서 사용자는 SupportsPostCommitTopology 토폴로지에 MergeFile 연산자를 추가해요. MergeFile 연산자는 작은 파일을 더 큰 파일로 병합해 파일 시스템 읽기 속도를 높일 수 있어요.