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

Apache Fluss: 실시간 레이크하우스를 위한 스트리밍 스토리지 (ASF TLP 2026.08)

요약

2026년 8월 6일, Apache Fluss가 Apache Software Foundation의 Top-Level Project(TLP)로 졸업했다. Incubator를 거쳐 2년 만에 정식 프로젝트가 된 Fluss는 레이크하우스(Lakehouse) 아키텍처를 위한 스트리밍 스토리지 엔진이다.

기존 레이크하우스는 Kafka 같은 메시지 큐로 스트리밍 데이터를 받아 Iceberg/Delta 같은 테이블 포맷으로 저장하는 2-티어 구조였다. Fluss는 이 둘을 하나로 합친다.

핵심 특성:

  • 컬럼 지향 Arrow 스토리지: 스트리밍 단계에서도 컬럼형으로 저장해 분석 쿼리 속도를 높인다.
  • Log Table + PK Table 이중 추상화: 순서 보장 이벤트 로그와 기본키 기반 업데이트를 하나의 시스템에서 처리.
  • Tiering Service: 스트리밍 데이터를 자동으로 Apache Iceberg / Paimon / Hudi로 오프로딩해 장기 보관과 대규모 배치 분석을 지원.
  • Flink, Spark, Doris 통합: 기존 레이크하우스 도구와 네이티브 커넥터로 연결.
  • 기존 행 지향 스트리밍 스토리지 대비 스트리밍 읽기 속도 10배 향상(컬럼 프루닝 기반).

배경: 레이크하우스의 스트리밍 문제

Kafka + Iceberg 2-티어 구조의 한계

현재 가장 널리 쓰이는 실시간 레이크하우스 아키텍처는 다음과 같다.

Producer → Kafka → (Flink/Spark 스트리밍 잡) → Iceberg/Delta → (쿼리 엔진)

이 구조에는 여러 운영 비용이 따른다.

  • 이중 저장: 같은 데이터가 Kafka와 Iceberg 두 곳에 존재한다.
  • 지연 증가: Kafka에서 Iceberg로 커밋하는 주기가 분 단위여서 최신 데이터 접근에 지연이 생긴다.
  • 메시지 큐의 컬럼 비효율: Kafka는 행 지향(row-oriented) 구조다. 분석 쿼리에서 특정 컬럼만 읽어도 전체 행을 역직렬화해야 한다.
  • Stateful 업데이트 어려움: 기본키 기반 upsert를 구현하려면 별도의 상태 저장소(RocksDB 등)와 복잡한 파이프라인이 필요하다.

Fluss가 노리는 위치

Fluss는 Kafka를 대체하거나 Iceberg를 없애는 것이 아니다. 스트리밍 레이어를 컬럼 지향으로 재설계하고, Iceberg/Paimon과는 Tiering Service로 연결해 단일 저장 계층처럼 보이도록 만드는 것이 목표다.

데이터 생산자
애플리케이션 / CDC / IoT 이벤트
↓ Apache Arrow 컬럼형 포맷으로 수신
Fluss 스트리밍 스토리지 계층
Log Table
append-only · 오프셋 순서 보장 · Arrow 컬럼형
PK Table
기본키 upsert · KV Store + Log Store 이중 구조
↓ Tiering Service (자동 오프로딩)
레이크하우스 테이블 포맷
Apache Iceberg
장기 보관 · 대규모 배치 분석
Apache Paimon / Hudi
추가 레이크하우스 포맷 지원
↓ 쿼리 엔진 통합
분석 레이어
Apache Flink · Apache Spark · Apache Doris
Apache Fluss 아키텍처 개요

핵심 개념: Log Table과 PK Table

Log Table

Log Table은 기존 Kafka 토픽과 가장 유사한 추상화다.

  • append-only: 데이터는 추가만 가능하며 삭제나 수정이 없다.
  • 오프셋 기반 순서 보장: 각 레코드는 단조 증가하는 오프셋으로 식별된다. 컨슈머는 특정 오프셋부터 재생(replay)할 수 있다.
  • Apache Arrow 컬럼형 저장: 기존 Kafka는 바이트 배열(행 직렬화)로 저장하지만, Fluss Log Table은 Arrow 컬럼형 청크로 저장한다. 분석 쿼리에서 특정 컬럼만 읽으면 해당 컬럼 데이터만 읽어 역직렬화한다.

이 구조 덕분에 컬럼 프루닝(column pruning)이 스트리밍 단계에서 적용된다. Flink/Spark SQL에서 SELECT event_time, user_id FROM events와 같이 일부 컬럼만 사용하는 쿼리라면, 모든 컬럼을 읽어올 필요가 없어진다.

PK Table

PK Table은 기본키(Primary Key) 기반 upsert를 지원하는 추상화다.

내부적으로 두 저장소를 결합한다.

  • KV Store: 최신 버전 데이터를 기본키로 빠르게 조회한다. 포인트 룩업과 범위 스캔을 지원한다.
  • Log Store: 변경 이력(change log)을 순서대로 보존한다. 다운스트림 컨슈머가 CDC(Change Data Capture) 스타일로 변경 사항을 구독할 수 있다.

이 이중 구조로 다음 두 가지 접근 패턴을 동시에 지원한다.

패턴사용 저장소예시
최신 상태 조회KV Store사용자 프로필 포인트 룩업
변경 스트림 구독Log Store실시간 집계, Flink 스트리밍 조인

Tiering Service: 스트리밍 → 레이크하우스 자동 전환

스트리밍 데이터는 영구 보관이나 대규모 분석에 적합하지 않다. Fluss의 Tiering Service는 이 문제를 자동으로 처리한다.

동작 방식

  1. Fluss 내부의 스트리밍 세그먼트가 설정된 시간 또는 크기 임계값에 도달하면 Tiering Service가 트리거된다.
  2. 세그먼트를 Iceberg 또는 Paimon 포맷으로 변환해 오브젝트 스토리지(S3, GCS, HDFS)에 기록한다.
  3. Fluss 내부에는 최근 N시간 또는 N일치 스트리밍 세그먼트만 보유한다.

쿼리 투명성

Flink/Spark SQL에서 Fluss 테이블과 Iceberg 테이블을 별도로 쿼리할 수도 있고, Fluss가 제공하는 유니온 뷰를 통해 최신 스트리밍 데이터와 오프로딩된 히스토리 데이터를 하나의 테이블처럼 쿼리할 수도 있다.

SELECT *
FROM fluss_catalog.orders
WHERE order_date >= '2026-08-01'

이 쿼리는 자동으로 "최근 데이터는 Fluss 스트리밍 세그먼트에서, 오래된 데이터는 Iceberg에서" 읽는다.


성능: 행 지향 대비 10배 스트리밍 읽기

Fluss 팀이 공개한 벤치마크(Flink 스트리밍 잡, 12컬럼 스키마 중 3컬럼 사용 패턴) 기준:

비교 항목행 지향 스트리밍 스토리지Fluss (Arrow 컬럼형)
스트리밍 읽기 처리량1× (기준)10×
네트워크 전송량 (3컬럼 조회)전체 행 전송선택 컬럼만 전송
역직렬화 비용전체 행 역직렬화선택 컬럼만 처리

주의: 모든 컬럼을 사용하는 쿼리에서는 행 지향 대비 이점이 크지 않다. 컬럼 선택률이 낮을수록(좁은 스키마 쿼리) 효과가 크다.


Flink, Spark, Doris 통합

Apache Flink

Fluss의 1차 통합 대상이자 가장 성숙한 연동이다. Jark Wu(Flink PMC, Alibaba Cloud) 주도로 개발되었다.

  • Flink SQL에서 Fluss 테이블을 Source / Sink 로 직접 사용한다.
  • PK Table의 Log Store는 Flink CDC Source로 동작해 변경 사항을 스트리밍으로 처리한다.
  • KV Store는 Flink Lookup Join의 Dimension Table로 사용 가능하다(포인트 룩업 최적화 포함).

Apache Spark

Spark Structured Streaming에서 Fluss를 소스로 사용할 수 있다. 배치 쿼리에서는 Iceberg 오프로딩 경유로 처리하는 것이 일반적이다.

Apache Doris

Doris(실시간 OLAP 엔진)는 Fluss PK Table의 KV Store를 직접 스캔할 수 있다. 수초 이내의 지연으로 최신 데이터를 Doris에서 집계 쿼리할 수 있다.


ASF TLP 졸업의 의미

TLP(Top-Level Project) 졸업은 Apache Software Foundation의 거버넌스 기준을 통과했음을 의미한다.

  • 커미터 다양성: 단일 회사(Alibaba Cloud)에서 시작했으나, 졸업 시점에는 여러 조직의 컨트리뷰터가 참여하는 독립 거버넌스 구조를 갖췄다.
  • 프로덕션 사용: 실제 프로덕션 환경에서 사용 중인 케이스(Case)가 PMC에 제출된 상태.
  • 장기 안정성 신호: TLP 지위는 장기 유지보수, 보안 릴리스, 브레이킹 체인지 최소화에 대한 약속을 수반한다.

Kafka, Flink, Iceberg가 모두 ASF TLP라는 점을 고려하면, Fluss가 같은 생태계에서 표준 스트리밍 스토리지로 자리잡기 위한 조건을 갖춘 셈이다.


Kafka와 어떻게 다른가

비교 항목Apache KafkaApache Fluss
저장 포맷행(bytes)컬럼(Apache Arrow)
upsert 지원없음(외부 상태 저장소 필요)PK Table 기본 지원
레이크하우스 연동별도 커넥터/잡 필요Tiering Service 내장
SQL 쿼리ksqlDB (별도)Flink SQL / Doris 네이티브
주요 강점메시지 큐 성숙도, 생태계 광대분석 워크로드 최적화, 레이크하우스 통합

Fluss는 Kafka를 대체하기보다 분석 워크로드 특화 스트리밍 스토리지로 포지셔닝한다. 범용 메시지 큐 역할보다 레이크하우스 파이프라인의 스트리밍 레이어 역할에 초점을 맞춘다.


운영 체크리스트

  • [ ] 기존 Kafka 파이프라인 중 분석 쿼리 비율 측정: Fluss 전환 효과는 분석 쿼리 비중이 높을수록 크다
  • [ ] 스키마 컬럼 수와 쿼리 선택률 파악: 좁은 스키마 쿼리(전체 컬럼의 25% 이하)에서 10× 효과 최대화
  • [ ] PK Table 필요 여부 판단: upsert가 없는 순수 이벤트 로그는 Log Table, CDC/upsert 패턴은 PK Table 선택
  • [ ] Tiering Service 대상 포맷 결정: 기존 Iceberg 사용 중이면 Iceberg로 오프로딩, Paimon 사용 시 Paimon으로
  • [ ] Flink 버전 호환성 확인: Fluss Flink 커넥터는 Flink 1.19 이상 권장 (커넥터 문서에서 버전 매트릭스 확인)
  • [ ] Doris 통합 계획: Doris를 OLAP 레이어로 사용한다면 Fluss KV Store 직접 스캔 방식 검토
  • [ ] 보존 기간 정책: 스트리밍 세그먼트 보존 기간(Tiering 트리거 기준)과 Iceberg 파티션 TTL 별도 설정
  • [ ] ASF TLP 릴리스 채널 구독: 보안 공지와 브레이킹 체인지 공지를 위해 [email protected] 구독

요점 정리

  • Apache Fluss는 2026년 8월 6일 ASF TLP로 졸업한 스트리밍 스토리지 엔진으로, 레이크하우스 파이프라인의 Kafka 역할을 컬럼형 Arrow 기반으로 재설계한다.
  • Log Table(이벤트 스트림)과 PK Table(KV Store + Log Store 이중 구조)로 순수 이벤트 로그와 upsert 패턴 모두를 커버한다.
  • Tiering Service가 스트리밍 세그먼트를 Iceberg/Paimon으로 자동 오프로딩해 레이크하우스 통합을 단순화한다.
  • 컬럼 프루닝으로 스트리밍 읽기 속도 10×, Flink/Spark/Doris 네이티브 통합, ASF 거버넌스 기반의 장기 안정성.

References

  • Apache Fluss 공식 사이트: https://fluss.apache.org/
  • Apache Fluss GitHub: https://github.com/apache/fluss
  • Apache Fluss ASF TLP 발표 (2026-08-06): https://blogs.apache.org/foundation/entry/the-apache-software-foundation-announces-apache-fluss-as-a-top-level-project
  • Apache Fluss 문서 (Flink 통합): https://fluss.apache.org/docs/engine-flink/
  • Apache Paimon 공식 사이트: https://paimon.apache.org/
  • Apache Iceberg 공식 사이트: https://iceberg.apache.org/