Taken en DAG-orkestratie

Automatisering van datapijplijnen in Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Oplossing voor handwerk: taken

  • Team draait elke ochtend handmatig een SQL-script
  • Eén gemiste run = geen zicht voor het bedrijf
  • Snowflake-taken automatiseren dit volledig

Voorbeeld:

-- De handmatige routine van vandaag (fragiel)
CALL logistics.refresh_ops_dashboard();
Automatisering van datapijplijnen in Snowflake

Wat moeten we weten over taken

  • Taken draaien SQL op schema — geen externe CRON nodig
  • Voer SQL-statements uit of roep stored procedures aan
  • Compute en planning draaien native in Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Automatisering van datapijplijnen in Snowflake

Een losstaande taak maken

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; 
Automatisering van datapijplijnen in Snowflake

Warehouse vs. serverless

Op warehouse gebaseerde

  • Gebruikt een benoemd virtueel warehouse
  • Minimaal 60 seconden kosten per run
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

Serverless

  • Laat de WAREHOUSE-clausule weg
  • Afrekening per seconde verbruik; geen idle-kosten
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Automatisering van datapijplijnen in Snowflake

DAG-gebaseerde taakorkestratie

DAG-diagram — ingest_raw (roottask, klokpictogram) → clean_events (draait na ingest_raw) → build_summary (draait na clean_events)

  • DAG's: Directed Acyclic Graph
  • Roottask bevat het CRON-schema
  • Kindtaken geven hun voorganger op met AFTER
  • Hele keten is gecentraliseerd
  • DAG's bepalen wanneer en in welke volgorde dingen draaien
Automatisering van datapijplijnen in Snowflake

Taakstatussen beheren

Statusdiagram — SUSPENDED → STARTED → SUCCEEDED of FAILED

-- Activeer een taak
ALTER TASK mytask RESUME; 
-- Pauzeer een taak
ALTER TASK mytask SUSPEND; 
-- Voer een taak uit
EXECUTE TASK mytask SUSPEND;
Automatisering van datapijplijnen in Snowflake

Taken en streams samen

CDC-pijplijn — delivery_events → stream legt wijzigingen vast → stream heeft data? JA → taak start → delivery_summary bijgewerkt / NEE → taak slaat over

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;
Automatisering van datapijplijnen in Snowflake

Laten we oefenen!

Automatisering van datapijplijnen in Snowflake

Preparing Video For Download...