Apache Paimon: 스트리밍 Lakehouse를 위한 LSM 기반 테이블 포맷과 CDC 운영 경계
요약
Apache Paimon은 Flink와 Spark에서 스트리밍과 배치 처리를 동시에 지원하는 Lakehouse 테이블 포맷이다. Iceberg나 Delta Lake가 주로 배치 커밋을 기준으로 설계된 것과 달리, Paimon은 LSM(Log-Structured Merge-Tree) 구조를 오브젝트 스토리지에 직접 적용해 스트리밍 업서트(upsert)를 테이블의 기본 동작으로 만들었다. CDC(Change Data Capture) 파이프라인에서 MySQL → Kafka → Paimon 경로가 특히 실용적이며, 테이블이 동시에 결과 테이블과 변경 스트림 소스 역할을 할 수 있다. 아파치 소프트웨어 재단 최상위 프로젝트로 2024년 졸업했고, 현재 1.x 시리즈가 진행 중이다.
1. 왜 Paimon이 필요한가: 기존 포맷의 한계
Iceberg와 Delta Lake는 스냅샷 기반 커밋 모델을 사용한다. 스트리밍 환경에서 Flink가 1초마다 체크포인트를 생성하면, 매 체크포인트마다 새 스냅샷과 메타데이터 파일이 생성된다. 24시간 후에는 수만 개의 작은 파일이 생기고, 컴팩션과 스냅샷 만료 작업을 별도로 관리해야 한다.
더 큰 문제는 업데이트 지원이다. CDC에서 MySQL의 UPDATE 이벤트를 Iceberg로 처리하려면 Merge-On-Read나 파티션 재작성이 필요하고, 신선한(fresh) 데이터를 읽으려면 복잡한 DELETE + INSERT 패턴이 필요하다.
Paimon의 접근 방식: LSM을 Lake에 직접 구현한다.
2. LSM 기반 아키텍처
Paimon의 핵심 구조는 기존 데이터베이스의 LSM Tree를 오브젝트 스토리지(S3, HDFS, GCS)에 구현한 것이다.
Binlog 이벤트
Paimon Sink
S3 / HDFS
변경 스트림 읽기
스냅샷 읽기
실시간 차원 조인
버킷과 LSM Tree
테이블은 파티션으로 나뉘고, 각 파티션은 버킷(bucket) 으로 세분된다. 버킷 수는 테이블 생성 시 고정(bucket = N)하거나 동적 버킷 모드를 사용한다. 각 버킷이 하나의 독립적인 LSM Tree를 가진다.
Flink가 체크포인트를 완료할 때 Writer는 L0 파일을 오브젝트 스토리지에 플러시한다. 이후 백그라운드 컴팩션이 L0 → L1 → L2로 파일을 병합하며 정렬 상태를 유지한다.
3. 두 가지 테이블 유형
Primary Key Table (업데이트 지원)
CREATE TABLE orders (
order_id BIGINT,
customer_id INT,
amount DECIMAL(10,2),
status STRING,
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'bucket' = '8',
'changelog-producer' = 'lookup'
);- 기본 키가 있는 레코드의 업서트(UPSERT)와 삭제를 지원
- LSM이 기본 키를 기준으로 정렬 → 동일 키의 최신 값을 효율적으로 조회
- CDC 파이프라인에서 MySQL UPDATE/DELETE를 네이티브로 처리
Append-Only Table (대용량 스트리밍 수집)
CREATE TABLE logs (
event_time TIMESTAMP,
user_id BIGINT,
action STRING
) WITH (
'bucket' = '-1' -- 동적 버킷
);- INSERT만 허용, 업데이트 없음
- 로그 데이터나 불변 이벤트 수집에 최적
- 배치 크기에 따라 자동으로 버킷 수 조정(동적 버킷 모드)
4. Changelog Producer: 변경 스트림의 원천
Paimon의 핵심 기능 중 하나는 테이블이 변경 스트림(changelog)을 동시에 방출한다는 점이다. 이를 통해 테이블을 데이터 소스이자 변경 스트림 소스로 사용할 수 있다.
세 가지 changelog producer 모드:
| 모드 | 작동 방식 | 적합한 상황 |
|---|---|---|
none (기본값) | 변경 로그 없음. 배치 읽기만 | 배치 ETL, 스냅샷 쿼리 |
input | 수집된 +I/-D/+U/-U 이벤트 그대로 전달 | 입력이 이미 CDC 형식일 때 |
lookup | 쓰기 시 기존 값을 조회해 before/after 생성 | 원본 데이터가 전체 행 업데이트일 때 |
full-compaction | 컴팩션 결과를 비교해 변경사항 추출 | 배치 쓰기 후 정확한 diff 필요 |
lookup 모드는 L0 인덱스를 메모리에 유지하고 LSM에서 이전 값을 조회하기 때문에 지연 시간이 약간 올라가지만, CDC 파이프라인에서 가장 흔하게 쓰인다.
5. Merge Engine: 업데이트 충돌 해결 전략
같은 기본 키로 여러 레코드가 들어올 때 어떻게 병합할지를 지정한다.
CREATE TABLE metrics (
metric_name STRING,
ts BIGINT,
value DOUBLE,
count BIGINT,
PRIMARY KEY (metric_name) NOT ENFORCED
) WITH (
'merge-engine' = 'aggregation',
'fields.value.aggregate-function' = 'sum',
'fields.count.aggregate-function' = 'sum'
);| Merge Engine | 동작 | 사용 예 |
|---|---|---|
deduplicate (기본) | 같은 PK의 마지막 레코드만 유지 | CDC 상태 테이블, 최신 값 |
partial-update | 지정한 컬럼만 업데이트, 나머지 유지 | 여러 소스의 부분 업데이트 병합 |
aggregation | 지정 함수로 값 집계 | 실시간 집계, 카운터, sum |
first-row | 처음 들어온 레코드만 유지, 이후 무시 | 이벤트 최초 발생 시점 기록 |
6. 세 가지 쓰기 모드: MOR, COW, MOW
LSM 기반이므로 데이터 신선도와 읽기 비용의 트레이드오프를 운영자가 선택한다.
MOR (Merge-On-Read) — 기본값
L0 파일이 쌓이고 컴팩션은 백그라운드에서만 발생. 쓰기 지연이 가장 낮지만, 읽기 시 여러 레이어를 병합해야 한다.
changelog-producer = lookup
# full-compaction은 필요 없음COW (Copy-On-Write)
매 쓰기마다 전체 컴팩션 수행. 읽기가 빠르지만 쓰기 지연이 높다.
'full-compaction.delta-commits' = '1'MOW (Merge-On-Write) + Deletion Vector
쓰기 시 삭제 벡터(bitset)를 생성해 이전 레코드를 논리적으로 표시. 읽기 시 삭제 벡터를 적용해 최신 상태를 반환. ClickHouse나 Delta Lake의 deletion vector와 유사한 방식.
'deletion-vectors.enabled' = 'true'7. 실전 CDC 파이프라인: MySQL → Paimon
가장 일반적인 사용 패턴은 MySQL CDC → Flink → Paimon이다.
UPDATE/DELETE)
+I / -D / +U / -U 이벤트
changelog-producer=input
+ 변경 로그
스트리밍 집계
스냅샷 읽기
임시 쿼리
체크포인트 주기(예: 30초)마다 Flink가 L0 파일을 플러시하므로, 최대 수십 초의 지연으로 MySQL 변경 사항이 Paimon에 반영된다.
8. 시간 여행과 스냅샷
Paimon은 Iceberg처럼 스냅샷 기반 time travel을 지원한다.
-- 특정 시점 스냅샷 읽기
SELECT * FROM orders /*+ OPTIONS('scan.timestamp-millis'='1714000000000') */;
-- 특정 스냅샷 ID 읽기
SELECT * FROM orders /*+ OPTIONS('scan.snapshot-id'='5') */;스냅샷은 Flink 체크포인트마다 생성되며, 보존 정책은 snapshot.num-retained.min/max와 snapshot.time-retained로 설정한다.
9. 운영 고려사항
컴팩션 모니터링
MOR 모드에서 L0 파일이 과도하게 쌓이면 읽기 성능이 저하된다. num-levels-to-file-counts를 모니터링하고, 컴팩션 지연이 발생하면 전용 컴팩션 작업을 별도로 실행한다.
-- 전용 컴팩션 작업 (Flink SQL)
CALL sys.compact(`table` => 'my_db.orders', where => 'dt=20260901');버킷 설계
- 버킷 키: 조인이나 집계에서 자주 사용하는 컬럼. 같은 버킷 내 레코드는 동일 노드에서 처리.
- 버킷 수: 너무 적으면 핫 파티션, 너무 많으면 메타데이터 오버헤드. 일반적으로 버킷당 1-2 GB 기준.
- 동적 버킷: 버킷 수를 사전에 결정하기 어려울 때 사용. 단, 정렬 특성이 보장되지 않아 Lookup Join 효율이 낮다.
Paimon vs 다른 Lakehouse 포맷
| 항목 | Apache Paimon | Apache Iceberg | Delta Lake |
|---|---|---|---|
| 스트리밍 업서트 | 네이티브 (LSM) | 제한적 (Merge-On-Read) | 제한적 |
| CDC 통합 | Flink CDC 네이티브 | 부가 작업 필요 | 부가 작업 필요 |
| 변경 스트림 방출 | 내장 (changelog-producer) | 별도 구성 필요 | DML change feed |
| 배치 쿼리 엔진 | Spark, Flink, Trino | Spark, Trino, Flink | Spark, Trino |
| 성숙도 | ASF TLP (2024 졸업) | ASF TLP (2018 졸업) | Linux Foundation |
10. Open Questions
- 대규모 상용 배포 사례가 Alibaba 에코시스템 외에 얼마나 있는지 확인이 필요하다.
- Paimon 1.x와 Flink 버전 호환성: 최신 Flink 2.x와의 연동은 공식 문서를 통해 확인 권장.
- Iceberg REST Catalog와의 통합 수준은 현재 개발 중이며, 1.x 시리즈에서 구체화될 예정이다.
References
- Apache Paimon 공식 문서 1.1
- Apache Paimon: the Streaming Lakehouse — Ververica
- Apache Paimon GitHub — apache/paimon
- Lakehouse Table Formats in 2026: Iceberg, Delta Lake, Hudi, Paimon, and DuckLake
- Apache Software Foundation Announces Apache Paimon as Top-Level Project
- Apache Paimon: Lakehouse Storage For Real-Time Analytics — Open Source For You (Aug 2026)
- Primary Key Table — Apache Paimon Docs