원천 OLTP 데이터베이스의 변경을 잡아 Kafka 로 보내고, 스트리밍 처리기가 Iceberg 테이블에 반영한다. 분석 엔진은 Iceberg 테이블을 읽는다.
Source DB (MySQL / Oracle / PostgreSQL)
│ binlog · redo · WAL
▼
CDC Connector (Debezium 등)
│ 변경 이벤트
▼
Kafka topic (key = PK, value = before/after)
│
▼
Sink (Kafka Connect Iceberg sink / Spark Structured Streaming / Flink)
│ append 또는 MERGE INTO
▼
Iceberg table (스냅샷 · 파티션 · 타임 트래블)
│
▼
Query engine (Trino · Impala · Spark SQL)
CDC 수집. 초기 적재(snapshot)와 이후 증분(incremental)을 어떻게 이어 붙일지 정한다. Debezium 은 스냅샷 모드를 여러 가지로 제공하며, 운영 중인 DB 에 부하를 주지 않도록 시간대와 병렬도를 조절한다. 원천 DB 쪽에는 binlog · 보충 로깅 · logical replication 설정이 선행돼야 한다.
Kafka 토픽 설계. 키를 기본키로 두면 같은 행의 변경이 같은 파티션으로 들어가 순서가 보장된다. 순서가 깨지면 갱신과 삭제가 뒤집히므로 중요하다. 이벤트 구조는 대체로 다음과 같다.
{
"op": "u",
"before": { "id": 1, "name": "old" },
"after": { "id": 1, "name": "new" },
"ts_ms": 1700000000000
}
보존 기간은 재처리 정책과 함께 정한다. 싱크가 며칠 멈춰도 복구할 수 있어야 한다면 그만큼 길게 둔다. 스키마 변경을 다루려면 스키마 레지스트리를 함께 쓴다.
싱크 선택. 단순 append 라면 Kafka Connect 의 Iceberg 싱크로 충분하다. 갱신과 삭제를 반영해 최신 상태를 유지해야 한다면 Spark Structured Streaming 이나 Flink 로 MERGE INTO 를 수행하는 구성이 필요하다. 마이크로배치 주기를 짧게 잡으면 지연은 줄지만 작은 파일이 많아진다.
Iceberg 테이블. 갱신과 삭제가 있으면 포맷 버전 2 로 만든다. 파티션은 조회 패턴에 맞춰 정하되, 잘게 쪼개면 파일 수가 폭증한다. 스냅샷과 delete 파일이 쌓이므로 compaction 과 스냅샷 만료를 처음부터 배치로 걸어 둔다. 이것을 빠뜨린 파이프라인은 몇 주 안에 조회가 느려진다.
작은 파일 문제가 가장 흔하다. 스트리밍 적재는 주기마다 파일을 만들므로 별도 정리가 없으면 수십만 개가 쌓인다.
ALTER TABLE iceberg.db.events EXECUTE optimize;
ALTER TABLE iceberg.db.events EXECUTE expire_snapshots(retention_threshold => '7d');
ALTER TABLE iceberg.db.events EXECUTE remove_orphan_files(retention_threshold => '7d');
중복 처리도 짚어야 한다. 장애 복구 과정에서 같은 이벤트가 다시 들어올 수 있으므로, 기본키 기준 MERGE 로 멱등하게 만들거나 이벤트 ID 로 중복을 걸러낸다.
삭제 이벤트(op=d)를 어떻게 다룰지도 정한다. 물리적으로 지울지, 삭제 플래그만 남길지에 따라 하류 분석 결과가 달라진다.
스키마 변경은 원천 DDL 이 바뀌는 순간 파이프라인 전체에 영향을 준다. 컬럼 추가는 대체로 안전하지만 타입 변경과 컬럼 삭제는 사전 합의가 필요하다.