Sarcini și orchestrare DAG

Automatizarea pipeline-urilor de date în Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Soluție pentru fluxuri manuale: sarcini

  • Echipa rulează manual un script SQL în fiecare dimineață
  • O rulare ratată — businessul pierde vizibilitate
  • Sarcinile Snowflake automatizează complet acest proces

Exemplu:

-- Today's manual routine (fragile)
CALL logistics.refresh_ops_dashboard();
Automatizarea pipeline-urilor de date în Snowflake

Ce trebuie să știm despre sarcini

  • Sarcinile rulează SQL conform unui program — fără CRON extern
  • Execută instrucțiuni SQL sau apelează proceduri stocate
  • Calcul și planificare native în Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Automatizarea pipeline-urilor de date în Snowflake

Crearea unei sarcini independente

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; 
Automatizarea pipeline-urilor de date în Snowflake

Warehouse vs. serverless

Bazată pe warehouse

  • Utilizează un warehouse virtual denumit
  • Facturare minimă de 60 de secunde per rulare
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

Serverless

  • Omiteți clauza WAREHOUSE
  • Facturare per secundă consumată; fără costuri inactive
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Automatizarea pipeline-urilor de date în Snowflake

Orchestrare a sarcinilor bazată pe DAG

Diagramă DAG — ingest_raw (sarcină rădăcină, pictogramă ceas) → clean_events (după ingest_raw) → build_summary (după clean_events)

  • DAG: Graf Aciclic Direcționat
  • Sarcina rădăcină deține programul CRON
  • Sarcinile copil declară predecesorul cu AFTER
  • Întregul lanț este centralizat
  • DAG-urile controlează când și în ce ordine rulează lucrurile
Automatizarea pipeline-urilor de date în Snowflake

Gestionarea stărilor sarcinilor

Diagramă de stare — SUSPENDED → STARTED → SUCCEEDED sau FAILED

-- Activate a task
ALTER TASK mytask RESUME; 
-- Pause a task
ALTER TASK mytask SUSPEND; 
-- Execute a task
EXECUTE TASK mytask SUSPEND;
Automatizarea pipeline-urilor de date în Snowflake

Combinarea sarcinilor și a stream-urilor

Diagramă pipeline CDC — delivery_events → stream captează modificările → stream are date? DA → sarcina se declanșează → delivery_summary actualizat / NU → sarcina este omisă

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;
Automatizarea pipeline-urilor de date în Snowflake

Să exersăm!

Automatizarea pipeline-urilor de date în Snowflake

Preparing Video For Download...