Orchestration des tâches et DAG

Automatisation des pipelines de données dans Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Solution pour les charges manuelles : tâches

  • L’équipe exécute un script SQL manuellement chaque matin
  • Un oubli = perte de visibilité pour l’entreprise
  • Les tâches Snowflake automatisent tout ça

Exemple :

-- Routine manuelle d’aujourd’hui (fragile)
CALL logistics.refresh_ops_dashboard();
Automatisation des pipelines de données dans Snowflake

À savoir sur les tâches

  • Les tâches exécutent du SQL selon un horaire — pas de CRON externe requis
  • Exécuter des instructions SQL ou appeler des procédures stockées
  • Calcul et horaire natifs dans Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Automatisation des pipelines de données dans Snowflake

Créer une tâche autonome

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; 
Automatisation des pipelines de données dans Snowflake

Entrepôt dédié vs sans serveur

Basé sur entrepôt

  • Utilise un entrepôt virtuel nommé
  • Facturation minimale de 60 s par exécution
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

Sans serveur (serverless)

  • Omettre la clause WAREHOUSE
  • Facturé à la seconde; aucun coût d’inactivité
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Automatisation des pipelines de données dans Snowflake

Orchestration de tâches avec DAG

Diagramme de DAG — ingest_raw (tâche racine, icône d’horloge) → clean_events (après ingest_raw) → build_summary (après clean_events)

  • DAG : graphe orienté acyclique
  • La tâche racine contient l’horaire CRON
  • Les tâches enfants déclarent leur prédécesseur avec AFTER
  • Chaîne entière centralisée
  • Les DAG contrôlent quand et dans quel ordre ça s’exécute
Automatisation des pipelines de données dans Snowflake

Gérer les états des tâches

Diagramme d’états — SUSPENDED → STARTED → SUCCEEDED ou FAILED

-- Activer une tâche
ALTER TASK mytask RESUME; 
-- Mettre une tâche en pause
ALTER TASK mytask SUSPEND; 
-- Exécuter une tâche
EXECUTE TASK mytask SUSPEND;
Automatisation des pipelines de données dans Snowflake

Combiner tâches et flux (streams)

Diagramme de pipeline CDC — delivery_events → le flux capte les changements → le flux a des données ? OUI → la tâche se déclenche → delivery_summary mis à jour / NON → la tâche saute

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;
Automatisation des pipelines de données dans Snowflake

Passons à la pratique !

Automatisation des pipelines de données dans Snowflake

Preparing Video For Download...