머티어리얼라이즈드 테이블 문

머티어리얼라이즈드 테이블 문 (Materialized Table Statements)

Flink SQL은 현재 다음의 머티어리얼라이즈드 테이블(Materialized Table) 문을 지원해요:

출처: 문서

본문

CREATE [OR ALTER] MATERIALIZED TABLE

CREATE [OR ALTER] MATERIALIZED TABLE [catalog_name.][db_name.]table_name

[(
    { <column_definition> | <computed_column_definition> | <watermark_definition> | <table_constraint> }[ , ...n]
    [ <primary_key> ]
    [ <metadata_column_definition> ][ , ...n]
)]

[COMMENT table_comment]

[ <distribution_definition> ]

[PARTITIONED BY (partition_column_name1, partition_column_name2, ...)]

[WITH (key1=val1, key2=val2, ...)]

[FRESHNESS = INTERVAL '' { SECOND[S] | MINUTE[S] | HOUR[S] | DAY[S] }]

[REFRESH_MODE = { CONTINUOUS | FULL }]

AS <select_statement>

<column_definition>:
  column_name column_type [ <column_constraint> ] [COMMENT column_comment]

<column_constraint>:
  [CONSTRAINT constraint_name] PRIMARY KEY NOT ENFORCED

<table_constraint>:
  [CONSTRAINT constraint_name] PRIMARY KEY (column_name, ...) NOT ENFORCED

<metadata_column_definition>:
  column_name column_type METADATA [ FROM metadata_key ] [ VIRTUAL ]

<computed_column_definition>:
  column_name AS computed_column_expression [COMMENT column_comment]

<watermark_definition>:
  WATERMARK FOR rowtime_column_name AS watermark_strategy_expression

<distribution_definition>:
{
    DISTRIBUTED BY [ { HASH | RANGE } ] (bucket_column_name1, bucket_column_name2, ...) [INTO n BUCKETS]
  | DISTRIBUTED INTO n BUCKETS
}

PRIMARY KEY

PRIMARY KEY는 테이블 내의 각 행을 고유하게 식별하는 선택적 컬럼 목록을 정의해요. 기본 키인 컬럼은 null이 아니어야 해요.

PARTITIONED BY

PARTITIONED BY는 머티어리얼라이즈드 테이블을 파티셔닝할 선택적 컬럼 목록을 정의해요. 이 머티어리얼라이즈드 테이블이 파일시스템 싱크로 사용되면 각 파티션에 대해 디렉터리가 생성돼요.

예시:

-- Create a materialized table and specify the partition field as `ds`.
CREATE MATERIALIZED TABLE my_materialized_table
    PARTITIONED BY (ds)
    FRESHNESS = INTERVAL '1' HOUR
    AS SELECT
        ds
    FROM
        ...

참고

  • 파티션 컬럼은 머티어리얼라이즈드 테이블의 쿼리 문에 포함되어야 해요.

WITH 옵션 (WITH Options)

WITH Options커넥터 옵션과 파티션 필드의 시간 포맷 옵션을 포함한 머티어리얼라이즈드 테이블 속성을 지정하는 데 사용돼요.

-- Create a materialized table, specify the partition field as 'ds', and the corresponding time format as 'yyyy-MM-dd'
CREATE MATERIALIZED TABLE my_materialized_table
    PARTITIONED BY (ds)
    WITH (
        'format' = 'json',
        'partition.fields.ds.date-formatter' = 'yyyy-MM-dd'
    )
    ...

위 예시에서 ds 파티션 컬럼에 대해 date-formatter 옵션을 지정했어요. 각 스케줄링 동안 스케줄링 시간은 ds 파티션 값으로 변환돼요. 예를 들어 스케줄링 시간이 2024-01-01 00:00:00이면, ds = '2024-01-01' 파티션만 리프레시돼요.

참고

FRESHNESS

FRESHNESS는 머티어리얼라이즈드 테이블의 데이터 신선도(freshness)를 정의해요.

FRESHNESS는 선택적이에요. 생략하면 시스템은 리프레시 모드에 따라 기본 신선도를 사용해요: CONTINUOUS 모드의 경우 materialized-table.default-freshness.continuous(기본값: 3분), FULL 모드의 경우 materialized-table.default-freshness.full(기본값: 1시간).

FRESHNESS와 Refresh Mode의 관계

FRESHNESS는 머티어리얼라이즈드 테이블의 내용이 기본 테이블의 업데이트보다 지연되어야 하는 최대 시간을 정의해요. 지정하지 않으면 리프레시 모드에 따라 구성에서 기본값을 사용해요. 두 가지를 수행해요: 첫째 구성을 통해 머티어리얼라이즈드 테이블의 리프레시 모드를 결정하고, 둘째 실제 데이터 신선도 요구 사항을 충족하기 위한 데이터 리프레시 빈도를 결정해요.

FRESHNESS 파라미터 설명

FRESHNESS 파라미터 범위는 INTERVAL '' { SECOND | MINUTE | HOUR | DAY }예요. ''은 양의 정수여야 하며, FULL 모드에서는 ''이 해당 시간 간격 단위의 공약수여야 해요.

예시: (materialized-table.refresh-mode.freshness-threshold가 30분이라고 가정)

-- The corresponding refresh pipeline is a streaming job with a checkpoint interval of 1 second
FRESHNESS = INTERVAL '1' SECOND

-- The corresponding refresh pipeline is a real-time job with a checkpoint interval of 1 minute
FRESHNESS = INTERVAL '1' MINUTE

-- The corresponding refresh pipeline is a scheduled workflow with a schedule cycle of 1 hour
FRESHNESS = INTERVAL '1' HOUR

-- The corresponding refresh pipeline is a scheduled workflow with a schedule cycle of 1 day
FRESHNESS = INTERVAL '1' DAY

기본 FRESHNESS 예시: (materialized-table.default-freshness.continuous가 3분, materialized-table.default-freshness.full이 1시간, materialized-table.refresh-mode.freshness-threshold가 30분이라고 가정)

-- FRESHNESS is omitted, uses the configured default of 3 minutes for CONTINUOUS mode
-- The corresponding refresh pipeline is a streaming job with a checkpoint interval of 3 minutes
CREATE MATERIALIZED TABLE my_materialized_table
    AS SELECT * FROM source_table;

-- FRESHNESS is omitted and FULL mode is explicitly specified, uses the configured default of 1 hour
-- The corresponding refresh pipeline is a scheduled workflow with a schedule cycle of 1 hour
CREATE MATERIALIZED TABLE my_materialized_table_full
    REFRESH_MODE = FULL
    AS SELECT * FROM source_table;

잘못된 FRESHNESS 예시:

-- Interval is a negative number
FRESHNESS = INTERVAL '-1' SECOND

-- Interval is 0
FRESHNESS = INTERVAL '0' SECOND

-- Interval is in months or years
FRESHNESS = INTERVAL '1' MONTH
FRESHNESS = INTERVAL '1' YEAR

-- In FULL mode, the interval is not a common divisor of the respective time range
FRESHNESS = INTERVAL '60' SECOND
FRESHNESS = INTERVAL '5' HOUR

참고

  • FRESHNESS가 지정되지 않으면 테이블은 리프레시 모드에 따라 기본 신선도 간격을 사용해요: CONTINUOUS 모드의 경우 materialized-table.default-freshness.continuous(기본값: 3분), FULL 모드의 경우 materialized-table.default-freshness.full(기본값: 1시간).
  • 머티어리얼라이즈드 테이블 데이터는 정의된 신선도 내에서 최대한 가깝게 리프레시되지만 완전한 충족을 보장할 수는 없어요.
  • CONTINUOUS 모드에서 너무 짧은 데이터 신선도 간격을 설정하면 체크포인트 간격과 일치하므로 작업 성능에 영향을 줄 수 있어요. 체크포인트 성능을 최적화하려면 enabling-changelog를 고려해 보세요.
  • FULL 모드에서는 데이터 신선도를 cron 표현식으로 번역해야 하므로, 현재 사전 정의된 시간 범위 내의 신선도 간격만 지원돼요. 이 설계는 cron의 능력과의 정렬을 보장해요. 구체적으로 다음 신선도를 지원해요:
    • 초(Second): 1, 2, 3, 4, 5, 6, 10, 12, 15, 20, 30.
    • 분(Minute): 1, 2, 3, 4, 5, 6, 10, 12, 15, 20, 30.
    • 시(Hour): 1, 2, 3, 4, 6, 8, 12.
    • 일(Day): 1.

REFRESH_MODE

REFRESH_MODE는 머티어리얼라이즈드 테이블의 리프레시 모드를 명시적으로 지정하는 데 사용돼요. 지정된 모드는 프레임워크의 자동 추론보다 우선하며, 특정 시나리오의 필요를 충족해요.

예시: (materialized-table.refresh-mode.freshness-threshold가 30분이라고 가정)

-- The refresh mode of the created materialized table is CONTINUOUS, and the job's checkpoint interval is 1 hour.
CREATE MATERIALIZED TABLE my_materialized_table
    FRESHNESS = INTERVAL '1' HOUR
    REFRESH_MODE = CONTINUOUS
    AS SELECT
       ...

-- The refresh mode of the created materialized table is FULL, and the job's schedule cycle is 10 minutes.
CREATE MATERIALIZED TABLE my_materialized_table
    FRESHNESS = INTERVAL '10' MINUTE
    REFRESH_MODE = FULL
    AS SELECT
       ...

AS

이 절은 머티어리얼라이즈드 테이블 데이터를 채우기 위한 쿼리를 정의하는 데 사용돼요. 업스트림 테이블은 머티어리얼라이즈드 테이블, 테이블 또는 뷰가 될 수 있어요. select 문은 모든 Flink SQL Queries를 지원해요.

예시:

CREATE MATERIALIZED TABLE my_materialized_table
    FRESHNESS = INTERVAL '10' SECOND
    AS SELECT * FROM kafka_catalog.db1.kafka_table;

OR ALTER

OR ALTER 절은 생성-또는-갱신(create-or-update) 의미론을 제공해요:

  • 머티어리얼라이즈드 테이블이 없으면: 지정된 옵션으로 새 머티어리얼라이즈드 테이블을 생성해요.
  • 머티어리얼라이즈드 테이블이 있으면: 쿼리 정의를 수정해요(ALTER MATERIALIZED TABLE AS처럼 동작).

이것은 머티어리얼라이즈드 테이블이 이미 존재하는지 확인하지 않고 원하는 상태를 정의하려는 선언적 배포(declarative deployment) 시나리오에서 특히 유용해요.

머티어리얼라이즈드 테이블이 존재할 때의 동작:

이 연산은 ALTER MATERIALIZED TABLE AS와 유사하게 머티어리얼라이즈드 테이블을 갱신해요:

Full 모드:

  1. 스키마와 쿼리 정의를 갱신해요.
  2. 다음 리프레시 작업이 트리거될 때 새 쿼리를 사용해 머티어리얼라이즈드 테이블이 리프레시돼요.

Continuous 모드:

  1. 현재 실행 중인 리프레시 작업을 일시 중지해요.
  2. 스키마와 쿼리 정의를 갱신해요.
  3. 처음부터 새 리프레시 작업을 시작해요.

자세한 내용은 ALTER MATERIALIZED TABLE AS를 참고하세요.

예시 (Examples)

materialized-table.refresh-mode.freshness-threshold가 30분이라고 가정해요.

10초의 데이터 신선도와 파생된 리프레시 모드 CONTINUOUS를 가진 머티어리얼라이즈드 테이블을 생성:

CREATE MATERIALIZED TABLE my_materialized_table_continuous
    PARTITIONED BY (ds)
    WITH (
        'format' = 'debezium-json',
        'partition.fields.ds.date-formatter' = 'yyyy-MM-dd'
    )
    FRESHNESS = INTERVAL '10' SECOND
    AS SELECT
        k.ds,
        k.user_id,
        COUNT(*) AS event_count,
        SUM(k.amount) AS total_amount,
        MAX(u.age) AS max_age
    FROM
        kafka_catalog.db1.kafka_table k
    JOIN
        user_catalog.db1.user_table u
    ON
        k.user_id = u.user_id
    WHERE
        k.event_type = 'purchase'
    GROUP BY
        k.ds, k.user_id

1시간의 데이터 신선도와 파생된 리프레시 모드 FULL을 가진 머티어리얼라이즈드 테이블을 생성:

CREATE MATERIALIZED TABLE my_materialized_table_full
    PARTITIONED BY (ds)
    WITH (
        'format' = 'json',
        'partition.fields.ds.date-formatter' = 'yyyy-MM-dd'
    )
    FRESHNESS = INTERVAL '1' HOUR
    AS SELECT
        p.ds,
        p.product_id,
        p.product_name,
        AVG(s.sale_price) AS avg_sale_price,
        SUM(s.quantity) AS total_quantity
    FROM
        paimon_catalog.db1.product_table p
    LEFT JOIN
        paimon_catalog.db1.sales_table s
    ON
        p.product_id = s.product_id
    WHERE
        p.category = 'electronics'
    GROUP BY
        p.ds, p.product_id, p.product_name

그리고 컬럼을 명시적으로 지정한 동일한 머티어리얼라이즈드 테이블:

CREATE MATERIALIZED TABLE my_materialized_table_full (
    ds, product_id, product_name, avg_sale_price, total_quantity)
    ...

컬럼의 순서는 쿼리에서와 같을 필요가 없어요. Flink는 필요한 경우 재정렬을 수행해요. 즉 이것도 유효해요.

CREATE MATERIALIZED TABLE my_materialized_table_full (
    product_id, product_name, ds, avg_sale_price, total_quantity)
    ...

다른 방법은 이름과 데이터 타입을 넣는 것이에요.

CREATE MATERIALIZED TABLE my_materialized_table_full (
    ds STRING, product_id STRING, product_name STRING, avg_sale_price DOUBLE, total_quantity BIGINT)
    ...

컬럼의 타입이 같지 않을 수 있으므로, 그 경우 암시적 캐스트(implicit cast)가 적용돼요. 일부 조합에 대해 암시적 캐스트가 지원되지 않으면 검증 오류가 던져져요. 또한 여기서도 재정렬을 할 수 있다는 점에 주목할 가치가 있어요.

두 번 실행되는 머티어리얼라이즈드 테이블 생성 또는 변경:

-- First execution: creates the materialized table
CREATE OR ALTER MATERIALIZED TABLE my_materialized_table
    FRESHNESS = INTERVAL '10' SECOND
    AS
    SELECT
        user_id,
        COUNT(*) AS event_count,
        SUM(amount) AS total_amount
    FROM
        kafka_catalog.db1.events
    WHERE
        event_type = 'purchase'
    GROUP BY
        user_id;

-- Second execution: alters the query definition (adds avg_amount column)
CREATE OR ALTER MATERIALIZED TABLE my_materialized_table
    FRESHNESS = INTERVAL '10' SECOND
    AS
    SELECT
        user_id,
        COUNT(*) AS event_count,
        SUM(amount) AS total_amount,
        AVG(amount) AS avg_amount  -- Add a new nullable column at the end
    FROM
        kafka_catalog.db1.events
    WHERE
        event_type = 'purchase'
    GROUP BY
        user_id;

참고

  • 기존 머티어리얼라이즈드 테이블을 변경할 때, 스키마 진화(schema evolution)는 현재 원래 머티어리얼라이즈드 테이블 스키마의 끝에 nullable 컬럼을 추가하는 것만 지원해요.
  • continuous 모드에서 새 리프레시 작업은 변경 시 원래 리프레시 작업의 상태에서 복원되지 않아요.
  • CREATE와 ALTER 연산의 모든 제한 사항이 적용돼요.

제한 사항 (Limitations)

  • 쿼리에서 사용되지 않는 물리적 컬럼을 명시적으로 지정하는 것은 지원하지 않아요.
  • select 쿼리에서 임시 테이블(temporary tables), 임시 뷰(temporary views), 임시 함수(temporary functions)를 참조하는 것은 지원하지 않아요.

ALTER MATERIALIZED TABLE

ALTER MATERIALIZED TABLE [catalog_name.][db_name.]table_name
    ADD { <column_definition> | ( <column_definition_list> ) | <table_constraint> }
    | MODIFY { <column_definition> | ( <column_definition_list> ) | <table_constraint> }
    | DROP {column_name | (column_name, column_name, ...) | PRIMARY KEY | CONSTRAINT constraint_name | WATERMARK | DISTRIBUTION }
    | SUSPEND | RESUME [WITH (key1=val1, key2=val2, ...)]
    | REFRESH [PARTITION partition_spec] |
    | AS <select_statement>

<column_position>:
  column_name  [FIRST | AFTER column_name]

<table_constraint>:
  [CONSTRAINT constraint_name] PRIMARY KEY (column_name, ...) NOT ENFORCED

<watermark_definition>:
  WATERMARK FOR rowtime_column_name AS watermark_strategy_expression

<column_definition>:
  { <physical_column_definition> | <metadata_column_definition> | <computed_column_definition> } [COMMENT column_comment]

<physical_column_definition>:
  column_type

<metadata_column_definition>:
  column_type METADATA [ FROM metadata_key ] [ VIRTUAL ]

<computed_column_definition>:
  AS computed_column_expression

<distribution_definition>:
{
    DISTRIBUTED BY [ { HASH | RANGE } ] (bucket_column_name1, bucket_column_name2, ...) [INTO n BUCKETS]
  | DISTRIBUTED INTO n BUCKETS
}

ALTER MATERIALIZED TABLE은 머티어리얼라이즈드 테이블을 관리하는 데 사용돼요. 이 명령은 사용자가 머티어리얼라이즈드 테이블의 리프레시 파이프라인을 일시 중지/재개하고 데이터 리프레시를 수동으로 트리거하며, 머티어리얼라이즈드 테이블의 쿼리 정의를 수정할 수 있게 해줘요.

ADD

ADD 절을 사용해 기존 머티어리얼라이즈드 테이블에 컬럼(computed와 metadata virtual 같은 비영속(non persisted) 컬럼만), 제약 조건, 워터마크, 분포를 추가해요.

지정된 위치에 컬럼을 추가하려면 FIRST 또는 AFTER col_name을 사용해요. 기본적으로 컬럼은 마지막에 추가돼요.

다음 예시들은 ADD 문의 사용법을 보여줘요.

-- add a new column 
ALTER MATERIALIZED TABLE MyMaterializedTable ADD category_id STRING METADATA VIRTUAL;

-- add columns, constraint, and watermark
ALTER MATERIALIZED TABLE MyMaterializedTable ADD (
    log_ts STRING METADATA VIRTUAL FIRST,
    ts AS TO_TIMESTAMP(log_ts) AFTER log_ts,
    PRIMARY KEY (id) NOT ENFORCED,
    WATERMARK FOR ts AS ts - INTERVAL '3' SECOND
);

-- add new distribution using a hash on uid into 4 buckets
ALTER MATERIALIZED TABLE MyMaterializedTable ADD DISTRIBUTION BY HASH(uid) INTO 4 BUCKETS;

-- add new distribution on uid into 4 buckets
ALTER MATERIALIZED TABLE MyMaterializedTable ADD DISTRIBUTION BY (uid) INTO 4 BUCKETS;

-- add new distribution on uid.
ALTER MATERIALIZED TABLE MyMaterializedTable ADD DISTRIBUTION BY (uid);

-- add new distribution into 4 buckets
ALTER MATERIALIZED TABLE MyMaterializedTable ADD DISTRIBUTION INTO 4 BUCKETS;

참고: 기본 키가 될 컬럼을 추가하면 해당 컬럼의 null 허용 여부(nullability)가 암시적으로 false로 변경돼요.

MODIFY

MODIFY 절을 사용해 기존 테이블의 컬럼 주석, 위치, 타입(computed와 metadata virtual 같은 비영속 컬럼만, 컬럼 참고), 기본 키 컬럼과 워터마크 전략을 변경해요.

기존 컬럼을 새 위치로 수정하려면 FIRST 또는 AFTER col_name을 사용해요. 기본적으로 위치는 변경되지 않아요.

다음 예시들은 MODIFY 문의 사용법을 보여줘요.

-- modify a column type, comment and position
ALTER MATERIALIZED TABLE MyMaterializedTable MODIFY measurement double METADATA COMMENT 'unit is bytes per second' AFTER `id`;

-- modify definition of column log_ts and ts, primary key, watermark. They must exist in table schema
ALTER MATERIALIZED TABLE MyMaterializedTable MODIFY (
    log_ts STRING METADATA COMMENT 'log timestamp string' AFTER `id`,  -- reorder columns
    ts AS TO_TIMESTAMP(log_ts) AFTER log_ts,
    PRIMARY KEY (id) NOT ENFORCED,
    WATERMARK FOR ts AS ts -- modify watermark strategy
);

참고: 기본 키가 될 컬럼을 수정하면 해당 컬럼의 null 허용 여부가 암시적으로 false로 변경돼요.

DROP

DROP 절을 사용해 기존 테이블의 컬럼(computed와 metadata virtual 같은 비영속 컬럼만, 컬럼 참고), 기본 키, 파티션, 워터마크 전략을 제거해요.

다음 예시들은 DROP 문의 사용법을 보여줘요.

-- drop a column
ALTER MATERIALIZED TABLE MyMaterializedTable DROP measurement;

-- drop columns
ALTER MATERIALIZED TABLE MyMaterializedTable DROP (col1, col2, col3);

-- drop primary key
ALTER MATERIALIZED TABLE MyMaterializedTable DROP PRIMARY KEY;

-- drop a watermark
ALTER MATERIALIZED TABLE MyMaterializedTable DROP WATERMARK;

-- drop distribution
ALTER MATERIALIZED TABLE MyMaterializedTable DROP DISTRIBUTION;

SUSPEND

ALTER MATERIALIZED TABLE [catalog_name.][db_name.]table_name SUSPEND

SUSPEND는 머티어리얼라이즈드 테이블의 백그라운드 리프레시 파이프라인을 일시 중지하는 데 사용돼요.

예시:

-- Specify SAVEPOINT path before pausing
SET 'execution.checkpointing.savepoint-dir' = 'hdfs://savepoint_path';

-- Suspend the specified materialized table
ALTER MATERIALIZED TABLE my_materialized_table SUSPEND;

참고

  • CONTINUOUS 모드의 테이블을 일시 중지할 때, 작업은 기본적으로 STOP WITH SAVEPOINT를 사용해 일시 중지돼요. 파라미터를 사용해 SAVEPOINT 저장 경로를 설정해야 해요.

RESUME

ALTER MATERIALIZED TABLE [catalog_name.][db_name.]table_name RESUME [WITH (key1=val1, key2=val2, ...)]

RESUME은 머티어리얼라이즈드 테이블의 리프레시 파이프라인을 재개하는 데 사용돼요. 머티어리얼라이즈드 테이블 동적 옵션은 WITH options 절을 통해 지정할 수 있으며, 이는 현재 리프레시된 파이프라인에만 적용되고 영속적이지 않아요.

예시:

-- Resume the specified materialized table
ALTER MATERIALIZED TABLE my_materialized_table RESUME;

-- Resume the specified materialized table and specify sink parallelism
ALTER MATERIALIZED TABLE my_materialized_table RESUME WITH ('sink.parallelism'='10');

REFRESH

ALTER MATERIALIZED TABLE [catalog_name.][db_name.]table_name REFRESH [PARTITION partition_spec]

REFRESH는 머티어리얼라이즈드 테이블의 리프레시를 능동적으로 트리거하는 데 사용돼요.

예시:

-- Refresh the entire table data
ALTER MATERIALIZED TABLE my_materialized_table REFRESH;

-- Refresh specified partition data
ALTER MATERIALIZED TABLE my_materialized_table REFRESH PARTITION (ds='2024-06-28');

참고

  • REFRESH 연산은 머티어리얼라이즈드 테이블 데이터를 리프레시하기 위해 Flink 배치 작업을 시작해요.

AS

ALTER MATERIALIZED TABLE [catalog_name.][db_name.]table_name AS <select_statement>

AS 절은 머티어리얼라이즈드 테이블을 리프레시하기 위한 쿼리 정의를 수정할 수 있게 해줘요. 먼저 새 쿼리에서 파생된 스키마를 사용해 테이블의 스키마를 진화시킨 다음, 새 쿼리를 사용해 테이블 데이터를 리프레시해요. 기본적으로 이 작업은 과거 데이터(historical data)에 영향을 주지 않는다는 점을 강조하는 것이 중요해요.

수정 과정은 머티어리얼라이즈드 테이블의 리프레시 모드에 따라 달라져요:

Full 모드:

  1. 머티어리얼라이즈드 테이블의 schemaquery definition을 갱신해요.
  2. 다음 리프레시 작업이 트리거될 때 새 쿼리 정의를 사용해 테이블이 리프레시돼요:
    • 파티션 테이블이고 partition.fields.#.date-formatter가 올바르게 설정되면, 최신 파티션만 리프레시돼요.
    • 그렇지 않으면 테이블 전체가 덮어써져요.

Continuous 모드:

  1. 현재 실행 중인 리프레시 작업을 일시 중지해요.
  2. 머티어리얼라이즈드 테이블의 schemaquery definition을 갱신해요.
  3. 머티어리얼라이즈드 테이블을 리프레시하기 위해 새 리프레시 작업을 시작해요:
    • 새 리프레시 작업은 처음부터 시작하며 이전 상태에서 복원하지 않아요.
    • 데이터 소스의 시작 오프셋은 커넥터의 기본 구현 또는 쿼리에 지정된 dynamic hint에 의해 결정돼요.

예시:

-- Definition of origin materialized table
CREATE MATERIALIZED TABLE my_materialized_table
    FRESHNESS = INTERVAL '10' SECOND
    AS
    SELECT
        user_id,
        COUNT(*) AS event_count,
        SUM(amount) AS total_amount
    FROM
        kafka_catalog.db1.events
    WHERE
        event_type = 'purchase'
    GROUP BY
        user_id;

-- Modify the query definition of materialized table
ALTER MATERIALIZED TABLE my_materialized_table
    AS
    SELECT
        user_id,
        COUNT(*) AS event_count,
        SUM(amount) AS total_amount,
        AVG(amount) AS avg_amount  -- Add a new nullable column at the end
    FROM
        kafka_catalog.db1.events
    WHERE
        event_type = 'purchase'
    GROUP BY
        user_id;

참고

  • 스키마 진화는 현재 원래 테이블 스키마의 끝에 nullable 컬럼을 추가하는 것만 지원해요.
  • continuous 모드에서 새 리프레시 작업은 원래 리프레시 작업의 상태에서 복원되지 않아요. 이로 인해 임시 데이터 중복 또는 손실이 발생할 수 있어요.

DROP MATERIALIZED TABLE

DROP MATERIALIZED TABLE [IF EXISTS] [catalog_name.][database_name.]table_name

머티어리얼라이즈드 테이블을 삭제할 때, 백그라운드 리프레시 파이프라인이 먼저 삭제되고, 그런 다음 머티어리얼라이즈드 테이블에 해당하는 메타데이터가 Catalog에서 제거돼요.

예시:

-- Delete the specified materialized table
DROP MATERIALIZED TABLE IF EXISTS my_materialized_table;

더 알아보기 (Learn more)