스트림과 변경 데이터 캡처

Snowflake 데이터 파이프라인 자동화

Emily Melhuish

Technical Curriculum Developer, Snowflake

변경된 데이터만 처리하기

CDC: 변경 데이터 캡처(Change Data Capture) 데이터 처리 방식 두 가지 비교

Snowflake 데이터 파이프라인 자동화

스트림

스트림의 기능

  • 소스 테이블의 모든 INSERT, UPDATE, DELETE를 추적합니다

Screenshot 2026-05-11 at 10.50.31 am.png

  • 데이터 중복 없이 변경 로그를 유지합니다
  • 소비되면 오프셋이 이동하며, 다음 읽기는 새로 시작됩니다
1 * Snowflake Learning Resource
Snowflake 데이터 파이프라인 자동화

스트림 유형

 

스트림 유형 캡처 대상 적합한 용도
Standard 모든 테이블 유형 및 뷰, 모든 DML 변경(삽입, 업데이트, 삭제) 임의 행이 변경될 수 있는 테이블(예: 배송)
Append-only 외부 테이블 제외 모든 테이블 유형 및 뷰, 행 삽입만 추적 삽입 전용 테이블(예: 배송 이벤트) - 더 효율적
Insert-only 외부 관리 Apache Iceberg 및 외부 테이블, 행 삽입만 추적 외부 테이블

디렉터리 테이블은 스테이지의 파일 메타데이터(이름, 크기, 마지막 수정 타임스탬프)를 제공합니다

Snowflake 데이터 파이프라인 자동화

스트림 생성

shipments 테이블의 Standard 스트림

CREATE STREAM shipments_stream
  ON TABLE logistics.shipments;

delivery_events 테이블의 Append-only 스트림

CREATE STREAM delivery_events_stream
  ON TABLE logistics.delivery_events
  APPEND_ONLY = TRUE;
Snowflake 데이터 파이프라인 자동화

스트림 메타데이터 열

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTION: INSERT 또는 DELETE
  • METADATA$ISUPDATE: 업데이트 쌍의 일부일 때 TRUE
  • METADATA$ROW_ID: 고유한 물리적 행 식별자
  • 업데이트는 DELETE + INSERT 쌍으로 나타나며, 둘 다 METADATA$ISUPDATE = TRUE로 표시됨

METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID 열과 샘플 데이터를 보여주는 스트림 메타데이터 열

Snowflake 데이터 파이프라인 자동화

스트림 오프셋

타임라인 다이어그램 - 스트림 생성(오프셋 시작) → 소스 테이블에서 변경 발생(스트림에 레코드 누적) → 트랜잭션에서 스트림 소비(오프셋이 현재 시점으로 이동)

Snowflake 데이터 파이프라인 자동화

파이프라인에서의 스트림 개요

  • 스트림은 태스크와 함께 사용됩니다 — 태스크는 SQL을 스케줄에 따라 실행하는 Snowflake 객체입니다
  • 태스크는 변경된 행만 읽습니다. 1천만 행 기준: 2초 vs 2분

Screenshot 2026-05-11 at 10.48.51 am.png

1 * Snowflake Learning Resource
Snowflake 데이터 파이프라인 자동화

파이프라인에서의 스트림 쿼리

  • 스트림은 태스크와 함께 사용됩니다 - 태스크는 SQL을 스케줄에 따라 실행하는 Snowflake 객체입니다
  • 태스크는 변경된 행만 읽습니다. 1천만 행 기준: 2초 vs 2분
CREATE TASK logistics.sync_shipments
  WAREHOUSE = compute_wh
  SCHEDULE = '5 MINUTE'
  WHEN SYSTEM$STREAM_HAS_DATA('logistics.staging_shipments_stream')
AS
  INSERT INTO logistics.shipments
  SELECT shipment_id, region, carrier, delivery_days
  FROM logistics.staging_shipments_stream
  WHERE METADATA$ACTION = 'INSERT';
Snowflake 데이터 파이프라인 자동화

연습해 봅시다!

Snowflake 데이터 파이프라인 자동화

Preparing Video For Download...