AI Briefing

Pinterest, 스트림 기반 DB 인제스션에 파티션 최종화 메커니즘 도입

·2026.09.26 00:01

핵심 내용

Pinterest이 Iceberg 스냅샷 요약을 활용해 다운스트림의 안전한 데이터 소비를 돕는 파티션 최종화 메커니즘을 도입했다.

1 / 15

자세히 보기

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가 자동으로 만들었습니다. 원문의 주장과 맥락은 원문에서 확인해 주세요. 저작권은 원저작자에게 있습니다.

AI 처리 방식을 확인하거나, 요약 오류와 출처 표기 문제, 삭제 요청을 문의 · 건의로 알려주세요.