Orkestrasi Tugas dan DAG

Otomatisasi Data Pipeline di Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Solusi untuk Beban Kerja Manual: Tasks

  • Tim menjalankan skrip SQL manual tiap pagi
  • Sekali terlewat — bisnis kehilangan visibilitas
  • Snowflake tasks mengotomatiskan sepenuhnya

Contoh:

-- Rutinitas manual hari ini (rapuh)
CALL logistics.refresh_ops_dashboard();
Otomatisasi Data Pipeline di Snowflake

Hal yang Perlu Diketahui tentang Tasks

  • Tasks menjalankan SQL terjadwal — tanpa CRON eksternal
  • Menjalankan pernyataan SQL atau memanggil stored procedure
  • Komputasi dan penjadwalan native di Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Otomatisasi Data Pipeline di Snowflake

Membuat Task Mandiri

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; 
Otomatisasi Data Pipeline di Snowflake

Berbasis Warehouse vs Serverless

Berbasis Warehouse

  • Menggunakan virtual warehouse bernama
  • Penagihan minimum 60 detik per run
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

Serverless

  • Hilangkan klausa WAREHOUSE
  • Ditagih per detik terpakai; tanpa biaya idle
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Otomatisasi Data Pipeline di Snowflake

Orkestrasi Tugas Berbasis DAG

Diagram DAG — ingest_raw (tugas root, ikon jam) → clean_events (jalan setelah ingest_raw) → build_summary (jalan setelah clean_events)

  • DAG: Directed Acyclic Graph
  • Tugas root memegang jadwal CRON
  • Tugas anak menyatakan pendahulu dengan AFTER
  • Seluruh rantai tersentralisasi
  • DAG mengatur kapan dan urutan eksekusi
Otomatisasi Data Pipeline di Snowflake

Mengelola Status Tugas

Diagram status — SUSPENDED → STARTED → SUCCEEDED atau FAILED

-- Aktifkan tugas
ALTER TASK mytask RESUME; 
-- Jeda tugas
ALTER TASK mytask SUSPEND; 
-- Eksekusi tugas
EXECUTE TASK mytask SUSPEND;
Otomatisasi Data Pipeline di Snowflake

Menggabungkan Tasks dan Streams

Diagram pipeline CDC — delivery_events → stream menangkap perubahan → stream ada data? YA → tugas jalan → delivery_summary diperbarui / TIDAK → tugas lewati

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;
Otomatisasi Data Pipeline di Snowflake

Ayo berlatih!

Otomatisasi Data Pipeline di Snowflake

Preparing Video For Download...