Úlohy a orchestrace DAG

Automatizace datových pipeline ve Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Řešení pro ruční pracovní zátěž: úlohy

  • Tým spouští SQL skript každé ráno ručně
  • Jedno vynechané spuštění — firma ztrácí přehled
  • Snowflake úlohy to zcela automatizují

Příklad:

-- Today's manual routine (fragile)
CALL logistics.refresh_ops_dashboard();
Automatizace datových pipeline ve Snowflake

Co potřebujeme vědět o úlohách

  • Úlohy spouštějí SQL dle plánu — bez externího CRONu
  • Spouštějí SQL příkazy nebo uložené procedury
  • Výpočet i plánování nativně ve Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Automatizace datových pipeline ve Snowflake

Vytvoření samostatné úlohy

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; 
Automatizace datových pipeline ve Snowflake

Warehouse-based vs. serverless

Warehouse-based

  • Používá pojmenovaný virtuální warehouse
  • Minimální fakturace 60 sekund za spuštění
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

Serverless

  • Vynechejte klauzuli WAREHOUSE
  • Fakturace za spotřebované sekundy; bez nákladů v nečinnosti
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Automatizace datových pipeline ve Snowflake

Orchestrace úloh pomocí DAG

Diagram DAG — ingest_raw (kořenová úloha, ikona hodin) → clean_events (spustí se po ingest_raw) → build_summary (spustí se po clean_events)

  • DAG: orientovaný acyklický graf
  • Kořenová úloha obsahuje CRON plán
  • Podřízené úlohy deklarují předchůdce pomocí AFTER
  • Celý řetězec je centralizován
  • DAG řídí, kdy a v jakém pořadí se věci spouštějí
Automatizace datových pipeline ve Snowflake

Správa stavů úloh

Diagram stavů — SUSPENDED → STARTED → SUCCEEDED nebo FAILED

-- Activate a task
ALTER TASK mytask RESUME; 
-- Pause a task
ALTER TASK mytask SUSPEND; 
-- Execute a task
EXECUTE TASK mytask SUSPEND;
Automatizace datových pipeline ve Snowflake

Kombinace úloh a streamů

Diagram CDC pipeline — delivery_events → stream zachycuje změny → stream má data? ANO → úloha se spustí → delivery_summary aktualizováno / NE → úloha přeskočena

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;
Automatizace datových pipeline ve Snowflake

Pojďme si procvičit!

Automatizace datových pipeline ve Snowflake

Preparing Video For Download...