Aufgaben und DAG-Orchestrierung

Automatisierung von Datenpipelines in Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Lösung für manuelle Workloads: Tasks

  • Team führt jeden Morgen manuell ein SQL-Skript aus
  • Ein Lauf verpasst – Business verliert Sichtbarkeit
  • Snowflake-Tasks automatisieren das komplett

Beispiel:

-- Heutige manuelle Routine (fragil)
CALL logistics.refresh_ops_dashboard();
Automatisierung von Datenpipelines in Snowflake

Was müssen wir über Tasks wissen?

  • Tasks führen SQL nach Zeitplan aus — kein externer CRON nötig
  • Führe SQL-Statements aus oder rufe Stored Procedures auf
  • Compute und Scheduling laufen nativ in Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Automatisierung von Datenpipelines in Snowflake

Eigenständige Task erstellen

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; 
Automatisierung von Datenpipelines in Snowflake

Warehouse-basiert vs. serverlos

Warehouse-basiert

  • Verwendet ein benanntes virtuelles Warehouse
  • 60 Sekunden Mindestabrechnung pro Lauf
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

Serverlos

  • WAREHOUSE-Klausel weglassen
  • Abrechnung pro Sekunde; keine Leerkosten
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Automatisierung von Datenpipelines in Snowflake

DAG-basierte Aufgabenorchestrierung

DAG-Diagramm — ingest_raw (Wurzelaufgabe, Uhrsymbol) → clean_events (läuft nach ingest_raw) → build_summary (läuft nach clean_events)

  • DAGs: gerichtete azyklische Graphen
  • Wurzelaufgabe enthält den CRON-Zeitplan
  • Kindaufgaben geben Vorgänger mit AFTER an
  • Ganze Kette ist zentralisiert
  • DAGs steuern, wann und in welcher Reihenfolge etwas läuft
Automatisierung von Datenpipelines in Snowflake

Aufgabenstatus verwalten

Zustandsdiagramm — SUSPENDED → STARTED → SUCCEEDED oder FAILED

-- Aufgabe aktivieren
ALTER TASK mytask RESUME; 
-- Aufgabe pausieren
ALTER TASK mytask SUSPEND; 
-- Aufgabe ausführen
EXECUTE TASK mytask SUSPEND;
Automatisierung von Datenpipelines in Snowflake

Aufgaben und Streams kombiniert

CDC-Pipeline — delivery_events → Stream erfasst Änderungen → hat Stream Daten? JA → Task wird ausgelöst → delivery_summary aktualisiert / NEIN → Task überspringt

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;
Automatisierung von Datenpipelines in Snowflake

Lass uns üben!

Automatisierung von Datenpipelines in Snowflake

Preparing Video For Download...