Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.
Associate + Professional
Lakeflow Jobs: DAG, tareas y triggers
Convierte ejecuciones manuales en workflows parametrizados, idempotentes y observables.
- Diseñar un DAG con paralelismo seguro
- Configurar tareas, parámetros y dependencias
- Elegir triggers temporales o dirigidos por datos
01Modelo mentalJobs, runs y task graph
Lakeflow Jobs expresa workflows como DAG de tareas con dependencias observables.
+
Jobs, runs y task graph
Lakeflow Jobs expresa workflows como DAG de tareas con dependencias observables.
- Objetivo
- Lakeflow Jobs expresa workflows como DAG de tareas con dependencias observables.
- Duración estimada
- 17 min aprox.
- Dificultad
- Associate + Professional
- Prerrequisitos
- m09
Cada tarea define tipo, compute, identidad y resultado; dependencias permiten paralelismo cuando no existe relación de datos.
Divide por unidades reintentables e idempotentes, no por cada línea de código.
Modelo mental
Lakeflow Jobs modela un workflow como un grafo dirigido acíclico de tareas. Cada nodo tiene tipo, compute, parámetros, identidad, timeout y política de retry; cada arista expresa una dependencia y condición de ejecución. El DAG no transporta automáticamente DataFrames ni garantiza idempotencia: coordina unidades que deben publicar resultados durables o valores pequeños. Las dependencias permiten paralelismo cuando no existe relación causal y bloquean downstream cuando una precondición falla. Un buen grafo refleja fronteras operativas: una tarea debe poder reintentarse, observarse y repararse sin repetir todo. El nombre Lakeflow Jobs es la superficie de orquestación; no debe confundirse con Spark Declarative Pipelines, framework para declarar datasets.
Grafo dirigido sin ciclos que representa tareas y dependencias causales dentro de un workflow ejecutable.
Permite planificar paralelismo, propagación de fallos y recuperación selectiva de manera explícita.Unidad observable de trabajo con entrada, compute, parámetros, resultado y política operacional propios dentro de un Job.
Define el nivel al que se reintenta, alerta, mide y repara una parte del proceso.Relación que condiciona la elegibilidad de una tarea al estado o salida de otra tarea anterior.
Evita ejecutar consumidores antes de que sus precondiciones durables y verificadas estén disponibles.tasks:
- task_key: bronze
- task_key: silver
depends_on: [{task_key: bronze}]
- task_key: gold
depends_on: [{task_key: silver}]Asigna compute y código en la definición completa.
Puntos clave
- DAG dirige orden
- Tareas independientes paralelizan
- Cada tarea debe ser reintentable
Evita
- DAG totalmente serial
- Estado oculto entre notebooks
Recuerdo activo
¿Qué permite ejecutar dos tareas a la vez?
Borrador privado · solo en este navegador02ImplementaciónNotebook, Python, SQL y pipeline tasks
Parámetros configuran runs y task values intercambian valores pequeños entre tareas.
+
Notebook, Python, SQL y pipeline tasks
Parámetros configuran runs y task values intercambian valores pequeños entre tareas.
- Objetivo
- Parámetros configuran runs y task values intercambian valores pequeños entre tareas.
- Duración estimada
- 17 min aprox.
- Dificultad
- Associate + Professional
- Prerrequisitos
- m09
Los job parameters pueden propagarse y los task parameters alimentan notebooks, Python o SQL.
taskValues transporta IDs, fechas o conteos; las tablas y archivos transportan datasets.
Modelo mental
Los parámetros configuran una ejecución antes o al iniciar tareas; los task values comunican pequeños valores calculados durante el run. Ninguno debe transportar datasets. Una fecha, ruta lógica, umbral o identificador de versión cabe en el contrato; millones de filas pertenecen a una tabla gobernada o un Volume. Los parámetros de Job centralizan valores compartidos y las referencias dinámicas aportan contexto como start time o salidas upstream. Un task value tiene clave, productor y consumidor explícitos y está sujeto a límites de tamaño. Convertir todos los valores a tipos de dominio y validar rangos evita que una cadena vacía cambie silenciosamente la partición procesada.
Valor declarado a nivel de workflow que configura un run y puede propagarse de manera consistente a varias tareas.
Evita editar código por fecha o ambiente y deja la configuración efectiva registrada junto con la ejecución.Par clave-valor pequeño publicado por una tarea para control o parametrización de otras tareas del mismo run.
Permite comunicar métricas y decisiones sin usar la orquestación como transporte de datasets completos.Referencia resuelta por Jobs a metadatos de ejecución o salidas disponibles cuando una tarea se prepara.
Conecta contexto y control flow de forma auditable sin valores copiados manualmente entre notebooks.dbutils.jobs.taskValues.set(key="validated_rows", value=validated.count())
rows = dbutils.jobs.taskValues.get(taskKey="validate", key="validated_rows")No uses count en definiciones declarativas de pipeline; este ejemplo es tarea Job clásica.
Puntos clave
- Parámetros son configuración
- Task values son pequeños
- Datos pasan por storage
Evita
- Pasar secretos como parámetro
- Serializar DataFrame en task value
Recuerdo activo
¿Es taskValues apropiado para un millón de filas?
Borrador privado · solo en este navegador03OperaciónParámetros y task values
Retries, if/else y for-each modelan recuperación y control flow sin duplicar lógica.
+
Parámetros y task values
Retries, if/else y for-each modelan recuperación y control flow sin duplicar lógica.
- Objetivo
- Retries, if/else y for-each modelan recuperación y control flow sin duplicar lógica.
- Duración estimada
- 17 min aprox.
- Dificultad
- Associate + Professional
- Prerrequisitos
- m09
Retry atiende fallos transitorios si la tarea es idempotente; if/else evalúa una condición y for-each repite sobre una lista acotada.
Una rama no sustituye validación de datos y un bucle masivo puede crear demasiadas tareas.
Modelo mental
Retries recuperan fallos transitorios; if/else elige una rama por una condición; for each repite una tarea parametrizada sobre una colección acotada. Son primitivas de control flow, no sustitutos de lógica de datos ni de un diseño idempotente. Un retry seguro presupone que volver a ejecutar la tarea no duplica efectos. La condición debe depender de un valor pequeño y estable, no del estado invisible de una sesión. Un loop necesita límite de concurrencia, identidad por elemento y estrategia para fallos parciales. Modelar estos caminos en Jobs hace visibles intentos y ramas; esconderlos en un gran notebook reduce observabilidad y obliga a repetir trabajo ya correcto.
Nuevo intento automático de una tarea fallida bajo una política de número, intervalo y timeout determinada.
Solo es seguro cuando los efectos de la tarea son idempotentes y el error puede ser transitorio.Nodo de control que evalúa dos operandos y determina qué dependencias de outcome quedan habilitadas.
Hace visible una decisión operacional como publicar, cuarentenizar o detener según una métrica pequeña.Control que ejecuta una tarea anidada por cada elemento de una colección parametrizada con concurrencia limitada.
Simplifica backfills o fan-out acotados sin confundir orquestación con paralelismo de filas de Spark.task_key: publish
max_retries: 2
min_retry_interval_millis: 60000
timeout_seconds: 1800Clasifica errores permanentes para no reintentarlos inútilmente.
Puntos clave
- Retry exige idempotencia
- If/else decide rutas
- For-each requiere límites
Evita
- Retry de append no idempotente
- For-each con miles de elementos
Recuerdo activo
¿Qué debe comprobarse antes de activar retry?
Borrador privado · solo en este navegador04DiagnósticoSchedule, file arrival y table update
Schedule, file arrival y table update disparan Jobs por tiempo o disponibilidad del dato.
+
Schedule, file arrival y table update
Schedule, file arrival y table update disparan Jobs por tiempo o disponibilidad del dato.
- Objetivo
- Schedule, file arrival y table update disparan Jobs por tiempo o disponibilidad del dato.
- Duración estimada
- 17 min aprox.
- Dificultad
- Associate + Professional
- Prerrequisitos
- m09
Schedule usa calendario y zona horaria; file arrival observa nuevas llegadas; table update reacciona a cambios de tablas compatibles.
Elige señal de datos cuando evita espera o runs vacíos; usa calendario cuando la obligación es temporal.
Modelo mental
Los triggers responden a dos clases de obligación: tiempo y disponibilidad de datos. Un schedule declara cuándo debe comenzar un run según cron y zona horaria. File arrival reacciona a nuevos archivos en una ubicación soportada y evita polling frecuente. Table update responde a actualizaciones de tablas compatibles y puede coordinar downstream a partir de cambios publicados. Elegir data-driven reduce ejecuciones vacías y latencia cuando la llegada es irregular; elegir calendario es correcto cuando el compromiso es un cierre temporal aunque no haya datos nuevos. Ningún trigger prueba que la entrada sea completa: el Job todavía necesita validación, idempotencia, concurrencia y política para múltiples eventos.
Disparador temporal que crea runs mediante una expresión cron interpretada en una zona horaria declarada.
Es apropiado cuando la obligación de proceso o cierre depende del reloj y no solo de datos nuevos.Disparador dirigido por datos que inicia un Job al detectar archivos nuevos en una ubicación compatible.
Reduce polling y runs vacíos para fuentes de archivos con llegadas irregulares o impredecibles.Disparador que responde a actualizaciones observadas en una o varias tablas soportadas por la plataforma.
Permite desacoplar productor y consumidor usando una publicación gobernada como señal operacional.schedule:
quartz_cron_expression: '0 0 2 * * ?'
timezone_id: UTC
pause_status: UNPAUSEDDocumenta cómo afecta horario de verano si no usas UTC.
Puntos clave
- Cron depende de zona horaria
- File arrival es data-driven
- Table update sigue tablas
Evita
- Cron sin timezone
- File trigger sobre ruta inestable
Recuerdo activo
¿Qué trigger reduce runs vacíos cuando llegan archivos irregularmente?
Borrador privado · solo en este navegador05Decisión de diseñoRetries, timeout y notificaciones
Run history, repairs y alertas convierten un DAG fallido en una recuperación auditable.
+
Retries, timeout y notificaciones
Run history, repairs y alertas convierten un DAG fallido en una recuperación auditable.
- Objetivo
- Run history, repairs y alertas convierten un DAG fallido en una recuperación auditable.
- Duración estimada
- 17 min aprox.
- Dificultad
- Associate + Professional
- Prerrequisitos
- m09
Repair run reejecuta tareas fallidas y dependientes sin repetir éxitos innecesarios; parameter override permite corregir una partición.
Las alertas deben indicar owner, run y acción. Spark UI atiende ejecución; run history atiende tendencia y estado del workflow.
Modelo mental
Run history explica el estado del workflow; Spark UI explica la ejecución distribuida dentro de tareas Spark; alertas movilizan a un owner; repair run recupera un subconjunto conservando éxitos previos. Estas superficies responden preguntas diferentes y deben formar un runbook. Primero se identifica job, run, task e intento; después se clasifica si el fallo es de orquestación, dependencia, datos, compute o plan. Repair no corrige una causa y no debe pulsarse antes de verificar idempotencia. Las system tables de Lakeflow permiten analizar tendencias y SLA más allá de la retención visual disponible. Una alerta útil incluye impacto, enlace, owner y primera acción, no solo el texto FAILED.
Reejecución selectiva de tareas fallidas o elegidas y de sus dependientes dentro de un run existente.
Reduce tiempo y riesgo al conservar trabajo exitoso, siempre que las fronteras de tareas sean idempotentes.Registro de ejecuciones, estados, parámetros, intentos y tiempos de Jobs disponible para diagnóstico y auditoría.
Localiza dónde falló el workflow antes de investigar detalles de compute o datos.Procedimiento operativo versionado que conecta síntomas, evidencia, owner, decisión de recuperación y criterios de cierre.
Hace que una alerta se convierta en una respuesta consistente en lugar de improvisación personal.SELECT job_id, run_id, result_state, period_start_time
FROM system.lakeflow.job_run_timeline
ORDER BY period_start_time DESC;Une jobs para obtener nombre y segmenta por workspace.
Puntos clave
- Repair preserva éxitos
- Override corrige el run
- Alertas deben ser accionables
Evita
- Relanzar todo tras una única tarea fallida
- Alerta sin enlace al run
Recuerdo activo