Orchestrazione di task e DAG

Automazione delle pipeline di dati in Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Soluzione per lavori manuali: Task

  • Il team esegue uno script SQL a mano ogni mattina
  • Se lo salti, l'azienda perde visibilità
  • I task Snowflake automatizzano tutto

Esempio:

-- Routine manuale di oggi (fragile)
CALL logistics.refresh_ops_dashboard();
Automazione delle pipeline di dati in Snowflake

Cosa sapere sui task

  • I task eseguono SQL a orari pianificati — nessun CRON esterno
  • Eseguono istruzioni SQL o chiamano procedure memorizzate
  • Calcolo e schedulazione nativi in Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Automazione delle pipeline di dati in Snowflake

Creare un task standalone

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; 
Automazione delle pipeline di dati in Snowflake

Warehouse vs serverless

Basati su warehouse

  • Usano un virtual warehouse nominato
  • Fatturazione minima 60 secondi per esecuzione
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

Serverless

  • Omessi la clausola WAREHOUSE
  • Addebito al secondo; nessun costo di idle
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Automazione delle pipeline di dati in Snowflake

Orchestrazione dei task basata su DAG

Diagramma DAG — ingest_raw (task radice, icona orologio) → clean_events (dopo ingest_raw) → build_summary (dopo clean_events)

  • DAG: Directed Acyclic Graph
  • Il task radice contiene la pianificazione CRON
  • I figli dichiarano il predecessore con AFTER
  • L'intera catena è centralizzata
  • I DAG controllano quando e in che ordine eseguire
Automazione delle pipeline di dati in Snowflake

Gestire gli stati dei task

Diagramma degli stati — SUSPENDED → STARTED → SUCCEEDED o FAILED

-- Attiva un task
ALTER TASK mytask RESUME; 
-- Metti in pausa un task
ALTER TASK mytask SUSPEND; 
-- Esegui un task
EXECUTE TASK mytask SUSPEND;
Automazione delle pipeline di dati in Snowflake

Task e stream insieme

Diagramma CDC — delivery_events → lo stream cattura le modifiche → lo stream ha dati? SÌ → il task parte → delivery_summary aggiornato / NO → il task salta

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;
Automazione delle pipeline di dati in Snowflake

Ayo berlatih!

Automazione delle pipeline di dati in Snowflake

Preparing Video For Download...