Orquestación de tareas y DAG

Automatización de canalizaciones de datos en Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Solución para cargas manuales: tareas

  • El equipo ejecuta un script SQL manual cada mañana
  • Si se salta una ejecución, el negocio pierde visibilidad
  • Las tareas de Snowflake lo automatizan todo

Ejemplo:

-- Rutina manual de hoy (frágil)
CALL logistics.refresh_ops_dashboard();
Automatización de canalizaciones de datos en Snowflake

Qué hay que saber sobre las tareas

  • Las tareas ejecutan SQL con horario; no se requiere CRON externo
  • Ejecutan sentencias SQL o llaman procedimientos almacenados
  • Cómputo y planificación nativos en Snowflake
CREATE TASK logistics.refresh_dashboard
  WAREHOUSE = harbr_wh
  SCHEDULE = 'USING CRON 0 6 * * * UTC'
AS CALL logistics.refresh_ops_dashboard();
Automatización de canalizaciones de datos en Snowflake

Crear una tarea independiente

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; 
Automatización de canalizaciones de datos en Snowflake

Con warehouse vs sin servidor

Con warehouse

  • Usa un warehouse virtual con nombre
  • Facturación mínima de 60 s por ejecución
CREATE TASK my_task
  WAREHOUSE = harbr_wh  -- named warehouse
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;

Sin servidor (serverless)

  • Omite la cláusula WAREHOUSE
  • Se cobra por segundo consumido; sin coste en inactividad
CREATE TASK my_serverless_task
  -- no WAREHOUSE clause
  SCHEDULE = '5 MINUTE'
AS INSERT INTO ...;
Automatización de canalizaciones de datos en Snowflake

Orquestación de tareas con DAG

Diagrama DAG — ingest_raw (tarea raíz, icono de reloj) → clean_events (corre tras ingest_raw) → build_summary (corre tras clean_events)

  • DAGs: grafos dirigidos acíclicos
  • La tarea raíz define el CRON
  • Las hijas declaran su predecesora con AFTER
  • Cadena completa centralizada
  • Los DAG controlan cuándo y en qué orden se ejecuta
Automatización de canalizaciones de datos en Snowflake

Gestionar estados de tarea

Diagrama de estados — SUSPENDED → STARTED → SUCCEEDED o FAILED

-- Activar una tarea
ALTER TASK mytask RESUME; 
-- Pausar una tarea
ALTER TASK mytask SUSPEND; 
-- Ejecutar una tarea
EXECUTE TASK mytask SUSPEND;
Automatización de canalizaciones de datos en Snowflake

Tareas y streams combinados

Diagrama CDC — delivery_events → stream capta cambios → ¿stream con datos? SÍ → se dispara la tarea → se actualiza delivery_summary / NO → se omite la tarea

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;
Automatización de canalizaciones de datos en Snowflake

¡Vamos a practicar!

Automatización de canalizaciones de datos en Snowflake

Preparing Video For Download...