다이나믹 테이블에 스트림(stream) 사용하기
다이나믹 테이블에 스트림(stream) 사용하기
스트림은 다이나믹 테이블에 대해 행 단위의 insert, update, delete를 포착해요. 읽을 때마다 이전 읽기 이후의 순 변경(net changes)을 반환하므로, 다이나믹 테이블이 스스로는 도달하지 못하는 대상(저장 프로시저, 외부 API 호출, task 안의 조건부 로직)으로 파이프라인을 연결하는 데 쓸 수 있어요.
출처: Snowflake 문서
본문
이 예시는 증분 리프레시(incremental refresh)를 사용하는 기존 다이나믹 테이블 <C>dt_orders</C>가 있다고 가정합니다. 설정 가이드는 다이나믹 테이블 만들기를 참고하세요.
중요 스트림을 만들 때는 ON DYNAMIC TABLE 구문을 사용하세요. ON TABLE은 다이나믹 테이블에서 동작하지 않고 오류를 반환합니다. 일반적인 오류의 예를 참고하세요.
CREATE OR REPLACE STREAM dt_orders_stream ON DYNAMIC TABLE dt_orders;
지원되는 스트림 유형
다이나믹 테이블의 스트림은 표준(delta) 유형만 지원합니다. 표준 스트림은 스트림이 마지막으로 읽힌 이후의 insert, update, delete를 순 변경으로 포착합니다.
CREATE STREAM ... ON DYNAMIC TABLE <name>
APPEND_ONLY는 다이나믹 테이블의 스트림에서 지원되지 않아요. insert-only 의미론이 필요하다면 표준 스트림을 소비하고 다운스트림 task에서 <C>METADATA$ACTION = 'INSERT'</C>로 필터링하세요. 또는 커스텀 증분 다이나믹 테이블이 <C>CHANGES(INFORMATION => APPEND_ONLY)</C> 절을 통해 append-only 의미론을 기본 지원합니다.
참고 다이나믹 테이블은 증분 리프레시를 사용해야 합니다. 전체 리프레시(full refresh)를 사용하는 다이나믹 테이블은 매 리프레시마다 테이블 출력 전체를 대체하므로 추적할 증분 변경 기록이 없어서 스트림을 지원하지 않아요. REFRESH_MODE가 AUTO로 설정되어 있다면, 스트림을 만들기 전에
<C>SHOW DYNAMIC TABLES</C>를 실행하고<C>refresh_mode</C>컬럼으로 해석된 모드를 확인하세요.
이 패턴을 언제 사용하나요
핵심 사용 사례는 다이나믹 테이블이 선언적 정의를 처리하고, 트리거된 task(triggered task)가 다이나믹 테이블이 표현할 수 없는 명령형 작업을 처리하는 하이브리드 파이프라인입니다:
다이나믹 테이블 -> 스트림 -> 트리거된 task -> 대상
일반적인 시나리오:
- 다이나믹 테이블이 아닌 대상(일반 테이블, 외부 스테이지)에 쓰기
- 부작용이 있는 저장 프로시저나 외부 함수 호출
- 변경 세트에 조건부 로직(IF/ELSE, 루프) 실행
task의 WHEN 절에 SYSTEM$STREAM_HAS_DATA를 사용해서 스트림에 소비되지 않은 변경이 있을 때만 task가 실행되게 하세요. task는 다이나믹 테이블이 리프레시되어 새 출력을 만든 뒤에 실행됩니다. 행을 변경하지 않는 리프레시는 task를 트리거하지 않아요.
-- 스트림을 소비하는 트리거된 task 생성
CREATE OR REPLACE TASK load_completed_orders
WAREHOUSE = transform_wh
WHEN SYSTEM$STREAM_HAS_DATA('dt_orders_stream')
AS
INSERT INTO completed_orders_archive
SELECT order_id, customer_id, order_date, product_name, line_total
FROM dt_orders_stream
WHERE order_status = 'completed';
여러 task가 같은 다이나믹 테이블을 소비한다면 각각에 별도의 스트림을 만드세요. 자세한 내용은 스트림 소개를 참고하세요.
재초기화 후 스트림 동작
다이나믹 테이블이 재초기화되면(예: 기본 테이블의 CREATE OR REPLACE, 마스킹 정책 변경, change-tracking 토글), 스트림은 유지되지만 다음 읽기는 일반 리프레시가 만드는 것보다 훨씬 큰 변경 행 세트를 반환할 수 있어요. 스트림은 계속 연결된 채로 마지막 읽기 오프셋과 재초기화된 출력 사이의 모든 행 수준 차이를 노출합니다.
경고 재초기화 후에는 이전에 update였던 행이 METADATA$ISUPDATE가 FALSE로 설정된 DELETE + INSERT 쌍으로 나타납니다. METADATA$ISUPDATE에 의존해 실제 delete와 update를 구분하는 다운스트림 로직은 오류 없이 데이터를 유실하거나 중복할 수 있어요. 안전하게 처리하려면 조건부 INSERT/DELETE 분기 대신 멱등(idempotent) MERGE 문을 사용하세요.
다음 MERGE는 정상 변경과 재초기화 후 나타나는 DELETE + INSERT 쌍을 모두 처리합니다. 세 개의 절이 기본 키로 매칭하면서 delete, update, insert를 각각 다룹니다:
MERGE INTO completed_orders_archive t
USING (SELECT order_id, customer_id, order_date, product_name, line_total,
METADATA$ACTION, METADATA$ISUPDATE
FROM dt_orders_stream
WHERE order_status = 'completed') s
ON t.order_id = s.order_id
WHEN MATCHED AND s.METADATA$ACTION = 'DELETE' AND NOT s.METADATA$ISUPDATE
THEN DELETE
WHEN MATCHED AND s.METADATA$ACTION = 'INSERT'
THEN UPDATE SET t.customer_id = s.customer_id,
t.order_date = s.order_date,
t.product_name = s.product_name,
t.line_total = s.line_total
WHEN NOT MATCHED AND s.METADATA$ACTION = 'INSERT'
THEN INSERT (order_id, customer_id, order_date, product_name, line_total)
VALUES (s.order_id, s.customer_id, s.order_date, s.product_name, s.line_total);
재초기화 후에는 살아남는 모든 행이 <C>METADATA$ISUPDATE = FALSE</C>인 DELETE + INSERT 쌍으로 나타납니다. MERGE는 각 쌍을 올바르게 처리해요: 절 1이 이전 행을 제거하고 절 3이 다시 삽입합니다. 순효과는 중복이나 데이터 유실 없이 같은 데이터입니다. 정상 변경의 경우 절 2가 제자리에서 update를 처리합니다.
다운스트림 task가 이런 큰 변경 세트를 우아하게 처리하도록 설계하세요. 멱등 insert나 일괄 처리(batch processing)를 사용해서 재초기화 후 읽기가 중복이나 누락 데이터를 만들지 않게 하세요.
제한: 스트림은 순 변경만 포착하고 전체 기록은 아님
다이나믹 테이블의 스트림을 감사 로그(audit log)로 사용하지 마세요. 한 리프레시 주기 안에서 행이 insert된 뒤 update되면, 스트림은 최종 상태만 보여 줍니다. 두 리프레시 사이에 스트림을 읽지 않으면 두 리프레시의 변경이 단일 순 결과로 합쳐집니다.
컴플라이언스나 감사를 위해 모든 상태 전이를 원한다면, 수집 계층에서 append-only 테이블에 직접 쓰세요. 스트림이 변경을 추적하는 방식에 대한 자세한 내용은 스트림 소개를 참고하세요.
스트림 확인하기
스트림을 만든 뒤, 연결되어 있고 stale하지 않은지 확인하세요:
SHOW STREAMS LIKE 'dt_orders_stream';
+--------------------+---------------+---------------+--------------+-------+-------+
| name | database_name | schema_name | source_type | mode | stale |
|--------------------+---------------+---------------+--------------+-------+-------|
| DT_ORDERS_STREAM | MY_DB | PUBLIC | Dynamic Table| DEFAULT| false |
+--------------------+---------------+---------------+--------------+-------+-------+
출력 컬럼은 가독성을 위해 잘랐습니다. <C>stale</C> 컬럼으로 스트림이 활성 상태인지 확인하세요. <C>stale</C>가 <C>true</C>면 스트림이 다이나믹 테이블의 변경 기록보다 뒤처진 것이므로 다시 만들어야 합니다. stale을 막으려면 데이터 보존 기간(data retention period) 안에 스트림을 정기적으로 읽으세요. 스트림 staleness에 대한 자세한 내용은 스트림 소개를 참고하세요.
일반적인 오류
ON DYNAMIC TABLE 대신 ON TABLE을 사용하면 오류가 반환됩니다:
-- 잘못된 예: ON TABLE은 다이나믹 테이블에서 동작하지 않습니다
CREATE OR REPLACE STREAM dt_orders_stream ON TABLE dt_orders;
002203 (42601): SQL compilation error:
Object found is of type 'DYNAMIC_TABLE', not specified type 'TABLE'.
다이나믹 테이블에 스트림을 만들려면 ON DYNAMIC TABLE 구문을 사용하세요:
-- 올바른 구문
CREATE OR REPLACE STREAM dt_orders_stream ON DYNAMIC TABLE dt_orders;
필요한 권한
다이나믹 테이블에 스트림을 만들려면 여러분의 역할에 다음이 필요합니다:
- 스키마에 대한 CREATE STREAM
- 다이나믹 테이블에 대한 SELECT
전체 액세스 제어 요구사항은 CREATE STREAM을 참고하세요.
다음 단계
- 다이나믹 테이블 리프레시 동작 방식을 배우려면 다이나믹 테이블 리프레시 모드를 참고하세요.
- 상위 다이나믹 테이블의 리프레시 건강 상태를 모니터링하려면 다이나믹 테이블 모니터링을 참고하세요.
- 스트림에 영향을 주는 리프레시 실패를 해결하려면 다이나믹 테이블 리프레시 문제 해결을 참고하세요.
- 전체 스트림 SQL 레퍼런스는 CREATE STREAM을 참고하세요.