LLM WikiAccess-protected knowledge portal
← 스터디 홈
163편 · 약 15분

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)에 구현한 것이다.

MySQL CDC
Binlog 이벤트
Flink Job
Paimon Sink
Paimon 테이블
S3 / HDFS
버킷 0 (LSM Tree)
L0 (Writer Buffer)
L1 files
L2 files (sorted)
버킷 N (LSM Tree)
L0 files
L1 files
L2 files (sorted)
스트리밍 소비자
변경 스트림 읽기
배치 소비자
스냅샷 읽기
Flink Lookup Join
실시간 차원 조인
각 버킷은 독립적인 LSM Tree. Flink 체크포인트가 L0 flush를 트리거.
Apache Paimon 아키텍처: LSM Tree on Object Storage

버킷과 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이다.

MySQL
Binlog (INSERT/
UPDATE/DELETE)
Flink CDC Source
flink-cdc-connectors
+I / -D / +U / -U 이벤트
Paimon Sink
PK Table
changelog-producer=input
S3 / HDFS
LSM 파일
+ 변경 로그
Flink 분석 쿼리
스트리밍 집계
Spark 배치 쿼리
스냅샷 읽기
Trino / Presto
임시 쿼리
하나의 Paimon 테이블이 실시간 스트리밍 소스와 배치 OLAP 소스 역할을 동시에 수행
MySQL CDC to Paimon Pipeline

체크포인트 주기(예: 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/maxsnapshot.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 PaimonApache IcebergDelta Lake
스트리밍 업서트네이티브 (LSM)제한적 (Merge-On-Read)제한적
CDC 통합Flink CDC 네이티브부가 작업 필요부가 작업 필요
변경 스트림 방출내장 (changelog-producer)별도 구성 필요DML change feed
배치 쿼리 엔진Spark, Flink, TrinoSpark, Trino, FlinkSpark, 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