Pinterest CDC 파이프라인으로 구축한 Kafka·Flink·Iceberg 레이크하우스
Pinterest의 CDC 수집 아키텍처를 Kafka, Flink, Spark, Iceberg 관점에서 분석하고 운영·마이그레이션 설계를 정리한다.
2026-08-14 · 최초 발행 2026-04-29
전체 덤프가 스트리밍 수집으로 바뀐 이유
기존 수집은 야간 배치에서 전체 테이블 덤프를 읽어 웨어하우스에 넣는 방식이었다. 구현은 단순하지만, 전체 스캔이 끝날 때까지 수십 분~수 시간이 걸리고 다음 실행 전에는 데이터가 stale 상태로 남는다. 테이블 수가 수천 개로 늘면 배치 윈도우도 길어져 SLA를 넘기기 쉽다.
프로덕션 DB에 가해지는 읽기 부하와 삭제 감지 문제도 있다. 소프트 딜리트가 없다면 전체 덤프만으로 삭제된 레코드를 알아내기 어렵다.
CDC(Change Data Capture)는 MySQL의 binlog나 TiDB·TiCDC의 변경 로그를 구독해 Insert, Update, Delete 이벤트를 캡처한다. 변경된 레코드만 다루므로 소스 부하를 줄이고, 삭제 이벤트도 전달할 수 있다.
| 항목 | 배치 수집 | CDC 스트리밍 |
|---|---|---|
| 지연 | 수 시간 ~ 24시간 | 5~15분 |
| 소스 DB 부하 | 높음 (풀 스캔) | 낮음 (binlog 읽기) |
| 삭제 이벤트 | 감지 어려움 | 정확히 캡처 |
| 구현 복잡도 | 낮음 | 높음 |
| 스키마 변경 대응 | 단순 | 자동화 필요 |
| 운영 오버헤드 | 낮음 | 높음 (상태 관리) |
변경 이벤트와 최신 스냅샷을 분리한 구조
Pinterest의 프레임워크는 스트리밍 레이어와 배치 레이어를 함께 둔다. Flink가 CDC 이벤트를 즉시 Iceberg CDC 테이블에 기록하고, Spark가 해당 테이블에서 최신 변경분을 읽어 Base 테이블에 Merge Into 한다. CDC 테이블은 이벤트 레저이고, Base 테이블은 분석 쿼리가 읽는 최신 스냅샷이다.
CDC Service는 MySQL binlog를 구조화된 변경 이벤트와 메타데이터로 바꾼다. TiDB는 TiCDC를 통해 같은 인터페이스로 추상화하며, 이벤트는 Avro 또는 JSON으로 직렬화해 Kafka 토픽에 발행한다. 소스에서 Kafka까지의 지연은 1초 미만이다.
Kafka는 테이블별 전용 토픽과 기본키(PK) 파티션 키를 사용한다. Flink 장애 뒤에도 재처리할 수 있는 내구성 있는 버퍼이며, 토픽 보존 기간을 설정하면 Bootstrap도 다시 실행할 수 있다.
Flink는 Kafka 이벤트를 역직렬화하고 스키마를 검증한 뒤 Iceberg FlinkSink로 CDC Iceberg 테이블에 Append 기록한다. 체크포인트 기반 at-least-once 시맨틱을 보장하고, Iceberg 파티션 전략에 맞춰 레코드를 라우팅한다.
CDC Iceberg 테이블은 모든 변경을 시계열로 저장한다. event_time 기준 시간별 파티션을 사용하며 이벤트는 일반적으로 5분 미만의 지연으로 사용할 수 있다. Spark는 15분~1시간 주기로 최신 변경분을 Base 테이블에 적용한다.
이벤트 레저의 스키마와 파일 관리
CDC 테이블에는 원본 레코드 필드와 함께 변경을 해석하기 위한 메타데이터가 필요하다.
cdc_event_time TIMESTAMP -- 이벤트 발생 시각
op STRING -- 'I' (Insert), 'U' (Update), 'D' (Delete)
pk_hash STRING -- 기본키 해시 (파티셔닝용)
before_image STRUCT<...> -- 변경 전 레코드 (Update/Delete 시)
after_image STRUCT<...> -- 변경 후 레코드 (Insert/Update 시)
Append-Only 테이블이므로 하나의 기본키에 여러 이벤트가 시간순으로 쌓인다. Spark Merge Into는 Iceberg의 ROW_NUMBER 윈도우 함수로 PK별 최신 이벤트를 고른 뒤 Base 테이블에 반영한다.
스트리밍 쓰기는 체크포인트마다 파일을 커밋하므로 소파일이 늘어날 수 있다. Pinterest는 Flink의 FlinkSink에서 write.target-file-size-bytes를 조정하고 체크포인트 간격을 늘려 파일당 데이터량을 키웠다. Spark에서는 Merge Into 전에 대상 CDC 파티션을 별도 compaction 작업으로 병합한다.
대용량 MERGE INTO에서 PK 조인을 일반적으로 수행하면 전체 셔플에 따른 네트워크와 메모리 비용이 커진다. CDC 테이블과 Base 테이블에 동일한 버킷 수와 PK 해시를 적용하면 Spark가 같은 버킷 번호끼리 로컬 조인을 수행해 셔플을 제거할 수 있다. 두 테이블의 파티션 스키마가 다를 때는 별도 워크어라운드가 필요하다.
이력 적재와 변경 스트림을 이어 붙이는 방법
CDC는 시작 이후의 변경만 캡처하므로 최초 가동 때는 Base 테이블을 전체 이력 스냅샷으로 채워야 한다. Pinterest의 Bootstrap은 read replica에서 전체 테이블 덤프를 추출하고, Spark로 Iceberg Base 테이블에 벌크 적재하는 순서로 시작한다.
덤프 시점의 binlog 포지션(GTID)을 기록한 뒤 그 포지션부터 CDC 스트리밍 Flink 잡을 시작한다. 이후 Spark Merge Into가 Base 테이블을 점진적으로 갱신한다. 이 연결이 어긋나면 데이터 불일치가 생기므로, Pinterest는 binlog 포지션을 Kafka 토픽 오프셋과 정확히 매핑했다.
백프레셔와 복구 상태를 운영 지표로 다루기
Kafka 처리 속도가 Iceberg 쓰기 속도를 앞서면 Flink에는 백프레셔가 발생한다. Pinterest는 토픽별 Consumer Group Lag을 Grafana에서 계속 추적하고, 인기 테이블에는 높은 병렬도, 소형 테이블에는 낮은 병렬도를 적용한다. 파일 닫기 트리거는 메모리 크기 기준으로 조정해 체크포인트 빈도와 균형을 맞춘다.
체크포인트 설정은 다음과 같다.
체크포인트 간격: 5분
최소 일시중지: 30초
최대 동시 체크포인트: 1
타임아웃: 10분
재시작 전략: 고정 지연 재시작 (3회, 60초 간격)
체크포인트가 끝난 뒤 Flink가 Iceberg에 커밋한다. 실패하면 마지막 성공 체크포인트부터 Kafka 오프셋을 다시 처리한다. CDC 이벤트가 중복 처리될 가능성은 있지만, Iceberg Merge Into의 Upsert 시맨틱으로 최종 결과의 정확성을 보장한다. 이는 Effectively Exactly-Once에 해당한다.
| 메트릭 | 임계값 | 측정 방법 |
|---|---|---|
| CDC 테이블 지연 | < 5분 | 이벤트 타임스탬프 vs 현재 시각 |
| Base 테이블 지연 | < 15분 | Spark 잡 완료 시각 |
| Kafka Consumer Lag | < 10만 건 | Kafka 내장 메트릭 |
| Flink 체크포인트 성공률 | > 99% | Flink UI |
| Spark Merge Into 소요 시간 | < 10분 | 잡 히스토리 |
기존 배치 테이블에서 읽기 트래픽을 옮길 때
Pinterest는 배치 수집과 CDC 파이프라인을 함께 실행하는 Shadow 운영으로 전환을 시작했다. Base 테이블과 기존 테이블의 레코드 카운트, 샘플 값을 비교하고 불일치가 발견되면 CDC 파이프라인을 디버깅한다.
읽기 트래픽은 비크리티컬 분석 쿼리부터 CDC Base 테이블로 옮긴다. SLA를 만족하지 못하면 즉시 배치 테이블로 폴백한다. CDC 파이프라인의 데이터 품질을 2주 이상 검증한 뒤 배치 잡을 비활성화하고, 모니터링 알럿도 CDC 기준으로 다시 설정한다.
전환 전에 확인할 항목은 다음과 같다.
- Bootstrap 완료 및 binlog 포지션 검증
- 스키마 호환성 확인 (Avro 스키마 레지스트리 등록)
- 소스 DB binlog 보존 기간 설정 (최소 7일 권장)
- Flink 잡 체크포인트 상태 저장소 구성
- 알럿 채널 설정 (Kafka Lag, Flink 체크포인트 실패)
- 롤백 절차 문서화
공개된 운영 규모와 지연
Pinterest는 수천 개의 파이프라인과 페타바이트 규모의 데이터를 이 프레임워크로 운영한다. 데이터 가용성 지연은 24시간 이상에서 15분으로 96% 단축됐고, CDC 테이블 지연은 통상 5분 미만이다.
전체 테이블 스캔을 없애 프로덕션 DB 영향을 최소화했으며, 신규 테이블은 YAML 설정 파일 작성만으로 파이프라인을 구성한다. MySQL, TiDB, KVStore를 통합 지원한다.
스트리밍 레이어까지 이어져야 하는 스키마 진화
Pinterest 엔지니어링 블로그는 다음 과제로 자동 스키마 진화(Schema Evolution)를 제시했다. 소스 테이블에 컬럼 추가·변경·삭제가 발생할 때 Kafka 토픽 스키마, Flink 처리 로직, Iceberg 테이블 스키마까지 안전하게 전파하는 자동화가 목표다.
Iceberg는 컬럼 추가·삭제·리네임을 포함한 스키마 진화를 지원한다. 남은 과제는 Kafka 직렬화 포맷과 Flink 역직렬화가 이 변화를 자동으로 따라가도록 연결하는 일이다.
Sources
- Pinterest Engineering: Next Generation DB Ingestion at Pinterest (Feb 2026)
- InfoQ: Pinterest's CDC-Powered Ingestion Slashes Database Latency from 24 Hours to 15 Minutes
- or1k.net: Next Generation DB Ingestion at Pinterest
- Apache Flink: From Stream to Lakehouse — Kafka Ingestion with the Flink Dynamic Iceberg Sink
- Dremio: A Guide to Change Data Capture (CDC) with Apache Iceberg
- Onehouse: Apache Kafka vs Apache Flink vs Apache Spark — Choosing the Right Ingestion Framework
- Kai Waehner: Data Streaming Meets Lakehouse — Apache Iceberg for Unified Real-Time and Batch Analytics
- Confluent Current 2025: Unified CDC Ingestion and Processing with Apache Flink and Iceberg