Puedes leer sin crear un espacio. Créalo solo cuando quieras guardar.
Guardar progresoLakeflow Jobs: DAG, tareas y triggers
Convierte ejecuciones manuales en workflows parametrizados, idempotentes y observables.
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.
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.
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.
¿Qué permite ejecutar dos tareas a la vez?
Profundiza
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.Resumen
Puntos clave
- DAG dirige orden
- Tareas independientes paralelizan
- Cada tarea debe ser reintentable
Evita
- DAG totalmente serial
- Estado oculto entre notebooks
02Implementació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.
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.
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.
¿Es taskValues apropiado para un millón de filas?
Profundiza
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.Resumen
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
03Operació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.
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.
task_key: publish
max_retries: 2
min_retry_interval_millis: 60000
timeout_seconds: 1800Clasifica errores permanentes para no reintentarlos inútilmente.
¿Qué debe comprobarse antes de activar retry?
Profundiza
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.Resumen
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
04Diagnó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.
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.
schedule:
quartz_cron_expression: '0 0 2 * * ?'
timezone_id: UTC
pause_status: UNPAUSEDDocumenta cómo afecta horario de verano si no usas UTC.
¿Qué trigger reduce runs vacíos cuando llegan archivos irregularmente?
Profundiza
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.Resumen
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
05Decisió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.
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.
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.
¿Qué evita repetir tareas ya correctas?
Profundiza
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.Resumen
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