Orchestration des tâches et des 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é métier
  • Les tâches Snowflake automatisent tout cela

Exemple :

-- Routine manuelle du jour (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 planning — pas de CRON externe requis
  • Exécutent des instructions SQL ou appellent des procédures stockées
  • Calcul et planification 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 serverless

Basées sur entrepôt

  • Utilisent 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)

  • Omettez la clause WAREHOUSE
  • Facturation à la seconde consommée ; pas de 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 basée sur des DAG

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

  • DAG : graphe orienté acyclique
  • La tâche racine porte la planification CRON
  • Les tâches enfants déclarent leur prédécesseur avec AFTER
  • Toute la chaîne est 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 en pause une tâche
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 streams

Schéma CDC — delivery_events → le stream capture les changements → le stream a des données ? OUI → la tâche se déclenche → delivery_summary mis à jour / NON → la tâche est ignorée

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...