Pinterest, 스트림 기반 DB 인제스션에 파티션 최종화 메커니즘 도입
핵심 내용
Pinterest이 Iceberg 스냅샷 요약을 활용해 다운스트림의 안전한 데이터 소비를 돕는 파티션 최종화 메커니즘을 도입했다.
자세히 보기
Pinterest은 Kafka, Flink, Iceberg를 활용해 배치에서 스트림 처리로 전환한 차세대 DB 인제스션 프레임워크에 데이터 완결성 문제를 해결하기 위한 파티션 최종화(Partition Finalization) 메커니즘을 도입했다. 기존 배치 시스템에서는 작업 완료 시 최종화가 암시적으로 처리됐지만, 새로운 스트림 아키텍처에서는 Flink가 변경 데이터 캡처(CDC) 이벤트를 지속적으로 쓰기 때문에 실제 시간(Wall-clock)이 지난 후에도 지연 데이터가 파티션에 유입될 수 있다. 이로 인해 다운스트림 소비자는 파티션이 언제 '최종화'되어 안전하게 읽을 수 있는지 파악하는 데 어려움을 겪었다.
이벤트 시간 기반 워터마킹
해결책은 실제 시간이 아닌 이벤트 시간(Event-time) 기반 접근법을 사용한다. Flink-to-Iceberg 싱크는 각 레코드에서 이벤트 시간을 추출하고 체크포인트 간 최소값, 최대값, 백분위수(예: 99th, 95th) 등의 통계를 수집한다. 이 통계는 t-digest 스케치를 사용해 Iceberg 스냅샷 요약에 저장되며, 레코드 양과 관계없이 고정된 작은 용량(수백 바이트)으로 분포를 근사한다. 이 스케치는 병합 가능하므로 병렬 서브태스크는 커밋 시 통계를 효율적으로 결합할 수 있다.
확장 가능한 싱크를 통한 구현
이 로직은 Flink-to-Iceberg 싱크의 두 가지 범용 확장 지점인 Accumulators와 CommitProcessor를 통해 구현된다.
- Accumulators: 체크포인트 윈도우별로 이벤트 시간 통계를 추적하는 라이터 서브태스크 내부의 경량 집계기다. 스냅샷은 커미터에 의해 병합되어 스냅샷 요약에 기록된다.
- CommitProcessor: 각 Iceberg 커밋 전후에 실행되는 훅이다. 커밋 후 훅은 병합된 통계를 읽어 파티션 최종화 워터마크를 계산하고 Iceberg 테이블 속성을 업데이트한다.
워터마크 알고리즘은 일반적으로 최근 커밋들의 커밋별 최소 이벤트 시간 중 최소값을 가져와 시간 경계로 내림 처리한다. 이 값은 단조 증가하므로 지연 데이터가 도착해도 이미 최종화된 시간은 최종화 상태를 유지하며, 이는 소비자 측에 알림을 발생시킨다.
전파 및 일반화
최종화 신호는 Iceberg 테이블 속성으로 저장되므로 쉽게 전파된다. Spark 업서트 잡은 CDC 테이블에서 워터마크를 읽어 커밋 시 베이스 테이블에 기록하며, 다운스트림 잡은 외부 서비스 없이 표준 메타데이터 센서를 통해 처리를 게이트할 수 있다. 이 메커니즘은 설정 가능하며, 테이블은 보수적 완결성(최소값 사용)과 신선도(백분위수 사용) 중 하나를 선택할 수 있다. 이 아키텍처는 DB 인제스션에 국한되지 않고 Pinterest의 차세대 Stream Ingestion 프레임워크에도 채택되고 있다.
이 한국어 요약은 AI가 자동으로 만들었습니다. 원문의 주장과 맥락은 원문에서 확인해 주세요. 저작권은 원저작자에게 있습니다.