Materialized Table 배포

Materialized Table 배포 (Deployment)

Materialized Table을 만들고 운영하는 것은 여러 컴포넌트의 협업을 필요로 해요. 이 문서는 Materialized Table의 완전한 배포 솔루션을 체계적으로 설명하며, 아키텍처 개요, 환경 준비, 배포 절차, 운영 관례를 다룹니다.

출처: Materialized Table Deployment

본문

Architecture Introduction

  • Client: Flink SQL Gateway와 상호작용할 수 있는 모든 클라이언트일 수 있어요. 예: SQL Client, Flink JDBC Driver 등.
  • Flink SQL Gateway: Materialized Table의 생성·변경·삭제를 지원해요. 또한 full mode Materialized Table을 주기적으로 refresh하는 임베디드 워크플로우 스케줄러 역할도 해요.
  • Flink Cluster: Materialized Table을 refresh하는 파이프라인이 Flink 클러스터에서 실행돼요.
  • Catalog: Materialized Table 메타데이터의 생성·조회·수정·삭제를 관리해요.
  • Catalog Store: Materialized Table 관련 작업에서 메타데이터를 가져오기 위해 카탈로그를 자동으로 초기화하도록 카탈로그 속성의 영속화를 지원해요.

Illustration of Flink Materialized Table Architecture

Deployment Preparation

Materialized Table refresh 잡은 현재 다음 클러스터 환경에서 실행을 지원해요:

Materialized Table은 SQL Gateway를 통해서만 만들 수 있어요. 메타데이터 영속화와 잡 스케줄링을 위한 특정 구성이 필요해요.

Configure Catalog Store

config.yaml에 카탈로그 속성을 영속화하도록 catalog store 구성을 추가하세요.

table:
  catalog-store:
    kind: file
    file:
      path: {path_to_catalog_store} # Replace with the actual path

자세한 내용은 Catalog Store를 참고하세요.

Configure Workflow Scheduler Plugin

주기적인 refresh 잡 스케줄링을 위해 config.yaml에 workflow scheduler 구성을 추가하세요. 현재는 embedded 스케줄러만 지원돼요:

workflow-scheduler:
  type: embedded

Start SQL Gateway

다음을 사용해 SQL Gateway를 시작하세요:

./sql-gateway.sh start

주의: Catalog가 materialized table 생성을 지원해야 하는데, 현재는 Paimon Catalog만 지원해요.

Operation Guide

Connecting to SQL Gateway

SQL Client를 사용하는 예:

./sql-client.sh gateway --endpoint {gateway_endpoint}:{gateway_port}

Creating Materialized Tables

Refresh Jobs Running on Standalone Cluster

Flink SQL> SET 'execution.mode' = 'remote';
[INFO] Execute statement succeeded.

FLINK SQL> CREATE MATERIALIZED TABLE my_materialized_table
> ... ;
[INFO] Execute statement succeeded.

Refresh Jobs Running in Session Mode

세션 모드의 경우 yarn-session이나 kubernetes-session에 문서화된 대로 세션 클러스터를 미리 생성하세요.

Kubernetes session mode:

Flink SQL> SET 'execution.mode' = 'kubernetes-session';
[INFO] Execute statement succeeded.

Flink SQL> SET 'kubernetes.cluster-id' = 'flink-cluster-mt-session-1';
[INFO] Execute statement succeeded.

Flink SQL> CREATE MATERIALIZED TABLE my_materialized_table
> ... ;
[INFO] Execute statement succeeded.

execution.modekubernetes-session으로 설정하고, 기존 Kubernetes 세션 클러스터에 해당하는 유효한 kubernetes.cluster-id를 지정하세요.

YARN session mode:

Flink SQL> SET 'execution.mode' = 'yarn-session';
[INFO] Execute statement succeeded.

Flink SQL> SET 'yarn.application.id' = 'application-xxxx';
[INFO] Execute statement succeeded.

Flink SQL> CREATE MATERIALIZED TABLE my_materialized_table
> ... ;
[INFO] Execute statement succeeded.

execution.modeyarn-session으로 설정하고, 기존 YARN 세션 클러스터에 해당하는 유효한 yarn.application.id를 지정하세요.

Refresh Jobs Running in Application Mode

Kubernetes application mode:

Flink SQL> SET 'execution.mode' = 'kubernetes-application';
[INFO] Execute statement succeeded.

Flink SQL> SET 'kubernetes.cluster-id' = 'flink-cluster-mt-application-1';
[INFO] Execute statement succeeded.

Flink SQL> CREATE MATERIALIZED TABLE my_materialized_table
> ... ;
[INFO] Execute statement succeeded.

execution.modekubernetes-application으로 설정하세요. kubernetes.cluster-id는 선택사항이며, 설정하지 않으면 자동으로 생성돼요.

YARN application mode:

Flink SQL> SET 'execution.mode' = 'yarn-application';
[INFO] Execute statement succeeded.

Flink SQL> CREATE MATERIALIZED TABLE my_materialized_table
> ... ;
[INFO] Execute statement succeeded.

execution.modeyarn-application으로 설정하세요. yarn.application.id는 설정할 필요가 없으며 제출 중에 자동으로 생성돼요.

Maintenance Operations

클러스터 정보(예: execution.mode 또는 kubernetes.cluster-id)는 이미 카탈로그에 영속화되어 있어서, Materialized Table의 refresh 잡을 일시 중지하거나 재개할 때 설정할 필요가 없어요.

Suspend Refresh Job

-- Suspend the MATERIALIZED TABLE refresh job
Flink SQL> ALTER MATERIALIZED TABLE my_materialized_table SUSPEND;
[INFO] Execute statement succeeded.

Resume Refresh Job

-- Resume the MATERIALIZED TABLE refresh job
Flink SQL> ALTER MATERIALIZED TABLE my_materialized_table RESUME;
[INFO] Execute statement succeeded.

Modify Query Definition

-- Modify the MATERIALIZED TABLE query definition
Flink SQL> ALTER MATERIALIZED TABLE my_materialized_table
> AS SELECT
> ... ;
[INFO] Execute statement succeeded.

더 알아보기 (Learn more)