머티리얼라이즈드 테이블 퀵스타트
머티리얼라이즈드 테이블 퀵스타트 (Materialized Table Quickstart Guide)
이 가이드는 머티리얼라이즈드 테이블(materialized table)을 빠르게 이해하고 시작하도록 도와줘요. 환경 설정과 CONTINUOUS·FULL 모드에서 머티리얼라이즈드 테이블 생성·변경·삭제를 포함해요.
출처: 문서
본문
환경 설정 (Environment Setup)
디렉터리 준비 (Directory Preparation)
아래 예제 경로는 실제 머신의 경로로 바꿔주세요.
- Catalog Store와 test-filesystem Catalog용 디렉터리를 만듭니다:
# Directory for File Catalog Store to save catalog information
mkdir -p {catalog_store_path}
# Directory for test-filesystem Catalog to save table metadata and table data
mkdir -p {catalog_path}
# Directory for the default database of test-filesystem Catalog
mkdir -p {catalog_path}/mydb
- Checkpoints와 Savepoints용 디렉터리를 만듭니다:
mkdir -p {checkpoints_path}
mkdir -p {savepoints_path}
리소스 준비 (Resource Preparation)
여기서의 방법은 로컬 설치에 기록된 단계와 비슷해요. Flink는 Linux, Mac OS X, Cygwin(Windows용) 같은 UNIX 계열 운영체제에서 실행할 수 있어요.
최신 Flink 바이너리 패키지를 다운로드하고 추출해요:
tar -xzf flink-*.tgz
test-filesystem 커넥터를 다운로드해 lib 디렉터리에 넣어요:
cp flink-table-filesystem-test-utils-{VERSION}.jar flink-*/lib/
구성 준비 (Configuration Preparation)
config.yaml 파일을 편집하고 다음 구성을 추가해요:
execution:
checkpoints:
dir: file://{checkpoints_path}
# Configure file catalog store
table:
catalog-store:
kind: file
file:
path: {catalog_store_path}
# Configure embedded scheduler
workflow-scheduler:
type: embedded
# Configure SQL gateway address and port
sql-gateway:
endpoint:
rest:
address: 127.0.0.1
port: 8083
Flink 클러스터 시작 (Start Flink Cluster)
다음 스크립트로 클러스터를 로컬에서 시작해요:
./bin/start-cluster.sh
SQL Gateway 시작 (Start SQL Gateway)
다음 스크립트로 SQL Gateway를 로컬에서 시작해요:
./bin/sql-gateway.sh start
SQL Client 시작 (Start SQL Client)
다음 스크립트로 SQL Client를 로컬에서 시작하고 SQL Gateway에 연결해요:
./bin/sql-client.sh gateway --endpoint http://127.0.0.1:8083
카탈로그와 소스 테이블 생성 (Create Catalog and Source Table)
- test-filesystem 카탈로그를 만듭니다:
CREATE CATALOG mt_cat WITH (
'type' = 'test-filesystem',
'path' = '{catalog_path}',
'default-database' = 'mydb'
);
USE CATALOG mt_cat;
- Source 테이블을 만듭니다:
-- 1. Create Source table and specify the data format is json
CREATE TABLE json_source (
order_id BIGINT,
user_id BIGINT,
user_name STRING,
order_created_at STRING,
payment_amount_cents BIGINT
) WITH (
'format' = 'json',
'source.monitor-interval' = '10s'
);
-- 2. Insert some test data
INSERT INTO json_source VALUES
(1001, 1, 'user1', '2024-06-19', 10),
(1002, 2, 'user2', '2024-06-19', 20),
(1003, 3, 'user3', '2024-06-19', 30),
(1004, 4, 'user4', '2024-06-19', 40),
(1005, 1, 'user1', '2024-06-20', 10),
(1006, 2, 'user2', '2024-06-20', 20),
(1007, 3, 'user3', '2024-06-20', 30),
(1008, 4, 'user4', '2024-06-20', 40);
CONTINUOUS 모드 머티리얼라이즈드 테이블 생성 (Create Continuous Mode Materialized Table)
머티리얼라이즈드 테이블 생성 (Create Materialized Table)
데이터 신선도(freshness)가 30초인 CONTINUOUS 모드 머티리얼라이즈드 테이블을 만듭니다. http://localhost:8081 페이지에서 머티리얼라이즈드 테이블을 지속 갱신하는 Flink 스트리밍 작업이 실행 중인 것을 볼 수 있어요. 체크포인트 간격은 30초예요.
CREATE MATERIALIZED TABLE continuous_users_shops
PARTITIONED BY (ds)
WITH (
'format' = 'debezium-json',
'sink.rolling-policy.rollover-interval' = '10s',
'sink.rolling-policy.check-interval' = '10s'
)
FRESHNESS = INTERVAL '30' SECOND
AS SELECT
user_id,
ds,
SUM (payment_amount_cents) AS payed_buy_fee_sum,
SUM (1) AS PV
FROM (
SELECT user_id, order_created_at AS ds, payment_amount_cents
FROM json_source
) AS tmp
GROUP BY user_id, ds;
머티리얼라이즈드 테이블 일시 중지 (Suspend Materialized Table)
머티리얼라이즈드 테이블의 갱신 파이프라인을 일시 중지해요. http://localhost:8081에서 머티리얼라이즈드 테이블을 지속 갱신하는 Flink 스트리밍 작업이 FINISHED 상태로 전환되는 것을 볼 수 있어요. 일시 중지 연산을 실행하기 전에 savepoint 경로를 설정해야 해요.
-- Set savepoint path before suspending
SET 'execution.checkpointing.savepoint-dir' = 'file://{savepoints_path}';
ALTER MATERIALIZED TABLE continuous_users_shops SUSPEND;
머티리얼라이즈드 테이블 조회 (Query Materialized Table)
머티리얼라이즈드 테이블 데이터를 조회해 데이터가 이미 쓰였는지 확인해요.
SELECT * FROM continuous_users_shops;
머티리얼라이즈드 테이블 재개 (Resume Materialized Table)
머티리얼라이즈드 테이블의 갱신 파이프라인을 재개해요. http://localhost:8081 페이지에서 머티리얼라이즈드 테이블을 지속 갱신하는 새 Flink 스트리밍 작업이 시작되고 지정된 savepoint 경로에서 상태를 복원하는 것을 볼 수 있어요.
ALTER MATERIALIZED TABLE continuous_users_shops RESUME;
머티리얼라이즈드 테이블 삭제 (Drop Materialized Table)
머티리얼라이즈드 테이블을 삭제하면 http://localhost:8081 페이지에서 머티리얼라이즈드 테이블을 지속 갱신하는 Flink 스트리밍 작업이 CANCELED 상태로 전환되는 것을 볼 수 있어요.
DROP MATERIALIZED TABLE continuous_users_shops;
FULL 모드 머티리얼라이즈드 테이블 생성 (Create Full Mode Materialized Table)
머티리얼라이즈드 테이블 생성 (Create Materialized Table)
데이터 신선도가 1분인 FULL 모드 머티리얼라이즈드 테이블을 만듭니다(여기서는 테스트 편의를 위해 freshness를 1분으로 설정). http://localhost:8081에서 머티리얼라이즈드 테이블을 주기적으로 갱신하는 Flink 배치 작업이 1분마다 스케줄링되는 것을 볼 수 있어요.
CREATE MATERIALIZED TABLE full_users_shops
PARTITIONED BY (ds)
WITH (
'format' = 'json',
'partition.fields.ds.date-formatter' = 'yyyy-MM-dd'
)
FRESHNESS = INTERVAL '1' MINUTE
REFRESH_MODE = FULL
AS SELECT
user_id,
ds,
SUM (payment_amount_cents) AS payed_buy_fee_sum,
SUM (1) AS PV
FROM (
SELECT user_id, order_created_at AS ds, payment_amount_cents
FROM json_source
) AS tmp
GROUP BY user_id, ds;
머티리얼라이즈드 테이블 조회 (Query Materialized Table)
오늘 파티션에 데이터를 삽입해요. 최소 1분 기다린 후 머티리얼라이즈드 테이블 결과를 조회하면 오늘 파티션 데이터만 갱신된 것을 알 수 있어요.
INSERT INTO json_source VALUES
(1001, 1, 'user1', CAST(CURRENT_DATE AS STRING), 10),
(1002, 2, 'user2', CAST(CURRENT_DATE AS STRING), 20),
(1003, 3, 'user3', CAST(CURRENT_DATE AS STRING), 30),
(1004, 4, 'user4', CAST(CURRENT_DATE AS STRING), 40);
SELECT * FROM full_users_shops;
과거 파티션 수동 갱신 (Manually Refresh Historical Partition)
파티션 ds='2024-06-20'을 수동으로 갱신하고 머티리얼라이즈드 테이블의 데이터를 확인해요. http://localhost:8081 페이지에서 현재 갱신 작업의 Flink 배치 작업을 찾을 수 있어요.
-- Manually refresh historical partition
ALTER MATERIALIZED TABLE full_users_shops REFRESH PARTITION(ds='2024-06-20');
-- Query materialized table data
SELECT * FROM full_users_shops;
머티리얼라이즈드 테이블 일시 중지 및 재개 (Suspend and Resume Materialized Table)
일시 중지·재개 연산으로 머티리얼라이즈드 테이블에 해당하는 갱신 작업을 제어할 수 있어요. 일시 중지 후에는 머티리얼라이즈드 테이블을 주기적으로 갱신하는 Flink 배치 작업이 스케줄링되지 않아요. 재개 후에는 다시 스케줄링돼요. http://localhost:8081 페이지에서 Flink 작업 스케줄링 상태를 확인할 수 있어요.
-- Suspend background refresh pipeline
ALTER MATERIALIZED TABLE full_users_shops SUSPEND;
-- Resume background refresh pipeline
ALTER MATERIALIZED TABLE full_users_shops RESUME;
머티리얼라이즈드 테이블 삭제 (Drop Materialized Table)
머티리얼라이즈드 테이블을 삭제하면 머티리얼라이즈드 테이블을 주기적으로 갱신하는 Flink 배치 작업이 다시 스케줄링되지 않아요. http://localhost:8081 페이지에서 확인할 수 있어요.
DROP MATERIALIZED TABLE full_users_shops;