Engineering Note
Apache Fluss 첫 실행: Flink Quickstart에서 확인할 것들
Apache Fluss를 로컬 실습 환경에서 Flink와 연결하고, Flink SQL로 카탈로그·테이블 모델·기본 키 조회와 갱신을 확인합니다.
Apache Fluss를 처음 실행할 때 먼저 확인할 것은 Docker Compose 설정 자체가 아닙니다. Flink가 Fluss를 카탈로그로 연결하고, 최신 상태와 이벤트 이력을 서로 다른 테이블로 다루는 방식입니다. 이 글은 Flink SQL로 그 흐름을 직접 확인합니다.
공식 Flink Quickstart는 이 실습에 필요한 Fluss, Flink, S3 호환 저장소를 한 대의 개발 장비에서 함께 시작합니다. 이 환경은 한 프로세스짜리 단독 실행이 아니라 여러 구성 요소가 실제 역할을 수행하는 작은 통합 환경입니다. 운영 배포를 대체하지는 않지만, Kafka·Iceberg·별도 클라우드 저장소를 먼저 준비하지 않고 테이블 생성·쓰기·기본 키 조회를 확인할 수 있습니다.
먼저 가져갈 결론
이 실습에서 직접 확인할 수 있는 것은 다음입니다.
- Fluss에는 Flink 카탈로그를 통해 데이터베이스와 테이블을 만듭니다.
PRIMARY KEY를 둔 테이블은 최신 키 상태를 저장하고 키로 조회할 수 있습니다.- Flink SQL 클라이언트 이미지에는 Fluss 커넥터가 포함되어 있어 JAR를 따로 내려받지 않아도 됩니다.
- RustFS는 로컬 S3 호환 원격 저장소 역할을 하고, ZooKeeper는 현재 Fluss 클러스터의 조정에 사용됩니다.
반대로 이 구성은 운영 환경의 보안, 여러 TabletServer, 고가용성, 실제 클라우드 자격 증명, Kafka 연동을 검증하지 않습니다. 실습 성공을 운영 설계의 결론으로 과장하지 않기 위해 이 경계를 먼저 두는 편이 좋습니다.
실습 환경 준비하기
이 글에서는 공식 실습 환경을 그대로 재현하기 위해 Docker와 Docker Compose 플러그인을 사용합니다. Compose는 Fluss의 기능을 설명하는 주제가 아니라, 필요한 구성 요소를 한 번에 시작하는 방법입니다. 공식 문서는 Docker 27.4.0과 Compose v2.30.3에서 예제를 검증했으며, 최근 Compose v2 사용을 권장합니다.
docker version
docker compose version
Java나 로컬 Flink 설치는 필요 없습니다. 이 실습 환경에는 Flink와 Fluss-Flink 커넥터가 이미 포함되어 있습니다.
공식 실습 환경 만들기
작업 디렉터리를 만듭니다.
mkdir fluss-quickstart-flink
cd fluss-quickstart-flink
그 안에 공식 문서의 docker-compose.yml 전체 구성을 그대로 저장합니다. 새 창에서 services:부터 마지막 volumes:까지 복사해 docker-compose.yml로 저장한 뒤 이 글로 돌아오면 됩니다. 이 글은 Compose 파일을 해설하는 대신, 시작된 구성 요소가 Fluss 실습에서 맡는 역할에 집중합니다. 이미지 태그와 저장소 설정이 바뀔 수 있으므로 긴 설정은 공식 Quickstart를 기준으로 삼습니다. 2026-07-29 기준 구성은 다음 서비스를 올립니다.
| 구분 | 서비스 | 역할 |
|---|---|---|
| 원격 저장소 | rustfs, rustfs-init | S3 호환 저장소와 fluss 버킷 초기화 |
| Fluss | coordinator-server, tablet-server, zookeeper | 메타데이터·조정과 데이터 관리 |
| Flink | jobmanager, taskmanager, sql-client | SQL 실행과 스트림·배치 작업 실행 |
이 구성은 apache/fluss:0.9.1-incubating, apache/fluss-quickstart-flink:1.20-0.9.1-incubating, zookeeper:3.9.2, RustFS 이미지를 사용합니다. Flink Quickstart 이미지는 Fluss-Flink 커넥터, flink-faker, S3 파일시스템 지원을 포함하므로 추가 JAR가 필요하지 않습니다.
rustfsadmin 자격 증명과 S3 설정은 로컬 체험 전용입니다. 운영 환경에서는 제공업체의 자격 증명 공급자, 비밀값 관리, TLS, 접근 제어와 저장소 정책을 따로 설계해야 합니다. 공식 Compose 파일의 data.dir: /tmp/fluss/data도 컨테이너 재생성에 견디는 운영용 로컬 디스크 설정이 아닙니다.
환경 시작하고 상태 확인하기
컨테이너를 시작합니다.
docker compose up -d
docker compose ps
rustfs-init은 버킷을 한 번 만들고 종료하므로 Exited (0)로 보여도 정상입니다. 나머지 서비스가 올라오는 동안에는 다음 로그가 가장 유용합니다.
docker compose logs -f coordinator-server tablet-server taskmanager
브라우저에서 http://localhost:8083를 열면 Flink 웹 UI를 볼 수 있습니다. http://localhost:9001은 RustFS 콘솔이며, 이 로컬 예제에서만 rustfsadmin / rustfsadmin으로 로그인합니다. 이미 8083·9000·9001 포트를 쓰고 있다면 Compose 파일의 호스트 쪽 포트 번호만 바꾸면 됩니다.
시작 전에 SQL 객체 계층 이해하기
이제부터 Flink SQL에서 객체를 만들고 사용합니다. 이름은 보통 카탈로그.데이터베이스.테이블 순서로 읽으면 됩니다.
| 객체 | 이 Quickstart에서의 의미 |
|---|---|
| 카탈로그 | Flink가 Fluss의 데이터베이스와 테이블을 인식하고 관리하도록 연결하는 카탈로그 구현입니다. fluss_catalog을 등록하면 Flink SQL에서 Fluss의 객체를 만들고 읽고 쓸 수 있습니다. |
| 데이터베이스 | 한 카탈로그 안에서 테이블을 나누는 이름 공간입니다. 이 글에서는 demo를 만듭니다. |
| 테이블 | Fluss에 실제로 데이터를 저장하고 읽는 논리적 단위입니다. 뒤에서 최신 상태용 customer_profile과 이력용 order_events를 만듭니다. |
따라서 뒤에서 만드는 customer_profile의 전체 이름은 fluss_catalog.demo.customer_profile입니다. USE CATALOG fluss_catalog와 USE demo를 실행하면 이후 SQL에서는 이 긴 이름 대신 테이블 이름만 쓸 수 있습니다. 반면 원본은 SQL 클라이언트 세션에 미리 만들어진 임시 테이블이며, default_catalog.default_database라는 논리 이름으로 참조합니다. 그래서 적재할 때는 그 전체 이름을 명시합니다.
Flink SQL 클라이언트 열기
다음 명령은 SQL 클라이언트를 실행할 때만 컨테이너를 하나 만들고 터미널을 연결합니다.
docker compose run --rm sql-client
이제 프롬프트에서 Fluss 카탈로그를 등록합니다. coordinator-server:9123은 Compose 네트워크 안에서 CoordinatorServer를 찾는 주소입니다. 호스트의 localhost를 쓰면 안 됩니다.
CREATE CATALOG fluss_catalog WITH (
'type' = 'fluss',
'bootstrap.servers' = 'coordinator-server:9123'
);
USE CATALOG fluss_catalog;
SHOW DATABASES;
카탈로그 설정은 기본적으로 SQL 클라이언트 세션 사이에 유지되지 않습니다. 클라이언트를 종료했다가 새로 열면 위의 CREATE CATALOG와 USE CATALOG를 다시 실행해야 합니다. 이는 Quickstart의 오류가 아니라 Flink 카탈로그 설정의 기본 동작입니다.
이 글에서는 삽입·갱신·조회가 한 번에 끝나는 실습을 만들기 위해 세션을 배치 모드와 동기 DML로 설정합니다. 기본값에서는 INSERT와 UPDATE가 비동기로 제출되어 바로 뒤의 조회가 작업 완료보다 먼저 실행될 수 있습니다.
SET 'sql-client.execution.result-mode' = 'tableau';
SET 'execution.runtime-mode' = 'batch';
SET 'table.dml-sync' = 'true';
원본 데이터를 이용해 두 테이블 만들기
Flink Quickstart 이미지의 SQL 클라이언트 세션에는 faker 커넥터로 만든 source_order, source_customer, source_nation 임시 테이블이 미리 만들어져 있습니다. 수동으로 두 행만 입력하는 대신, 이 원본을 Fluss 테이블에 넣어 봅시다. 먼저 원본 정의를 확인합니다.
SHOW CREATE TABLE `default_catalog`.`default_database`.source_customer;
SHOW CREATE TABLE `default_catalog`.`default_database`.source_order;
출력에서 source_order의 rows-per-second = 10과 number-of-rows = 10000을 확인해 둡니다. 뒤에서 주문 10,000건 적재가 느린 이유를 이해하는 기준이 됩니다.
이제 고객의 최신 상태를 담는 기본 키 테이블(Primary Key Table)과, 주문 이력을 계속 추가하는 로그 테이블(Log Table)을 만듭니다.
CREATE DATABASE demo;
USE demo;
CREATE TABLE customer_profile (
customer_id INT NOT NULL,
name STRING,
membership STRING,
account_balance DECIMAL(15, 2),
PRIMARY KEY (customer_id) NOT ENFORCED
) WITH (
'bucket.num' = '1'
);
CREATE TABLE order_events (
order_id BIGINT,
customer_id INT NOT NULL,
total_price DECIMAL(15, 2),
ordered_on DATE,
order_priority STRING,
clerk STRING
) WITH (
'bucket.num' = '1'
);
customer_profile의 PRIMARY KEY는 같은 고객의 새 값을 최신 상태로 만듭니다. 반면 order_events에는 기본 키가 없으므로 주문이 한 건씩 쌓이는 Log Table입니다. order_events.order_id는 주문을 식별하기 위한 일반 열일 뿐 기본 키가 아니며, 같은 값이 다시 들어와도 중복을 막거나 기존 행을 대체하지 않습니다. 두 모델을 나란히 만들어야 Fluss가 이벤트와 최신 상태를 어떻게 구분하는지 체감할 수 있습니다.
생성된 고객·주문 데이터 적재하기
Quickstart가 미리 만든 source_order는 초당 10건(rows-per-second = 10)을 생성하도록 설정되어 있습니다. 그래서 원본 10,000건 전체를 그대로 적재하면 데이터 생성만 약 1,000초가 걸립니다. Fluss 쓰기 성능의 문제가 아니라, 실습용 faker 원본의 의도적인 속도 제한입니다.
첫 실행에서는 고객 전체와 주문 500건만 적재합니다. 앞에서 table.dml-sync를 켰으므로 각 INSERT가 끝난 뒤 다음 명령으로 넘어갑니다.
INSERT INTO customer_profile
SELECT cust_key, name, mktsegment, acctbal
FROM `default_catalog`.`default_database`.source_customer;
INSERT INTO order_events
SELECT order_key, cust_key, total_price, order_date, order_priority, clerk
FROM `default_catalog`.`default_database`.source_order
LIMIT 500;
대안: 10,000건을 빠르게 생성하기
처음부터 주문 10,000건을 모두 확인하고 싶다면, 위의 order_events 적재문은 실행하지 말고 아래 대안을 선택합니다. fast_order_source는 같은 종류의 주문 데이터를 초당 2,000건으로 생성하고, 정확히 10,000건에서 멈춥니다. order_events가 비어 있는 새 실습 환경에서 실행해야 예상 건수가 맞습니다.
rows-per-second는 원본 생성 속도의 상한일 뿐, Fluss에 쓰는 전체 시간이 정확히 그 속도를 보장하지는 않습니다. 그래도 로컬 실습에서 생성 단계가 병목이 되는 일은 피할 수 있습니다.
CREATE TEMPORARY TABLE fast_order_source (
order_id BIGINT,
customer_id INT,
total_price DECIMAL(15, 2),
ordered_on DATE,
order_priority STRING,
clerk STRING
) WITH (
'connector' = 'faker',
'number-of-rows' = '10000',
'rows-per-second' = '2000',
'fields.order_id.expression' = '#{number.numberBetween ''0'',''100000000''}',
'fields.customer_id.expression' = '#{number.numberBetween ''0'',''20''}',
'fields.total_price.expression' = '#{number.randomDouble ''3'',''1'',''1000''}',
'fields.ordered_on.expression' = '#{date.past ''100'' ''DAYS''}',
'fields.order_priority.expression' = '#{regexify ''(low|medium|high){1}''}',
'fields.clerk.expression' = '#{regexify ''(Clerk1|Clerk2|Clerk3|Clerk4){1}''}'
);
INSERT INTO order_events
SELECT order_id, customer_id, total_price, ordered_on, order_priority, clerk
FROM fast_order_source;
다음 조회로 실제 적재량과 주문 데이터를 확인합니다. 기본 경로의 order_count는 500이고, 빠른 생성 대안을 택했다면 10,000이어야 합니다. Fluss의 현재 배치 읽기에서는 Log Table에 대한 COUNT(*)와 LIMIT 미리보기를 지원하므로, 이 두 쿼리를 사용합니다.
SELECT COUNT(*) AS customer_count FROM customer_profile;
SELECT COUNT(*) AS order_count FROM order_events;
SELECT
order_id,
customer_id,
total_price,
ordered_on,
order_priority
FROM order_events
LIMIT 10;
Log Table에 이벤트 추가하기
order_events에는 기본 키가 없으므로 같은 키를 찾아 현재 값을 바꾸는 대신, 새 주문 이벤트를 한 행 더 추가합니다. 아래 행은 이미 적재한 주문을 대체하지 않고 이력에 더해집니다.
INSERT INTO order_events VALUES (
999999999,
999999,
CAST(42.00 AS DECIMAL(15, 2)),
DATE '2026-07-29',
'quickstart',
'Quickstart Clerk'
);
SELECT COUNT(*) AS order_count FROM order_events;
기본 경로에서는 order_count가 501, 빠른 생성 대안에서는 10,001이 됩니다. 이처럼 Log Table은 이벤트를 계속 쌓는 데 맞고, 기본 키로 한 행을 찾는 점 조회는 Primary Key Table에만 지원됩니다. 그래서 이 실습에서는 로그 테이블의 추가를 행 수 증가로 확인합니다.
이제 방금 추가한 이벤트를 제자리에서 바꾸려 시도해 봅니다. 아래 문은 실패해야 정상입니다. Log Table은 추가만 지원하며 갱신·삭제를 지원하지 않기 때문입니다. 오류 문구는 Flink와 Fluss 버전에 따라 달라질 수 있습니다.
UPDATE order_events
SET order_priority = 'corrected'
WHERE order_id = 999999999;
고객별 주문 수나 결제액처럼 Log Table 전체를 GROUP BY하는 분석은 이 배치 읽기 경로의 대상이 아닙니다. 이를 계속 계산하려면 Flink 스트리밍 작업으로 집계 테이블을 만들거나, 다음 단계에서 레이크하우스 티어링과 분석 엔진을 사용합니다.
Primary Key Table에서 같은 키의 최신 상태 바꾸기
이제 주문 이력은 그대로 둔 채, 고객 한 명의 최신 상태만 바꿉니다. 생성 원본에 어떤 등급이 들어 있는지에 의존하지 않도록, 실습 전용 키 999999에 변경 전 값을 먼저 기록합니다.
INSERT INTO customer_profile VALUES (
999999,
'Quickstart User',
'quickstart-before',
CAST(0.00 AS DECIMAL(15, 2))
);
이제 변경 전 상태를 확인합니다.
SELECT customer_id, name, membership, account_balance
FROM customer_profile
WHERE customer_id = 999999;
이제 같은 키 999999에 새 상태를 다시 INSERT하고 조회합니다. Primary Key Table에서는 같은 기본 키의 새 행이 이전 행과 함께 남지 않고 최신 상태가 됩니다.
INSERT INTO customer_profile VALUES (
999999,
'Quickstart User',
'quickstart-after',
CAST(0.00 AS DECIMAL(15, 2))
);
SELECT customer_id, name, membership, account_balance
FROM customer_profile
WHERE customer_id = 999999;
변경 전 결과와 마지막 결과를 비교하면 실습용 행의 membership이 quickstart-before에서 quickstart-after로 바뀐 것을 확인할 수 있습니다. 같은 키로 두 번 INSERT했어도 조회 결과가 한 행인 점이 Log Table과의 차이입니다. 이 흐름으로 Flink가 계산·SQL 실행 계층이고, Fluss가 이벤트 이력과 기본 키 테이블의 최신 상태를 서로 다른 모델로 저장·조회시키는 계층임을 확인할 수 있습니다.
다음 단계: 조회 조인으로 주문을 보강하기
공식 예제는 여기서 고객·국가 Primary Key Table을 더 만들고, 주문 이벤트에 두 테이블을 조회 조인합니다. 주문 처리 시점의 고객·국가 정보를 붙이는 실제 스트리밍 경로를 보려면 공식 Flink Quickstart의 스트리밍 적재와 조회 조인 단계를 이어서 실행하면 됩니다.
여기까지의 구성에는 Kafka가 없습니다. Fluss와 Flink 사이의 테이블 작성·조회 경로만 보려는 목적에는 Kafka가 필요하지 않기 때문입니다. 기존 Kafka 이벤트를 Fluss 테이블에 연결하는 문제는 별도의 데이터 경로와 장애·재처리 설계가 필요합니다.
다음 단계는 Lakehouse Quickstart다
이번 실습은 Fluss의 실시간 테이블과 Flink SQL에 집중합니다. 티어링과 통합 읽기를 보려면 Streaming Lakehouse Quickstart로 넘어가면 됩니다. 이 예제는 Paimon 또는 Iceberg와 RustFS를 추가하고, 레이크하우스 계층을 활성화한 테이블을 만듭니다.
복잡도도 같이 늘어납니다. Paimon 경로에는 Fluss 서버 측 S3 플러그인이, Iceberg 경로에는 Iceberg·JDBC 관련 JAR와 PostgreSQL 카탈로그가 추가됩니다. 따라서 처음에는 이 글의 Flink Quickstart로 테이블 모델과 SQL 연결을 확인한 뒤, 정말 긴 이력과 통합 읽기가 필요한 경우에만 Lakehouse 구성을 이어가는 편이 좋습니다.
정리하기
SQL 클라이언트에서는 quit;로 나옵니다. 실습 데이터를 포함해 Compose 볼륨까지 지우려면 다음을 실행합니다.
docker compose down -v
-v는 rustfs-data 볼륨의 데이터를 삭제합니다. 나중에 같은 데이터를 다시 확인하고 싶다면 docker compose down만 사용하세요.
결론
Flink Quickstart의 핵심은 Compose 자체가 아닙니다. Fluss, Flink, S3 호환 저장소, 조정 서비스가 맞물린 환경에서 Flink SQL로 테이블 모델을 확인할 수 있다는 점입니다.
실습 환경을 시작한다.
Flink SQL에서 Fluss 카탈로그를 등록한다.
기본 키 테이블을 만들고 최신 상태를 조회한다.
그 다음에야 스트림 적재, 조회 조인, Lakehouse 티어링으로 범위를 넓힌다.
이 순서라면 Kafka나 Iceberg를 처음부터 한꺼번에 판단하지 않고, Fluss가 해결하려는 저장소와 테이블 모델을 먼저 직접 확인할 수 있습니다.
함께 읽기 좋은 글
- Apache Fluss란: Kafka·Flink·Iceberg 사이의 스트리밍 레이크하우스 - Quickstart에서 실행한 구성 요소가 데이터 경로에서 맡는 역할을 비교합니다.
- Kafka의 디스크는 어디로 가고 있나: Object Storage와 새로운 스트리밍 계층 - 원격 저장소와 스트리밍 저장소의 설계 변화를 함께 봅니다.
- 메트릭 수집의 Pull과 Push는 무엇이 다른가 - 실습 환경을 운영 수준으로 넓힐 때 필요한 관측 가능성의 책임을 정리합니다.