
Flink와 Iceberg의 메타데이터를 극대화하여 다운스트림 배치 작업의 지연과 데이터 누락 문제를 완벽하게 제어하는 아키텍처를 공개합니다.
본 아티클은 핀터레스트가 배치 기반 데이터 인제스션을 실시간 CDC 스트리밍 방식으로 마이그레이션하면서 겪은 데이터 완결성 판단 문제를 해결하는 과정을 다룹니다. 플링크와 아이스버그의 싱크 단에 직접 확장 기능인 Accumulators와 CommitProcessor를 구현하여 성능 저하 없이 파티션 완료 상태를 추적합니다. 최종적으로 데이터 신선도와 완결성 사이의 조율이 가능한 설정형 정책을 제공하여 다양한 다운스트림 시스템의 요구사항을 충족합니다.
실시간 CDC 스트리밍 환경에서 일관된 일 또는 시간 단위 배치 다운스트림 처리가 필요하거나, 이벤트 지연 데이터로 인해 배치 파이프라인 시작 시점을 결정하는 데 어려움을 겪는 데이터 엔지니어들에게 Flink-Iceberg 메타데이터 활용 패턴으로 적극 추천합니다.
기존 배치 기반 시스템에서는 전체 DB를 한 번에 읽어 파티션을 쉽게 완료 처리할 수 있었으나, 실시간 스트리밍 기반의 차세대 CDC 프레임워크 도입 이후에는 이벤트 시간 지연으로 인해 특정 파티션의 데이터가 언제 완료되는지 판단하기 어려워졌습니다.
아파치 플링크와 아이스버그 싱크에 Accumulators와 CommitProcessor라는 확장 포인트를 도입하여, 체크포인트별로 이벤트 시간 통계(t-digest 스케치)를 수집하고 커밋 후 비퇴행성 워터마크 알고리즘을 수행하여 파티션 완료 여부를 판단하도록 구현했습니다.
수백 바이트 수준의 초경량 메타데이터만으로 데이터를 자가 기술하며, 다운스트림 배치 작업이 센서를 통해 메타데이터를 확인하고 안전한 시점에 처리를 시작할 수 있게 되었으며, 이 신호는 스파크 업서트 작업을 통해 베이스 테이블까지 전파됩니다.
Trade-off
데이터 완결성과 대기 시간 사이의 트레이드오프가 존재하며, 신속성을 위해 퍼센타일 기반 정책을 선택할 경우 매우 늦게 도착하는 일부 데이터가 유실되거나 파티션 완료 마킹 이후에 도착하여 얼럿을 발생시킬 수 있습니다.
대량의 데이터 분포를 고정된 작은 크기의 메모리 공간 내에서 고정밀도로 근사하는 컴팩트하고 병합 가능한 스케치 알고리즘입니다.
아파치 아이스버그의 메타데이터 레이어로, 매 커밋마다 쓰여진 파일 수 등 쓰기 작업에 대한 요약 정보를 자체 보관하는 영역입니다.
시간의 경과에 따라 절대 뒤로 후퇴하지 않고 오직 앞으로만 가도록 보장된 논리적 시간의 마일스톤을 추적하는 알고리즘입니다.




