LLM WikiAccess-protected knowledge portal
← 스터디 홈
45편 · 약 14분

Apache Paimon 1.3: 스트리밍과 배치를 LSM 테이블 하나로 통합하는 레이크하우스 포맷

왜 지금 봐야 하나

전통적인 데이터 레이크에서는 스트리밍 수집과 배치 분석이 서로 다른 경로를 탄다. Kafka에서 받은 변경 이벤트를 Parquet 파일로 쓰면, 그 순간부터 파일은 불변이다. 업데이트나 삭제가 필요하면 전체 파티션을 재작성하거나, 별도의 Delta 레코드를 쌓아 두다가 읽기 시점에 병합해야 한다. 이 구조는 소규모 데이터에서는 괜찮지만, 대량 업데이트가 잦은 CDC 파이프라인이나 upsert 워크로드에서는 운영 복잡성이 빠르게 늘어난다.

Apache Paimon은 이 문제를 레이크 포맷 레이어에서 해결한다. LSM(Log-Structured Merge-tree) 구조를 파일 스토리지에 적용해, S3나 HDFS 위에서도 INSERT·UPDATE·DELETE를 네이티브하게 처리한다. 단일 Paimon 테이블이 메시지 큐처럼 스트리밍으로 소비되는 동시에 배치 분석 테이블처럼 쿼리된다. Flink 2.x 시대에 스트리밍 레이크하우스 패턴의 핵심 구현체 중 하나로 자리를 잡았다.

Paimon은 원래 Apache Flink 커뮤니티에서 Flink Table Store라는 이름으로 출발해 2024년 4월 ASF 최상위 프로젝트(TLP)로 졸업했다. 2026년 5월 기준으로 1.3.x 계열이 안정 릴리스이며, 1.4 스냅샷이 개발 중이다.


핵심 개념 한눈에 보기

개념설명운영 의미
Primary Key Table버킷별 LSM 트리, INSERT/UPDATE/DELETE 지원upsert, CDC 수신 테이블에 적합
Append Tableprimary key 없음, 단순 append로그, 이벤트 수집, Kafka 교체
Merge EngineDeduplicate / Partial Update / Aggregation / First Row충돌 처리 전략 선택
Table ModeMOR / COW / MOW (Deletion Vectors)읽기·쓰기 비용 트레이드오프
Changelog Producer다운스트림 스트리밍 소비용 변경 로그 생성Flink 스트리밍 조인·집계에 필수
Bucket파티션 안의 물리적 분산 단위, 각자 독립 LSM병렬 쓰기 제어

Paimon 아키텍처 전체 구조

Apache Paimon 아키텍처: 스트리밍 수집 → LSM 저장 → 다중 엔진 쿼리 데이터 소스 Flink CDC (MySQL/PG) Kafka 스트림 Flink DataStream API Spark 배치 Insert Iceberg/Hive 마이그레이션 초기 적재(batch load) Primary Key Table INSERT·UPDATE·DELETE 지원 / LSM 기반 / 버킷별 독립 트리 Bucket 0 Level 0 (MemTable flush → sorted run) Level 1 (sorted runs, minor compaction) Level N (major compaction, 완전 병합) 파티션 키로 라우팅 → primary key로 정렬 각 sorted run: primary key 범위 겹침 없음 읽기 시 여러 Level 병합 필요 (MOR 기준) Bucket 1 … N-1 Merge Engine Deduplicate: 최신 행 유지 (기본값) Table Mode MOR(기본) / COW / MOW(Deletion Vector) Changelog Producer None / Input / Lookup / Full Compaction Append Table primary key 없음 / 순서 보장 / 로그·이벤트에 최적 수집 경로 쓰기 즉시 새 파일 생성. 정렬 없음. Flink Sink / Spark DataFrameWriter 동일 지원 자동 소파일 병합 `write.target-file-size` 기준 백그라운드 compaction 수동 개입 없이 작은 파일 문제 완화 스트리밍 소비 Flink Source로 append-only 스트림 소비 만료되지 않는 히스토리 데이터 접근 오브젝트 스토리지 / HDFS S3, GCS, Azure Blob, HDFS — Parquet/ORC 데이터 파일 + schema/snapshot 메타데이터 쿼리 엔진 Flink SQL (스트리밍·배치) Spark SQL (배치) Trino / Presto (ad-hoc) StarRocks / Doris Hive (호환 레이어) DuckDB (로컬)
Apache Paimon — 레이크 포맷 위의 LSM 기반 스트리밍·배치 통합 테이블

기본 아키텍처: 왜 LSM 트리인가

Parquet 파일 기반 레이크 포맷(Iceberg, Delta Lake 등)은 COW(Copy On Write) 방식으로 업데이트를 처리한다. 레코드 하나를 바꾸려면 그 레코드가 속한 파일 전체를 재작성해야 한다. 업데이트 빈도가 높은 CDC 워크로드에서는 과도한 쓰기 증폭(write amplification)이 발생한다.

LSM 트리는 이 문제를 다른 방식으로 푼다.

  • 쓰기는 항상 새 파일(sorted run)로 내려간다. 기존 파일을 수정하지 않는다.
  • 병합(compaction)은 백그라운드에서 점진적으로 일어난다. 마이너 compaction은 작은 파일들을 합치고, 메이저 compaction은 여러 Level을 통합한다.
  • 읽기는 여러 Level을 병합하는 비용이 발생한다. 이 트레이드오프를 테이블 모드로 조정한다.

Paimon은 이 LSM 구조를 테이블 파티션 안의 버킷(Bucket) 단위로 독립적으로 운영한다. 각 버킷은 독립된 LSM 트리를 가지므로 병렬 쓰기가 가능하고, primary key는 버킷 내에서만 유일성이 보장된다(파티션 key + primary key 조합이 전역 유일성을 결정한다).


Primary Key Table 심화

버킷 구조와 데이터 배치

레코드가 들어오면 파티션 키로 파티션을 결정하고, 그 안에서 primary key의 해시로 버킷을 결정한다. 버킷 수는 테이블 생성 시 고정되므로, 버킷이 너무 적으면 한 버킷 안에 데이터가 과도하게 쌓이고, 너무 많으면 소파일 문제가 생긴다. Paimon 1.3에서는 버킷 수 재조정을 위한 rescale 작업을 지원한다.

Merge Engine 종류

버킷 안에서 같은 primary key를 가진 여러 레코드가 존재할 때, Merge Engine이 어떻게 처리할지 결정한다.

Merge Engine동작적합한 케이스
deduplicate같은 key 중 가장 최신 행만 유지기본값. upsert, CDC last-write-wins
partial-update각 컬럼을 독립적으로 upsert여러 소스가 각자 컬럼을 채우는 경우
aggregation합계·최댓값 등 집계 함수로 병합실시간 집계 테이블
first-row같은 key의 첫 번째 행만 유지중복 제거, 최초 이벤트 보존

partial-update는 특히 유용하다. 예를 들어 주문 테이블에서 결제 서비스가 payment_status 컬럼만 업데이트하고, 배송 서비스가 delivery_status 컬럼만 업데이트하는 패턴에서 서로의 컬럼을 덮어쓰지 않는다.


Table Mode: 읽기·쓰기 비용 트레이드오프

Paimon primary key table에는 세 가지 모드가 있다.

MOR (Merge On Read, 기본값)

쓰기 시 마이너 compaction만 수행하고, 읽기 시 여러 Level의 sorted run을 실시간으로 병합한다.

  • 쓰기 비용: 낮음 (새 파일 추가, 소규모 compaction)
  • 읽기 비용: 높음 (Level 간 병합 필요)
  • 적합한 케이스: 쓰기 많고 읽기 상대적으로 적은 CDC 수집 파이프라인

COW (Copy On Write)

full-compaction.delta-commits=1 설정 시 모든 쓰기마다 전체 compaction을 동기적으로 수행한다. 읽기는 단일 완전 병합 파일만 본다.

  • 쓰기 비용: 매우 높음 (전체 compaction)
  • 읽기 비용: 낮음 (병합 불필요)
  • 적합한 케이스: 쓰기 드물고 읽기 집중적인 분석 테이블

MOW with Deletion Vectors

deletion-vectors.enabled=true 설정 시 MOR의 읽기 비용을 줄이는 중간 경로다. 쓰기 시 각 파일에 대해 삭제된 레코드의 위치를 RoaringBitmap으로 기록한 Deletion Vector 파일을 생성한다. 읽기 시 "데이터 파일 + deletion vector 필터링"으로 전체 병합 없이 답을 낼 수 있다.

  • 쓰기 비용: MOR보다 약간 높음 (LSM 조회 + deletion vector 파일 생성)
  • 읽기 비용: MOR보다 낮음 (단일 파일 + bitmap 필터)
  • 적합한 케이스: 균형 잡힌 읽기·쓰기, 대화형 쿼리가 섞인 파이프라인

Changelog Producer: 다운스트림 스트리밍을 위한 변경 로그

Paimon primary key table의 독특한 기능 중 하나는 Changelog Producer다. Flink에서 집계·조인 작업은 입력이 단순 append가 아니라 +I(insert), -U(update-before), +U(update-after), -D(delete) 형태의 변경 스트림이어야 한다. Paimon이 이 변경 로그를 직접 생산한다.

Changelog Producer동작비용
none변경 로그 없음최저 (스트리밍 집계 불가)
input입력으로 들어온 변경만 그대로 전달낮음 (CDC 입력 시 유용)
lookup쓰기 전 LSM에서 이전 값을 조회해 before/after 생성중간 (추가 LSM 조회)
full-compaction전체 compaction 결과로 변경 로그 생성높음 (지연 있음, 정확성 최고)

실시간 Flink 집계나 스트리밍 조인 다운스트림이 있다면 lookup 또는 full-compaction을 사용해야 한다. none이나 input은 단순 append만 받는 다운스트림에 적합하다.


Append Table: 로그·이벤트 수집에 최적화

Primary key가 없는 Append Table은 Kafka의 토픽과 비슷하게 동작하되, 쿼리 가능성과 저렴한 스토리지 비용을 같이 제공한다.

  • Flink Source로 스트리밍 소비: 데이터를 메시지 큐처럼 구독한다. 히스토리 데이터가 만료되지 않아 임의 시점부터 재소비할 수 있다.
  • 자동 소파일 병합: 백그라운드 compaction이 write.target-file-size 기준으로 작은 파일들을 자동으로 합친다.
  • 배치 쿼리: Spark나 Trino에서 일반 Parquet 테이블처럼 읽는다.

Kafka를 대체하는 것이 아니라, "스트리밍으로 쓰고 배치로도 읽어야 하는" 워크로드에서 Kafka 토픽의 데이터를 레이크 포맷으로 흘리는 브리지 역할로 자주 쓰인다.


Flink / Spark 통합

Flink

Paimon의 최강 통합 파트너다. Flink Table API에서 직접 CREATE TABLE ... WITH ('connector' = 'paimon') 구문으로 선언하고, INSERT INTO, MERGE INTO, UPDATE, DELETE를 SQL로 실행할 수 있다.

-- Paimon 테이블 생성
CREATE TABLE orders (
  order_id BIGINT,
  user_id BIGINT,
  status STRING,
  updated_at TIMESTAMP(3),
  PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
  'connector' = 'paimon',
  'path' = 's3://bucket/warehouse/orders',
  'changelog-producer' = 'lookup',
  'merge-engine' = 'deduplicate'
);

-- Flink 스트리밍 CDC 수신
INSERT INTO orders SELECT * FROM mysql_cdc_source;

Lookup Join은 Paimon primary key table을 차원 테이블로 사용해 실시간 스트림을 enrichment하는 패턴이다. RocksDB 상태 대신 Paimon 파일을 직접 조회하므로 상태 크기 제한 없이 대규모 차원 테이블을 처리할 수 있다.

Spark

Spark 3.x 이상에서 paimon-spark 라이브러리를 추가하면 Spark SQL과 DataFrame API로 읽기·쓰기가 가능하다. 주로 초기 대량 적재(bulk load), 스케줄된 배치 변환, BI 도구 쿼리에 활용된다.


운영 고려사항

Compaction 전략

Compaction은 Paimon의 핵심 운영 포인트다.

  • 자동 compaction(기본): Flink 스트리밍 작업 안에 compaction이 내장되어 별도 작업 없이 동작한다.
  • 전용 compaction 작업: 분리된 Flink 작업으로 compaction을 실행해 메인 수집 작업과 자원을 분리할 수 있다.
  • Spark compaction: CALL paimon.sys.compact(...) 절차로 Spark에서 compaction을 수동 트리거한다.
-- Spark에서 특정 파티션 수동 compaction
CALL paimon.sys.compact(table => 'orders', partitions => "dt='2026-07-25'");

버킷 수 선택과 rescale

버킷 수가 너무 작으면 핫 버킷이 생기고 병렬성이 제한된다. 너무 많으면 소파일 문제가 심해진다. 일반적인 가이드라인은 버킷당 1~2GB의 데이터를 목표로 한다. Paimon 1.3은 온라인 rescale(서비스 중단 없이 버킷 수 변경)을 지원하기 시작했지만 아직 실험적 기능이다.

메타데이터와 snapshot 관리

Paimon은 각 쓰기마다 스냅샷 파일을 생성한다. snapshot.num-retained.min / snapshot.num-retained.max로 보존할 스냅샷 수를 제어하고, snapshot.time-retained로 최소 보존 시간을 설정한다. 오래된 스냅샷은 자동 만료되어 삭제된다.

한계

  • 트랜잭션 격리 수준: Paimon은 Snapshot Isolation을 제공하지만, 다중 테이블 ACID 트랜잭션은 지원하지 않는다.
  • Iceberg/Delta 비교: 생태계 성숙도 면에서 Iceberg와 Delta Lake에 비해 주변 도구(카탈로그 연동, BI 도구 직접 지원)가 아직 부족하다.
  • LSM compaction 비용: 쓰기 집중 워크로드에서 compaction이 과도한 I/O를 발생시킬 수 있다. 전용 compaction 작업 분리와 자원 모니터링이 필요하다.
  • 스키마 변경: 컬럼 추가는 온라인 가능하지만 primary key 변경은 불가능하다.

운영 체크리스트

  • [ ] 테이블 종류를 결정했는가? (upsert/CDC → Primary Key, 로그/이벤트 → Append)
  • [ ] Primary Key Table의 Merge Engine을 데이터 특성에 맞게 선택했는가?
  • [ ] Table Mode를 선택했는가? (쓰기 중심 → MOR, 읽기 중심 → COW, 균형 → MOW/Deletion Vectors)
  • [ ] 다운스트림 Flink 집계·조인이 있다면 Changelog Producer를 lookup 또는 full-compaction으로 설정했는가?
  • [ ] 버킷 수를 예상 데이터 크기 기준으로 결정했는가? (목표: 버킷당 1~2GB)
  • [ ] Compaction 전략을 자동/전용 작업 분리 중 선택했는가?
  • [ ] 스냅샷 보존 정책(snapshot.num-retained.*, snapshot.time-retained)을 설정했는가?
  • [ ] Time Travel 쿼리가 필요한 경우 충분한 스냅샷 보존 기간을 확보했는가?
  • [ ] Flink Lookup Join 사용 시 차원 테이블 업데이트 빈도와 캐시 TTL을 검토했는가?
  • [ ] 운영 환경 쿼리 엔진(Trino, StarRocks 등)의 Paimon 커넥터 버전 호환성을 확인했는가?

References

  • Apache Paimon, "Primary Key Table Overview", Paimon Docs 1.3. https://paimon.apache.org/docs/1.3/primary-key-table/overview/
  • Apache Paimon, "Table Mode — MOR/COW/MOW", Paimon Docs master. https://paimon.apache.org/docs/master/primary-key-table/table-mode/
  • Apache Paimon, "Deletion Vectors", Paimon Docs. https://paimon.apache.org/docs/0.8/primary-key-table/deletion-vectors/
  • Apache Paimon, "Changelog Producer", Paimon Docs master. https://paimon.apache.org/docs/master/primary-key-table/changelog-producer/
  • Polyzos, G., "Apache Paimon — Table format for Streaming", Dev Genius, Medium. https://blog.devgenius.io/apache-paimon-table-format-for-streaming-1671598dd1f6
  • Polyzos, G., "Apache Paimon: Introducing Deletion Vectors", Medium. https://medium.com/@ipolyzos_/apache-paimon-introducing-deletion-vectors-584666ee90de
  • Ververica, "Apache Paimon: the Streaming Lakehouse". https://www.ververica.com/blog/apache-paimon-the-streaming-lakehouse
  • CelerData, "Apache Paimon Explained: An In-Depth Guide". https://celerdata.com/glossary/apache-paimon
  • Apache Software Foundation, "Apache Software Foundation Announces New Top-Level Project Apache® Paimon", press release 2024-04-16. https://news.apache.org/foundation/entry/apache-software-foundation-announces-new-top-level-project-apache-paimon
  • Merced, A., "Lakehouse Table Formats in 2026: Iceberg, Delta Lake, Hudi, Paimon, and DuckLake", DEV Community. https://dev.to/alexmercedcoder/lakehouse-table-formats-in-2026-iceberg-delta-lake-hudi-paimon-and-ducklake-how-they-work-p1k