Zadania i orkiestracja DAG

Automatyzacja potoków danych w Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Rozwiązanie dla ręcznych obciążeń: zadania

  • Zespół ręcznie uruchamia skrypt SQL każdego ranka
  • Jedna pominięta operacja — brak widoczności w biznesie
  • Zadania Snowflake w pełni automatyzują ten proces

Przykład:

-- Today's manual routine (fragile)
CALL logistics.refresh_ops_dashboard();
Automatyzacja potoków danych w Snowflake

Co należy wiedzieć o zadaniach

  • Zadania wykonują SQL według harmonogramu — bez zewnętrznego CRON
  • Wykonują instrukcje SQL lub wywołują procedury składowane
  • Obliczenia i harmonogramowanie działają natywnie w Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Automatyzacja potoków danych w Snowflake

Tworzenie samodzielnego zadania

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; 
Automatyzacja potoków danych w Snowflake

Magazyn wirtualny vs. bezserwerowe

Na bazie magazynu

  • Używa nazwanego wirtualnego magazynu
  • Minimalne rozliczenie: 60 sekund na uruchomienie
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

Bezserwerowe

  • Pominąć klauzulę WAREHOUSE
  • Rozliczenie per sekunda; brak kosztów bezczynności
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Automatyzacja potoków danych w Snowflake

Orkiestracja zadań za pomocą DAG

Diagram DAG — ingest_raw (zadanie główne, ikona zegara) → clean_events (uruchamiane po ingest_raw) → build_summary (uruchamiane po clean_events)

  • DAG: skierowany graf acykliczny
  • Zadanie główne zawiera harmonogram CRON
  • Zadania podrzędne deklarują poprzednika za pomocą AFTER
  • Cały łańcuch jest scentralizowany
  • DAG kontroluje kiedy i w jakiej kolejności działają zadania
Automatyzacja potoków danych w Snowflake

Zarządzanie stanami zadań

Diagram stanów — SUSPENDED → STARTED → SUCCEEDED lub FAILED

-- Activate a task
ALTER TASK mytask RESUME; 
-- Pause a task
ALTER TASK mytask SUSPEND; 
-- Execute a task
EXECUTE TASK mytask SUSPEND;
Automatyzacja potoków danych w Snowflake

Zadania i strumienie razem

Diagram potoku CDC — delivery_events → strumień rejestruje zmiany → strumień ma dane? TAK → zadanie uruchamiane → delivery_summary zaktualizowane / NIE → zadanie pomijane

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;
Automatyzacja potoków danych w Snowflake

Czas na ćwiczenia!

Automatyzacja potoków danych w Snowflake

Preparing Video For Download...