스트림과 태스크에서 동적 테이블로 마이그레이션
스트림과 태스크에서 동적 테이블로 마이그레이션
동적 테이블은 명령형 파이프라인 코드(INSERT, MERGE, 저장 프로시저 호출)를 선언형 모델로 대체해요. 원하는 결과를 설명하는 SELECT를 작성하면 Snowflake가 그것을 신선하게 유지해줘요.
이 페이지는 스트림-and-태스크 파이프라인을 한 태스크씩 동적 테이블로 변환하는 방법을 다뤄요.
출처: Snowflake 문서
본문
개념 지도: 스트림과 태스크에서 동적 테이블로
| 스트림과 태스크 개념 | 동적 테이블 동등물 | 참고 |
|---|---|---|
| 기본 테이블 | 기본 테이블 | 같은 객체, 변경 불필요. |
| 기본 테이블의 스트림 | (내장) | Snowflake가 변경을 자동으로 추적. 별도 스트림을 만들지 않아도 돼요. |
태스크 본문 (INSERT INTO ... SELECT ...) |
CREATE DYNAMIC TABLE … AS SELECT의 정의 |
SELECT만 작성. INSERT INTO나 MERGE 같은 DML 문 없음. |
태스크 스케줄 (SCHEDULE = '10 MINUTE') |
TARGET_LAG = '10 minutes' |
TARGET_LAG는 고정 간격이 아니라 신선도 목표예요. Snowflake가 언제 갱신할지 결정해요. |
| 태스크 DAG (루트 태스크 + 자식 태스크) | 동적 테이블 파이프라인 | 중간 동적 테이블은 TARGET_LAG = DOWNSTREAM을 사용. 종단 동적 테이블이 신선도 목표를 가져요. |
| 루트 태스크 | 종단 동적 테이블 (또는 컨트롤러 동적 테이블) | 종단 동적 테이블이 전체 파이프라인의 갱신 스케줄을 주도해요. |
MERGE ON key WHEN MATCHED UPDATE … |
QUALIFY ROW_NUMBER() OVER (...) = 1 |
동적 테이블은 "키당 최신 상태"를 DML이 아니라 윈도우 함수로 표현해요. SCD Type 1 예제 참조. |
SYSTEM$STREAM_HAS_DATA() 검사 |
(불필요) | 기본 테이블이 변경되지 않으면 Snowflake가 갱신을 자동으로 건너뛰어요. |
마이그레이션 전략 선택
동적 테이블은 두 가지 운영 모델을 지원해요. 팀 워크플로우에 맞는 것을 선택하세요.
자가 갱신 동적 테이블 (Snowflake가 타이밍 관리)
TARGET_LAG 값을 설정하면 Snowflake가 갱신 스케줄을 처리해요. 가장 단순한 모델이며 새 파이프라인이나 운영 오버헤드를 줄이려는 팀에 적합해요.
CREATE OR REPLACE DYNAMIC TABLE dt_orders
TARGET_LAG = '10 minutes'
WAREHOUSE = transform_wh
AS
SELECT ...
FROM raw_orders;
다음에 가장 좋아요:
- 아직 오케스트레이션을 갖추지 않은 새 파이프라인.
- Snowflake가 갱신 타이밍과 의존성 순서를 관리하도록 하려는 팀.
- 신선도 목표(예: "10분 이상 지연되지 않아야 함")가 올바른 요구 사항인 파이프라인.
오케스트레이터 관리 동적 테이블 (SCHEDULER = DISABLE)
이미 오케스트레이터(로 dbt, Airflow, 또는 태스크)를 실행 중이라면 SCHEDULER=DISABLE로 동적 테이블을 만드세요. Snowflake가 갱신을 자동으로 예약하지 않아요. 각 갱신을 수동으로 또는 오케스트레이터에서 트리거하세요. 갱신은 업스트림이나 다운스트림으로 계단식(cascade)되지 않아요.
-- 스케줄링을 비활성화한 동적 테이블 생성.
CREATE OR REPLACE DYNAMIC TABLE dt_orders
WAREHOUSE = transform_wh
SCHEDULER = DISABLE
AS
SELECT ...
FROM raw_orders;
-- 수동으로 또는 오케스트레이터에서 갱신 트리거.
ALTER DYNAMIC TABLE dt_orders REFRESH;
주요 동작:
SCHEDULER = DISABLE를 사용할 때는TARGET_LAG를 설정할 수 없어요.SCHEDULER = DISABLE를 사용한 갱신은 업스트림이나 다운스트림으로 계단식되지 않아요. 각 동적 테이블이 독립적으로 갱신돼요.- 오케스트레이터가 순서, 타이밍, 오류 처리를 제어해요.
- 2026년 3월부터 GA.
- 다음에 가장 좋아요:
- 선언형 정의를 채택하면서 기존 오케스트레이터를 유지하려는 팀.
- 갱신 순서나 조건부 갱신 로직에 대한 명시적 제어가 필요한 파이프라인.
태스크를 마이그레이션할지 결정
각 태스크를 다음 기준에 개별적으로 평가해 마이그레이션할지 결정하세요.
| 기준 | 동적 테이블로 마이그레이션 | 스트림과 태스크 유지 |
|---|---|---|
| 파이프라인 로직 | 순수 SQL: SELECT, JOIN, GROUP BY, CASE WHEN, 윈도우 함수, 커스텀 UDF |
루프, 커서, 다중 문 프로시저 블록 |
| DML 연산 | 기본 테이블에서 삽입·업데이트·삭제가 증분 갱신을 통해 전파됨. | 명시적 DELETE, TRUNCATE, 또는 삭제-후-삽입 패턴 |
| 데이터 소스 | Snowflake 테이블, 뷰 | 외부 테이블, 디렉터리 테이블(동적 테이블 입력으로 지원되지 않음) |
| 지연 요구 사항 | 최소 1분(TARGET_LAG 하한). 실제 신선도는 갱신 기간과 파이프라인 깊이에 따라 달라짐. |
현 1분 이하 오케스트레이션(태스크는 1초 스케줄 지원) |
비결정적 함수(RANDOM) |
REFRESH_MODE = AUTO로 전체 갱신을 강제. 명시적 REFRESH_MODE = INCREMENTAL이면 실패. |
기본 지원 |
| 부작용 또는 외부 API 호출 | 불가능. 동적 테이블은 순수 SELECT 문. |
지원: 외부 함수, 알림, 저장 프로시저 호출 |
| 한 트랜잭션의 다중 테이블 쓰기 | 지원되지 않음. 각 동적 테이블은 정확히 하나의 출력을 생성. | 한 태스크가 단일 트랜잭션에서 여러 테이블에 쓸 수 있음. |
| 다중 문 트랜잭션 | 동적 테이블당 하나의 정의(단일 SELECT) |
태스크가 임의의 다중 문 트랜잭션을 실행할 수 있음. |
| SQL로 표현 가능한 비즈니스 규칙 | JOIN, CASE WHEN, UNION, 조건부 집계가 모두 정의에서 작동. |
런타임 평가 절차적 로직(커서, 동적 SQL 디스패치) |
🔴
RANDOM()같은 비결정적 함수는REFRESH_MODE=AUTO일 때 전체 갱신을 강제해요.REFRESH_MODE=INCREMENTAL을 명시적으로 설정하면CREATE문이 실패해요. 마이그레이션 전에 비결정적 함수를 결정적 대안으로 교체하세요. 시퀀스 함수(NEXTVAL)는 증분 갱신 모드에서 지원되지 않아요. 전체 호환성 매트릭스는 동적 테이블 지원 쿼리 문서를 참조하세요.
커스텀 증분 동적 테이블
스트림+태스크 파이프라인이 MERGE 또는 INSERT DML을 사용한다면, 커스텀 증분 동적 테이블이 더 적은 오케스트레이션 오버헤드로 같은 패턴을 받아들여요. 매핑은 거의 일대일이에요:
| 스트림 + 태스크 | 커스텀 증분 동적 테이블 |
|---|---|
| 기본 테이블의 스트림 | CHANGES() 절 |
SYSTEM$STREAM_HAS_DATA() 검사 |
자동 (Snowflake가 처리) |
| 태스크 스케줄 / CRON | TARGET_LAG |
MERGE INTO target |
MERGE INTO SELF |
INSERT INTO target |
INSERT INTO SELF |
스트림 메타데이터 컬럼 (METADATA$ACTION, METADATA$ISUPDATE) |
CHANGES()의 같은 컬럼 |
| 수동 스트림 진행 | 동적 테이블의 자동 진행 |
| 태스크 DAG (선행자) | 파이프라인 (TARGET_LAG = DOWNSTREAM) |
커스텀 증분 동적 테이블은 스트림, 태스크, 그리고 그 오케스트레이션을 별도로 관리할 필요를 없애요. REFRESH USING 안에 같은 MERGE 또는 INSERT 로직을 작성하면 Snowflake가 스케줄링, 재시도, 트랜잭션 보장을 처리해요.
전체 세부 사항은 커스텀 증분화 문서를 참조하세요.
잘 마이그레이션되는 패턴
다음 패턴은 스트림과 태스크에서 동적 테이블로 깔끔하게 변환돼요.
추가 전용(append-only) 수집
이것은 가장 단순한 마이그레이션이에요. 대상 테이블에 새 행을 삽입하는 스트림과 태스크가 단일 동적 테이블이 돼요.
먼저 기본 테이블과 샘플 데이터를 설정하세요. 이 예시는 첫 동적 테이블 만들기에서 사용한 것과 같은 raw_orders 테이블을 사용해요. 이미 그 튜토리얼을 실행했다면 설정 단계를 건너뛰어도 돼요.
Before: stream and task
CREATE OR REPLACE STREAM raw_orders_stream ON TABLE raw_orders
APPEND_ONLY = TRUE;
CREATE OR REPLACE TABLE clean_orders (
order_id INT,
customer_id INT,
order_date TIMESTAMP_NTZ,
product_name VARCHAR,
quantity INT,
unit_price DECIMAL(10,2),
line_total DECIMAL(12,2),
order_status VARCHAR
);
CREATE OR REPLACE TASK load_clean_orders
WAREHOUSE = transform_wh
SCHEDULE = '10 MINUTE'
WHEN SYSTEM$STREAM_HAS_DATA('raw_orders_stream')
AS
INSERT INTO clean_orders
SELECT
order_id,
customer_id,
order_date,
TRIM(UPPER(product_name)) AS product_name,
quantity,
unit_price,
quantity * unit_price AS line_total,
order_status
FROM raw_orders_stream
WHERE order_status != 'returned';
ALTER TASK load_clean_orders RESUME;
After: dynamic table
CREATE OR REPLACE DYNAMIC TABLE dt_clean_orders
TARGET_LAG = '10 minutes'
WAREHOUSE = transform_wh
REFRESH_MODE = INCREMENTAL
AS
SELECT
order_id,
customer_id,
order_date,
TRIM(UPPER(product_name)) AS product_name,
quantity,
unit_price,
quantity * unit_price AS line_total,
order_status
-- remaining columns omitted for brevity
FROM raw_orders
WHERE order_status != 'returned';
스트림, 대상 테이블 DDL, 태스크, 그리고 INSERT가 모두 하나의 CREATE DYNAMIC TABLE 문으로 합쳐져요. SELECT가 기본 테이블에서 직접 읽어요. 스트림이 필요 없어요.
🔴 "Before" 버전의
APPEND_ONLY스트림은 새로 삽입된 행만 캡처했어요. 동적 테이블은WHERE필터와 일치하는 기본 테이블의 현재 상태를 정확히 나타내요. 기본 테이블이 INSERT만 받으면(원시 랜딩 테이블에서 흔함) 결과는 동등해요. 기본 테이블이UPDATE나DELETE도 받으면 동적 테이블도 그 변경을 반영해요. 초기 갱신이 완료된 후 동적 테이블 내용을 검증하세요.
📌 스트림-정적 조인(태스크가 스트림을 차원 테이블과 조인하는 것)은 SELECT 기반 동적 테이블에서 지원되지 않지만, 커스텀 증분 동적 테이블은 이 패턴을 기본적으로 지원해요. 커스텀 증분화 문서 참조.
SELECT * FROM dt_clean_orders ORDER BY order_id;
+----------+-------------+---------------------+--------------+----------+------------+------------+--------------+
| ORDER_ID | CUSTOMER_ID | ORDER_DATE | PRODUCT_NAME | QUANTITY | UNIT_PRICE | LINE_TOTAL | ORDER_STATUS |
+----------+-------------+---------------------+--------------+----------+------------+------------+--------------+
| 1001 | 1 | 2025-01-15 08:30:00 | WIDGET A | 3 | 29.99 | 89.97 | completed |
| 1002 | 2 | 2025-01-15 09:45:00 | WIDGET B | 1 | 49.99 | 49.99 | completed |
| 1003 | 1 | 2025-01-15 14:20:00 | WIDGET A | 2 | 29.99 | 59.98 | pending |
| 1004 | 3 | 2025-01-16 10:00:00 | GADGET X | 5 | 12.50 | 62.50 | completed |
+----------+-------------+---------------------+--------------+----------+------------+------------+--------------+
SCD Type 1 upsert
스트림과 태스크 파이프라인에서 SCD Type 1(각 레코드의 최신 버전만 유지)은 일반적으로 MERGE 문을 사용해요. 동적 테이블에서는 키당 가장 최근 행을 고르는 윈도우 함수로 같은 로직을 표현해요.
Before: stream and task with MERGE
CREATE OR REPLACE TABLE customer_updates (
customer_id INT,
customer_name VARCHAR,
region VARCHAR,
segment VARCHAR,
updated_at TIMESTAMP_NTZ
);
INSERT INTO customer_updates VALUES
(1, 'Acme Corp', 'US-West', 'Enterprise', '2025-01-10 09:00:00'),
(2, 'Globex Inc', 'US-East', 'Mid-Market', '2025-01-10 09:00:00'),
(1, 'Acme Corp', 'US-West', 'Strategic', '2025-01-15 14:00:00'),
(3, 'Initech LLC', 'EU-West', 'Startup', '2025-01-12 11:00:00');
CREATE OR REPLACE STREAM customer_updates_stream ON TABLE customer_updates;
CREATE OR REPLACE TABLE dim_customers_scd1 (
customer_id INT,
customer_name VARCHAR,
region VARCHAR,
segment VARCHAR,
updated_at TIMESTAMP_NTZ
);
CREATE OR REPLACE TASK merge_customers
WAREHOUSE = transform_wh
SCHEDULE = '10 MINUTE'
WHEN SYSTEM$STREAM_HAS_DATA('customer_updates_stream')
AS
MERGE INTO dim_customers_scd1 AS tgt
USING (
SELECT customer_id, customer_name, region, segment, updated_at
FROM customer_updates_stream
QUALIFY ROW_NUMBER() OVER (
PARTITION BY customer_id ORDER BY updated_at DESC
) = 1
) AS src
ON tgt.customer_id = src.customer_id
WHEN MATCHED THEN UPDATE SET
tgt.customer_name = src.customer_name,
tgt.region = src.region,
tgt.segment = src.segment,
tgt.updated_at = src.updated_at
WHEN NOT MATCHED THEN INSERT
(customer_id, customer_name, region, segment, updated_at)
VALUES
(src.customer_id, src.customer_name, src.region, src.segment, src.updated_at);
ALTER TASK merge_customers RESUME;
After: dynamic table
CREATE OR REPLACE DYNAMIC TABLE dt_dim_customers_scd1
TARGET_LAG = '10 minutes'
WAREHOUSE = transform_wh
REFRESH_MODE = INCREMENTAL
AS
SELECT
customer_id,
customer_name,
region,
segment,
updated_at
FROM customer_updates
QUALIFY ROW_NUMBER() OVER (
PARTITION BY customer_id ORDER BY updated_at DESC
) = 1;
MERGE 로직이 완전히 사라져요. 동적 테이블은 항상 각 고객의 최신 버전을 만들어내요. customer_updates에 더 새로운 updated_at을 가진 새 행이 도착하면 다음 갱신이 자동으로 그것을 잡아요.
SELECT * FROM dt_dim_customers_scd1 ORDER BY customer_id;
+-------------+---------------+---------+------------+---------------------+
| CUSTOMER_ID | CUSTOMER_NAME | REGION | SEGMENT | UPDATED_AT |
+-------------+---------------+---------+------------+---------------------+
| 1 | Acme Corp | US-West | Strategic | 2025-01-15 14:00:00 |
| 2 | Globex Inc | US-East | Mid-Market | 2025-01-10 09:00:00 |
| 3 | Initech LLC | EU-West | Startup | 2025-01-12 11:00:00 |
+-------------+---------------+---------+------------+---------------------+
고객 1은 Enterprise가 아니라 Strategic(더 새로운 세그먼트)을 보여줘요. ROW_NUMBER 윈도우 함수가 이전에 MERGE가 수행했던 "최신 것이 이긴다" 로직을 처리해요.
다단계 파이프라인 (태스크 DAG에서 동적 테이블 파이프라인으로)
중간 동적 테이블은 TARGET_LAG=DOWNSTREAM을 사용해요. 즉 시간 기반 지연이 있는 다운스트림 테이블이 신선한 데이터가 필요할 때만 갱신된다는 뜻이에요. 종단 동적 테이블이 전체 파이프라인의 스케줄을 주도해요.
Before: task DAG
-- Stage 1: raw orders 정리
CREATE OR REPLACE TASK stage1_clean
WAREHOUSE = transform_wh
SCHEDULE = '30 MINUTE'
AS
INSERT INTO dt_orders
SELECT order_id, customer_id, order_date,
TRIM(UPPER(product_name)) AS product_name,
quantity, unit_price,
quantity * unit_price AS line_total,
order_status
FROM raw_orders_stream
WHERE order_status != 'returned';
-- Stage 2: 고객과 조인 (stage1_clean에 의존)
CREATE OR REPLACE TASK stage2_enrich
WAREHOUSE = transform_wh
AFTER stage1_clean
AS
INSERT INTO enriched_orders
SELECT s.order_id, s.order_date, s.product_name, s.line_total,
c.customer_name, c.region, c.segment
FROM dt_orders s
JOIN dim_customers c ON s.customer_id = c.customer_id;
-- Stage 3: 일일 집계 (stage2_enrich에 의존)
CREATE OR REPLACE TASK stage3_aggregate
WAREHOUSE = transform_wh
AFTER stage2_enrich
AS
INSERT INTO dt_orders_daily
SELECT DATE_TRUNC('day', order_date) AS order_day,
region, segment,
COUNT(*) AS order_count,
SUM(line_total) AS daily_revenue
FROM enriched_orders
GROUP BY ALL;
-- 루트 태스크 재개가 전체 DAG를 시작합니다. 자식 태스크가 자동으로 실행됩니다.
ALTER TASK stage1_clean RESUME;
After: dynamic table pipeline
이 예시는 동적 테이블 만들기에서 사용한 dim_customers 테이블을 사용해요. 만들지 않았다면 해당 페이지의 설정 SQL을 참조하세요.
-- Stage 1: raw orders 정리 (스트림이 아닌 기본 테이블에서 읽기)
CREATE OR REPLACE DYNAMIC TABLE dt_orders
TARGET_LAG = DOWNSTREAM
WAREHOUSE = transform_wh
REFRESH_MODE = INCREMENTAL
AS
SELECT
order_id, customer_id, order_date,
TRIM(UPPER(product_name)) AS product_name,
quantity, unit_price,
quantity * unit_price AS line_total,
order_status
FROM raw_orders
WHERE order_status != 'returned';
-- Stage 2: 고객 데이터로 보강
CREATE OR REPLACE DYNAMIC TABLE dt_enriched_orders
TARGET_LAG = DOWNSTREAM
WAREHOUSE = transform_wh
REFRESH_MODE = INCREMENTAL
AS
SELECT
s.order_id, s.order_date, s.product_name, s.line_total,
c.customer_name, c.region, c.segment
FROM dt_orders s
JOIN dim_customers c ON s.customer_id = c.customer_id;
-- Stage 3: 일일 집계 (종단 -- 신선도 목표를 가짐)
CREATE OR REPLACE DYNAMIC TABLE dt_orders_daily
TARGET_LAG = '30 minutes'
WAREHOUSE = transform_wh
REFRESH_MODE = INCREMENTAL
AS
SELECT
DATE_TRUNC('day', order_date) AS order_day,
region, segment,
COUNT(*) AS order_count,
SUM(line_total) AS daily_revenue
FROM dt_enriched_orders
GROUP BY ALL;
AFTER 의존성 순서가 동적 테이블 사이의 암시적 데이터 의존성으로 대체돼요. Snowflake가 파이프라인 그래프를 읽고 dt_orders_daily를 갱신하기 전에 dt_orders와 dt_enriched_orders를 자동으로 갱신해요.
SELECT * FROM dt_orders_daily ORDER BY order_day, region;
+------------+---------+------------+-------------+---------------+
| ORDER_DAY | REGION | SEGMENT | ORDER_COUNT | DAILY_REVENUE |
+------------+---------+------------+-------------+---------------+
| 2025-01-15 | US-East | Mid-Market | 1 | 49.99 |
| 2025-01-15 | US-West | Enterprise | 2 | 149.95 |
| 2025-01-16 | EU-West | Startup | 1 | 62.50 |
+------------+---------+------------+-------------+---------------+
마이그레이션되지 않는 패턴
다음 패턴은 스트림과 태스크가 필요해요. 동적 테이블에 억지로 넣지 마세요.
외부 API 호출과 부작용
태스크가 외부 함수를 호출하거나, 알림을 보내거나, 대상 테이블에 쓰는 것 외의 부작용을 만들면 그 로직은 동적 테이블 안에서 실행될 수 없어요. 동적 테이블은 순수 SELECT 문이에요.
루프와 커서 기반 절차 코드
동적 테이블은 다중 문 절차 블록, 루프, 커서 기반 반복을 포함할 수 없어요. 태스크 본문이 커서를 열고, 결과를 반복하고, 각 행에 조건부로 다른 DML을 적용한다면 그 단계는 스트림과 태스크에 남아 있어야 해요.
SQL로 표현할 수 있는 조건부 로직(CASE WHEN으로 매핑하는 IF/ELSE, COALESCE와 함께 OUTER JOIN으로 매핑하는 분기 JOIN)은 동적 테이블 정의에서 잘 작동해요.
GDPR과 삭제 권리(right-to-erasure)
기본 테이블의 삭제는 다음 갱신에 동적 테이블을 통해 전파돼요. 소스에서 고객 데이터를 삭제하면 동적 테이블이 갱신 후 그 삭제를 반영해요. 이것은 대부분의 규정 준수 시나리오에서 작동해요.
frozen region이 있는 동적 테이블은 frozen region에서의 삭제를 전파하지 않아요. 규정 준수 워크플로우가 frozen 경계 안의 데이터를 대상으로 하면 기본 테이블에서 삭제된 후에도 그 행이 남아 있어요. frozen region에 대한 자세한 내용은 Frozen regions와 backfill 문서를 참조하세요.
규정 준수를 위해 의존하기 전에 작은 데이터 세트로 테스트해 특정 정의에 대한 삭제 전파 동작을 검증하세요. 규정 준수 워크플로우에 안정적인 삭제 전파가 필요하다면 REFRESH_MODE=FULL을 사용하거나 삭제 단계를 스트림과 태스크에 유지하세요.
시퀀스 기반 서로게이트 키
동적 테이블은 증분 갱신 모드에서 SEQUENCE.NEXTVAL을 지원하지 않아요. 전체 갱신 모드에서는 시퀀스가 매 갱신마다 새 값을 생성하므로 서로게이트 키가 안정적이지 않아요.
대안:
- 결정적 해시:
SHA2(CONCAT(order_id, '|', customer_id)) AS surrogate_key - 데이터의 자연 키를 기반으로 한 시스템 파생 고유 식별자.
- 충돌 확률이 허용 가능한 곳에서 행 수준 정체성을 위한
HASH(*). RELY제약 조건이 있는 복합 기본 키.
-- SEQUENCE.NEXTVAL 대신 결정적 해시 사용:
SHA2(CONCAT(order_id, '|', customer_id)) AS surrogate_key
최상의 성능을 위해 서로게이트 키를 파이프라인에서 가능한 한 일찍 구체화하세요. 다운스트림 동적 테이블에서 키를 계산하면 매 갱신마다 불필요한 오버헤드가 추가돼요.
단일 태스크 변환
한 번에 한 태스크씩 변환하려면 다음 단계를 따르세요. 파이프라인에서 가장 단순한 태스크부터 시작해 검증한 후 다음으로 넘어가세요. 이 접근법은 폭발 반경을 제한하고 롤백을 간단하게 만들어요.
- 마이그레이션 전략을 선택하세요. 자가 갱신(
TARGET_LAG설정)과 오케스트레이터 관리(SCHEDULER = DISABLE설정) 중에서 결정하세요. 점진적으로 마이그레이션하면서 지금은 기존 오케스트레이션을 유지하려면SCHEDULER = DISABLE이 스케줄링 방식을 바꾸지 않고 선언형 모델을 채택하게 해줘요. - 갱신 모드를 선택하세요. 지원되는 연산자가 있는 append-heavy 워크로드에는
REFRESH_MODE = INCREMENTAL을 사용하세요. 비결정적 함수나 지원되지 않는 연산자가 있는 정의에는REFRESH_MODE = FULL을 사용하세요. 생성 시점에 Snowflake가 결정하도록 하려면REFRESH_MODE = AUTO를 사용하세요. 자세한 내용은 동적 테이블 갱신 모드 문서를 참조하세요. - 출력을 검증하세요. 동적 테이블의 내용을 이전 대상 테이블과 비교하세요. 갱신 상태를 확인하고 행 수를 비교하세요:
-- 초기 갱신이 성공적으로 완료됐는지 확인.
SELECT name, refresh_mode, scheduling_state
FROM TABLE(INFORMATION_SCHEMA.DYNAMIC_TABLES())
WHERE name = 'DT_CLEAN_ORDERS';
-- 이전 테이블과 새 테이블의 행 수 비교.
SELECT (SELECT COUNT(*) FROM clean_orders_old) AS old_count,
(SELECT COUNT(*) FROM dt_clean_orders) AS new_count;
- 이전 태스크를 일시 중단하세요. 검증 후 이전 태스크를 일시 중단하고 두 파이프라인을 폐기 전에 최소 한 번의 전체 비즈니스 주기 동안 병렬로 실행하세요. 새 동적 테이블에 확신이 들 때까지 이전 대상 테이블이나 스트림을 삭제하지 마세요. 롤백이 필요하면 일시 중단된 태스크를 재개하면 스트림 변경의 백로그가 처리돼요.
마이그레이션 중 동적 테이블 일시 중단
기본 테이블이나 업스트림 객체에 대한 스키마 변경은 전환 중 활성 동적 테이블이 반복적으로 실패하게 할 수 있어요. 시작 전에 영향을 받는 동적 테이블을 일시 중단하고 완료 후 재개하세요.
일반적인 워크플로우는: 잎에서 뿌리 순서로 일시 중단하고, 변경을 적용한 다음, 뿌리에서 잎 순서로 재개하는 것이에요.
전체 SUSPEND/RESUME 워크플로우(의존성 발견 포함)는 동적 테이블 수정 문서를 참조하세요.
하이브리드 파이프라인: 부분 마이그레이션
전체 파이프라인을 한 번에 마이그레이션할 필요는 없어요. 동적 테이블과 스트림-and-태스크가 같은 파이프라인에서 공존할 수 있어요.
예를 들어 3단계 파이프라인의 처음 두 단계를 동적 테이블로 변환하면서 마지막 단계는 스트림과 태스크에 유지할 수 있어요(저장 프로시저를 호출하기 때문). 동적 테이블에서 태스크로 변경을 읽으려면 동적 테이블에 스트림을 만드세요.
-- 동적 테이블에 스트림 생성 (ON TABLE이 아닌 ON DYNAMIC TABLE 사용).
CREATE OR REPLACE STREAM dt_orders_stream
ON DYNAMIC TABLE dt_orders;
-- 태스크가 평소처럼 이 스트림에서 읽습니다.
CREATE OR REPLACE TASK final_step
WAREHOUSE = transform_wh
SCHEDULE = '10 MINUTE'
WHEN SYSTEM$STREAM_HAS_DATA('dt_orders_stream')
AS
CALL my_procedure_with_side_effects(dt_orders_stream);
🔴 동적 테이블에 스트림을 만들 때는
ON DYNAMIC TABLE을 사용해야 해요.ON TABLE을 사용하면 오류가 발생해요. 또한 트리거된 태스크(AFTER절과 스트림 구동 트리거 사용)는 동적 테이블 스트림에서 작동하지 않아요. 대신WHEN SYSTEM$STREAM_HAS_DATA()검사와 함께 예약된 태스크를 사용하세요.
마이그레이션 중 흔한 오류
다음은 동적 테이블로 마이그레이션할 때 가장 자주 발생하는 문제예요. 동적 테이블 오류 조건의 전체 목록은 동적 테이블 갱신 문제 해결 문서를 참조하세요.
비결정적 함수가 전체 갱신을 강제
정의에 비결정적 함수가 포함되고 REFRESH_MODE=AUTO를 설정하면 Snowflake가 FULL로 결정해요. 생성 후 SHOW DYNAMIC TABLES로 결정된 모드를 확인하고 refresh_mode와 refresh_mode_reason 컬럼을 확인하세요. 자세한 내용은 갱신 모드와 동적 테이블 지원 쿼리 문서를 참조하세요.
스키마 변경이 동적 테이블에 미치는 영향
SELECT *를 사용하는 동적 테이블은 대부분의 기본 테이블 스키마 변경(컬럼 추가 또는 삭제)에 자동으로 적응해요. 호환되지 않는 유형 변경(예: TEXT → INT)만 갱신 실패를 일으켜요. 명시적 컬럼 목록이 있는 동적 테이블은 업스트림 스키마 변경을 반영하려면 CREATE OR REPLACE 또는 CREATE OR ALTER가 필요해요. 스키마 진화 동작에 대한 자세한 내용은 동적 테이블 수정 문서를 참조하세요.
다음 단계
마이그레이션 후 파이프라인을 효과적으로 운영하는 방법을 배우세요:
- 동적 테이블 안에서
MERGE또는INSERT로직을 사용하려면 커스텀 증분화 문서 참조. - 갱신 상태를 모니터링하려면 동적 테이블 모니터링 문서 참조.
- 갱신 성능을 최적화하려면 증분 갱신용 쿼리 최적화 문서 참조.
- 비용을 관리하려면 동적 테이블 비용 이해 문서 참조.
- 갱신 실패를 해결하려면 동적 테이블 갱신 문제 해결 문서 참조.