Uppgifter och DAG-orkestrering

Automatisering av datapipelines i Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Lösning för manuella arbetsflöden: uppgifter

  • Teamet kör ett SQL-skript manuellt varje morgon
  • En missad körning — verksamheten tappar insyn
  • Snowflake-uppgifter automatiserar detta helt

Exempel:

-- Today's manual routine (fragile)
CALL logistics.refresh_ops_dashboard();
Automatisering av datapipelines i Snowflake

Vad behöver vi veta om uppgifter

  • Uppgifter kör SQL enligt schema — ingen extern CRON krävs
  • Kör SQL-satser eller anropar lagrade procedurer
  • Beräkning och schemaläggning sker nativt i Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Automatisering av datapipelines i Snowflake

Skapa en fristående uppgift

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 av datapipelines i Snowflake

Lagerbaserad vs. serverlös

Lagerbaserad

  • Använder ett namngivet virtuellt lager
  • Minst 60 sekunders debitering per körning
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

Serverlös

  • Utelämna WAREHOUSE-satsen
  • Debiteras per förbrukad sekund; ingen vilokostnad
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Automatisering av datapipelines i Snowflake

DAG-baserad uppgiftsorkestrering

DAG-diagram — ingest_raw (rotuppgift, klockikon) → clean_events (körs efter ingest_raw) → build_summary (körs efter clean_events)

  • DAG: Directed Acyclic Graph
  • Rotuppgiften innehåller CRON-schemat
  • Underuppgifter deklarerar föregångaren med AFTER
  • Hela kedjan är centraliserad
  • DAG:ar styr när och i vilken ordning saker körs
Automatisering av datapipelines i Snowflake

Hantera uppgiftstillstånd

Tillståndsdiagram — SUSPENDED → STARTED → SUCCEEDED eller FAILED

-- Activate a task
ALTER TASK mytask RESUME; 
-- Pause a task
ALTER TASK mytask SUSPEND; 
-- Execute a task
EXECUTE TASK mytask SUSPEND;
Automatisering av datapipelines i Snowflake

Uppgifter och strömmar kombinerade

CDC-pipelinediagram — delivery_events → stream fångar ändringar → stream har data? JA → uppgift körs → delivery_summary uppdateras / NEJ → uppgift hoppas över

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 av datapipelines i Snowflake

Laten we oefenen!

Automatisering av datapipelines i Snowflake

Preparing Video For Download...