タスクとDAGオーケストレーション

Snowflakeにおけるデータパイプラインの自動化

Emily Melhuish

Technical Curriculum Developer, Snowflake

手動ワークロードの解決策:タスク

  • チームは毎朝SQLスクリプトを手動実行
  • 1回の実行漏れで業務の可視性が失われる
  • 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...