任務與 DAG 編排

在 Snowflake 進行資料管線自動化

Emily Melhuish

Technical Curriculum Developer, Snowflake

手動工作量的解方:Tasks

  • 團隊每天早上手動執行 SQL 指令碼
  • 漏跑一次,業務就失去可見性
  • Snowflake 任務可完全自動化

範例:

-- 今天的手動流程(脆弱)
CALL logistics.refresh_ops_dashboard();
在 Snowflake 進行資料管線自動化

Tasks 的重點

  • 任務按排程執行 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 進行資料管線自動化

建立獨立任務(Standalone Task)

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 進行資料管線自動化

Warehouse 型 vs 無伺服器

以 Warehouse 為基礎

  • 使用具名虛擬 warehouse
  • 每次執行最少計費 60 秒
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

無伺服器(Serverless)

  • 省略 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

-- 啟用任務
ALTER TASK mytask RESUME; 
-- 暫停任務
ALTER TASK mytask SUSPEND; 
-- 立即執行任務
EXECUTE TASK mytask SUSPEND;
在 Snowflake 進行資料管線自動化

結合 Tasks 與 Streams

CDC 管線圖 — delivery_events → stream 擷取變更 → stream 有資料?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...