태스크와 DAG 오케스트레이션

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

Emily Melhuish

Technical Curriculum Developer, Snowflake

수동 작업의 해결책: 태스크

  • 팀이 매일 아침 SQL 스크립트를 수동으로 실행
  • 한 번 누락되면 비즈니스 가시성 손실
  • Snowflake 태스크로 완전 자동화 가능

예시:

-- Today's manual routine (fragile)
CALL logistics.refresh_ops_dashboard();
Snowflake 데이터 파이프라인 자동화

태스크에 대해 알아야 할 사항

  • 태스크는 스케줄에 따라 SQL 실행 — 외부 CRON 불필요
  • SQL 문 실행 또는 저장 프로시저 호출
  • 컴퓨팅과 스케줄링이 Snowflake 내에서 기본 제공
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Snowflake 데이터 파이프라인 자동화

독립형 태스크 생성

CREATE OR REPLACE TASK transform_delivery_summary 
  WAREHOUSE = harbr_wh 
  SCHEDULE = 'USING CRON 5 * * * * UTC' 
  AS INSERT INTO delivery_summary 

SELECT shipment_id, status, updated_at 
FROM delivery_events 
WHERE processed = FALSE; 
Snowflake 데이터 파이프라인 자동화

웨어하우스 기반 vs 서버리스

웨어하우스 기반

  • 지정된 가상 웨어하우스 사용
  • 실행당 최소 60초 요금 청구
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

서버리스

  • WAREHOUSE 절 생략
  • 사용한 초 단위로 청구, 유휴 비용 없음
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Snowflake 데이터 파이프라인 자동화

DAG 기반 태스크 오케스트레이션

DAG 다이어그램 — ingest_raw(루트 태스크, 시계 아이콘) → clean_events(ingest_raw 이후 실행) → build_summary(clean_events 이후 실행)

  • DAG: 방향성 비순환 그래프
  • 루트 태스크가 CRON 스케줄을 보유
  • 자식 태스크는 AFTER로 선행 태스크를 선언
  • 전체 체인이 중앙에서 관리됨
  • DAG는 실행 시점과 순서를 제어
Snowflake 데이터 파이프라인 자동화

태스크 상태 관리

상태 다이어그램 — SUSPENDED → STARTED → SUCCEEDED 또는 FAILED

-- Activate a task
ALTER TASK mytask RESUME; 
-- Pause a task
ALTER TASK mytask SUSPEND; 
-- Execute a task
EXECUTE TASK mytask SUSPEND;
Snowflake 데이터 파이프라인 자동화

태스크와 스트림 결합

CDC 파이프라인 다이어그램 — delivery_events → 스트림이 변경 사항 캡처 → 스트림에 데이터 있음? YES → 태스크 실행 → delivery_summary 업데이트 / NO → 태스크 건너뜀

CREATE OR REPLACE TASK logistics.process_delivery_events
    WAREHOUSE = compute_wh
    SCHEDULE = '5 minute'
    WHEN SYSTEM$STREAM_HAS_DATA('logistics.delivery_events_stream')
AS
INSERT INTO logistics.processed_events
SELECT * FROM logistics.delivery_events_stream;
Snowflake 데이터 파이프라인 자동화

연습해 봅시다!

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

Preparing Video For Download...