Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.
Professional
Lakeflow Jobs avanzado: control flow y repairs
Orquesta decisiones, bucles y recuperaciones sin convertir el DAG en lógica opaca.
- Usar branching y for-each con límites
- Aplicar retries y repairs correctamente
- Transferir parámetros y task values
01Modelo mentalIf/else task
El DAG de Lakeflow Jobs expresa dependencias de ejecución y permite paralelismo solo cuando tareas y datos son realmente independientes.
+
If/else task
El DAG de Lakeflow Jobs expresa dependencias de ejecución y permite paralelismo solo cuando tareas y datos son realmente independientes.
- Objetivo
- El DAG de Lakeflow Jobs expresa dependencias de ejecución y permite paralelismo solo cuando tareas y datos son realmente independientes.
- Duración estimada
- 17 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m19
Cada task tiene una responsabilidad, parámetros y resultado observable. Dos ingestas regionales pueden ejecutarse en paralelo si escriben particiones o targets independientes; la publicación gold depende de ambas. Introducir dependencias innecesarias alarga el critical path, pero eliminar una dependencia de datos crea carreras.
El DAG no debe ocultar lógica de transformación dentro de docenas de notebooks. Jobs coordina unidades desplegables —pipeline, wheel, SQL o notebook— y las tareas comparten información mediante parámetros, task values o outputs, no variables de memoria del driver.
Modelo mental
Un Lakeflow Job es un grafo de tareas, no una lista visual de notebooks. Cada arista declara una condición de dependencia y el scheduler ejecuta en paralelo únicamente las ramas cuyos prerrequisitos están satisfechos. El DAG debe reflejar dependencias de datos y efectos reales: dos tareas que escriben la misma tabla no son independientes aunque no se lean entre sí, y una arista innecesaria desperdicia paralelismo. La unidad de retry, timeout, compute, parámetros y observabilidad es la tarea; por eso conviene que sea cohesionada e idempotente. Un Job puede orquestar notebooks, scripts Python, pipelines, SQL y otros tipos, pero no convierte su contenido en transaccional de extremo a extremo. La arquitectura separa producir, validar y publicar para que un fallo no exponga datos parciales. Los nombres y task keys son contratos operativos porque aparecen en referencias dinámicas, repair runs, alertas y system tables.
Identificador estable y único de una tarea dentro del Job, utilizado por dependencias, referencias dinámicas, métricas y operaciones de reparación.
Cambiarlo sin planificación puede romper parámetros downstream y comparabilidad histórica aunque el nombre visible parezca equivalente.Secuencia dependiente de tareas cuya duración acumulada determina el tiempo mínimo posible para completar el run completo.
Ayuda a optimizar donde realmente reduce SLA, en lugar de acelerar ramas que ya terminan antes.Relación no visible solo por lecturas, creada cuando tareas compiten por el mismo target, recurso externo o publicación.
Debe representarse o eliminarse mediante aislamiento para impedir carreras y resultados no deterministas.tasks:
- task_key: ingest_eu
python_wheel_task:
package_name: commerce
entry_point: ingest
- task_key: ingest_us
python_wheel_task:
package_name: commerce
entry_point: ingest
- task_key: publish_gold
depends_on:
- task_key: ingest_eu
- task_key: ingest_us
pipeline_task:
pipeline_id: ${resources.pipelines.orders.id}En el bundle real usa la sintaxis de sustitución `${resources.pipelines.orders.id}`; aquí se escapa el signo para mantener el ejemplo como texto.
Puntos clave
- Dependencias representan requisitos de datos/estado, no preferencia visual.
- Tareas independientes pueden usar compute y retries separados.
- El critical path determina la latencia mínima del workflow.
Evita
- Serializar tareas independientes y aumentar tiempo/coste sin mejorar corrección.
- Ejecutar en paralelo tareas que sobrescriben el mismo rango de la tabla.
Recuerdo activo
¿Qué determina el critical path de un Job?
Borrador privado · solo en este navegador02ImplementaciónFor each task
Job parameters, task values y referencias dinámicas trasladan contexto entre tareas sin acoplarlas a estado de notebook.
+
For each task
Job parameters, task values y referencias dinámicas trasladan contexto entre tareas sin acoplarlas a estado de notebook.
- Objetivo
- Job parameters, task values y referencias dinámicas trasladan contexto entre tareas sin acoplarlas a estado de notebook.
- Duración estimada
- 17 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m19
Los job parameters describen entradas del run y se propagan a tareas compatibles. Una tarea puede publicar un valor pequeño mediante `dbutils.jobs.taskValues.set`; downstream lo referencia como `{{tasks.<task>.values.<key>}}`. Las referencias se sustituyen como texto, no evalúan expresiones.
Los task values sirven para contadores, rutas o listas pequeñas, no para transportar DataFrames. Los datos voluminosos se materializan en tablas/Volumes y se pasa un identificador. Un nombre mal escrito puede tratarse como literal, así que la validación del bundle y una prueba smoke son esenciales.
Modelo mental
Los parámetros describen intención de ejecución; los task values transportan resultados pequeños calculados durante esa ejecución; las referencias dinámicas enlazan ambos sin copiar estado a notebooks. Un job parameter como `business_date` o `environment` debe tener tipo y validación conceptual, aunque llegue como texto. Puede propagarse a tareas mediante configuración, mientras `dbutils.jobs.taskValues.set` publica un valor para una tarea downstream o un `If/else`. No es un almacén de datos: payloads grandes pertenecen a tablas, Volumes u object storage y se pasan por referencia. Las referencias `{{...}}` se resuelven por el servicio antes de ejecutar la tarea y algunas no fallan si se escriben mal, por lo que deben revisarse y probarse. Secretos nunca viajan como parámetros visibles. El contrato incluye default seguro, timezone, formato, origen y comportamiento de rerun para que un repair use el mismo intervalo lógico.
Entrada definida en el ámbito del Job que configura una ejecución y puede propagarse de forma consistente a múltiples tareas.
Centraliza fecha, entorno o modo y hace reproducibles runs normales, manuales, backfills y reparaciones.Valor pequeño producido durante una tarea y expuesto por clave a condiciones o tareas posteriores dentro del mismo run.
Permite comunicar decisiones y referencias sin acoplar notebooks a variables globales o archivos temporales implícitos.Plantilla resuelta por Lakeflow Jobs con contexto del run, trigger, parámetros, tareas o metadatos disponibles oficialmente.
Conecta configuración declarativa con valores de ejecución, pero exige sintaxis y alcance verificados para evitar literales accidentales.invalid_rows = spark.table("main.ops.validation_results").where(
"run_id = :run_id AND is_valid = false"
).count()
total_rows = spark.table("main.ops.validation_results").where(
"run_id = :run_id"
).count()
ratio = invalid_rows / total_rows if total_rows else 1.0
dbutils.jobs.taskValues.set(key="invalid_ratio", value=ratio)El downstream referencia `{{tasks.validate.values.invalid_ratio}}`; conserva el detalle de filas en una tabla, no en el task value.
Puntos clave
- Job parameter identifica el run; task value comunica un resultado pequeño de upstream.
- Referencias dinámicas usan doble llave y no ejecutan código.
- Datos grandes se comparten mediante almacenamiento gobernado, no mediante valores del DAG.
Evita
- Pasar miles de filas como JSON en un task value y alcanzar límites de tamaño.
- Usar una referencia dinámica inexistente y no detectar que quedó como texto literal.
Recuerdo activo
¿Cómo consume downstream el valor `invalid_ratio` publicado por `validate`?
Borrador privado · solo en este navegador03OperaciónRun job y modularidad
If/else decide por un valor; Run if decide por el estado de tareas upstream y ambos resuelven problemas distintos.
+
Run job y modularidad
If/else decide por un valor; Run if decide por el estado de tareas upstream y ambos resuelven problemas distintos.
- Objetivo
- If/else decide por un valor; Run if decide por el estado de tareas upstream y ambos resuelven problemas distintos.
- Duración estimada
- 17 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m19
Una If/else task compara parámetros, valores dinámicos o task values con operadores como `>`, `==` o `!=`. Por ejemplo, publica si `invalid_ratio <= 0.01` y envía a cuarentena en caso contrario. Las tareas de cada rama declaran el outcome requerido.
`Run if` se configura sobre una dependencia para ejecutar limpieza, notificación o recuperación según estados como `ALL_SUCCESS`, `AT_LEAST_ONE_FAILED` o `ALL_DONE`. No se debe codificar un valor de negocio como estado de tarea ni usar If/else para saber si upstream lanzó una excepción.
Modelo mental
`If/else` y `Run if` controlan dimensiones diferentes. La tarea If/else compara un valor —parámetro, referencia dinámica o task value— con un operador y abre una rama verdadera o falsa. `Run if` evalúa estados terminales de dependencias, como todos correctos, al menos uno fallido o todos terminados, y decide si una tarea downstream es aplicable. Confundirlos produce DAGs frágiles: comprobar `row_count > 0` es decisión por valor; ejecutar limpieza aunque upstream falle es decisión por estado. Las tareas omitidas adquieren estados que influyen en dependientes, por lo que se diseña y prueba cada ruta, incluida la ausencia de datos. Una rama condicional no reemplaza validación transaccional: publicar porque una bandera dice true requiere confiar en quién calculó esa bandera y conservar evidencia. Las tareas de cleanup usan `All done`, pero deben ser idempotentes y no ocultar el fallo original.
Tarea de control que compara un valor disponible con un operador soportado y habilita una de dos ramas del DAG.
Expresa decisiones de negocio o de datos, como publicar únicamente cuando un conteo supera un umbral.Condición asociada a dependencias que decide ejecución según estados de tareas upstream, incluidos éxito, fallo o finalización.
Permite cleanup, notificación y tolerancia parcial sin convertir estados técnicos en valores inventados.Conjunto de tareas que el scheduler no ejecuta porque una condición eligió otra rama o sus dependencias no aplican.
Debe probarse porque su estado influye en downstream y puede ocultar que nunca se validó una alternativa rara.{
"task_key": "quality_gate",
"depends_on": [{"task_key": "validate"}],
"condition_task": {
"op": "LESS_THAN_OR_EQUAL",
"left": "{{tasks.validate.values.invalid_ratio}}",
"right": "0.01"
}
}Las tareas downstream de publicación o cuarentena dependen de `quality_gate` con el outcome true o false correspondiente.
Puntos clave
- If/else evalúa datos o parámetros; Run if evalúa resultado de ejecución.
- Una tarea de cleanup suele usar `ALL_DONE` para ejecutarse incluso tras fallos.
- Las ramas deben converger con condiciones que acepten outcomes esperados.
Evita
- Usar If/else para capturar fallos técnicos cuando `Run if` ya modela estados de upstream.
- Olvidar una rama o convergencia y dejar el Job aparentemente correcto pero incompleto.
Recuerdo activo
¿Qué condición usarías para liberar un recurso tanto si upstream tuvo éxito como si falló?
Borrador privado · solo en este navegador04DiagnósticoJob repairs y parameter overrides
For each ejecuta una tarea anidada por elemento con concurrencia limitada y exige que cada iteración sea aislable e idempotente.
+
Job repairs y parameter overrides
For each ejecuta una tarea anidada por elemento con concurrencia limitada y exige que cada iteración sea aislable e idempotente.
- Objetivo
- For each ejecuta una tarea anidada por elemento con concurrencia limitada y exige que cada iteración sea aislable e idempotente.
- Duración estimada
- 17 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m19
La lista puede venir de un job parameter, task value o salida SQL y cada elemento se referencia con `{{input...}}`. Procesar regiones o fechas en paralelo reduce latencia, pero la concurrencia debe respetar límites del sistema origen, compute y destino.
Cada iteración escribe una partición o clave independiente. Si todas hacen `overwrite` de la tabla completa, el loop crea carreras. Para miles de elementos, agrupar en rangos o usar una transformación Spark distribuida suele ser mejor que crear miles de task runs.
Modelo mental
La tarea `For each` expande una colección en iteraciones de una tarea anidada y limita cuántas se ejecutan simultáneamente. Es adecuada cuando cada elemento —fecha, región, tabla— puede procesarse de forma aislada e idempotente. No es un sustituto general de paralelismo Spark: lanzar miles de tasks para particiones de un mismo DataFrame añade overhead de scheduler y compute que un único Job distribuido resolvería mejor. La colección debe ser acotada, validada y suficientemente pequeña para los límites de Jobs y referencias dinámicas. Cada iteración recibe el elemento actual y debe escribir a un namespace o clave que evite colisiones. La concurrencia se fija según cuotas de API, capacidad de warehouse y targets, no según el máximo disponible. Si una iteración falla, reparación y reintentos deben poder repetir solo esa unidad sin alterar las ya confirmadas.
Definición ejecutable que For each instancia una vez por elemento, con parámetros resueltos y estado observable para esa iteración.
Concentra la lógica repetible y permite reparar una unidad concreta sin duplicar toda la orquestación.Máximo de iteraciones de For each que Lakeflow Jobs permite ejecutar simultáneamente dentro del run activo.
Protege servicios y compute downstream y determina equilibrio entre duración, coste y riesgo de throttling.Propiedad por la que una iteración lee y escribe recursos identificables sin competir ni depender implícitamente de otra.
Es requisito para paralelismo seguro, idempotencia y reparación selectiva de iteraciones fallidas.{
"task_key": "process_regions",
"for_each_task": {
"inputs": "{{tasks.discover.values.regions}}",
"concurrency": 4,
"task": {
"task_key": "process_region",
"notebook_task": {
"notebook_path": "/Workspace/commerce/process_region",
"base_parameters": {"region": "{{input}}"}
}
}
}
}`regions` debe ser JSON válido y pequeño; el notebook escribe solo el rango de la región recibida.
Puntos clave
- For each contiene exactamente una tarea anidada que recibe el elemento actual.
- Concurrency limita iteraciones simultáneas y protege dependencias externas.
- El cuerpo debe poder reintentarse por elemento sin duplicar efectos.
Evita
- Configurar concurrencia igual al número de elementos y saturar la API o base fuente.
- Usar For each para millones de filas que Spark puede procesar en un único DataFrame distribuido.
Recuerdo activo
¿Qué hace segura una reparación parcial de un For each?
Borrador privado · solo en este navegador05Decisión de diseñoConcurrency y queueing
Retries corrigen fallos transitorios de una tarea; repair runs reejecutan el subconjunto fallido tras corregir la causa sin repetir trabajo exitoso.
+
Concurrency y queueing
Retries corrigen fallos transitorios de una tarea; repair runs reejecutan el subconjunto fallido tras corregir la causa sin repetir trabajo exitoso.
- Objetivo
- Retries corrigen fallos transitorios de una tarea; repair runs reejecutan el subconjunto fallido tras corregir la causa sin repetir trabajo exitoso.
- Duración estimada
- 17 min aprox.
- Dificultad
- Professional
- Prerrequisitos
- m19
Una retry policy limita intentos e intervalo y debe reservarse para fallos plausiblemente transitorios. Reintentar una violación de esquema determinista solo consume tiempo. Las tareas con efectos deben ser idempotentes porque tanto retry como repair pueden ejecutarlas de nuevo.
Un repair run mantiene el contexto del run original y permite reejecutar tareas fallidas o omitidas, con parámetros corregidos cuando corresponda. Antes de reparar se inspecciona error, input y versión; después se valida el resultado downstream. Una nueva ejecución completa es preferible si cambió el alcance o no puede garantizarse coherencia con tareas exitosas anteriores.
Modelo mental
Un retry repite automáticamente una tarea ante un fallo que se presume transitorio; un repair run se inicia después para reejecutar tareas fallidas o omitidas y sus dependientes necesarios dentro de un run existente. Ninguno corrige lógica no idempotente. La política de retry especifica número, intervalo y, cuando aplica, backoff; debe ser corta para errores de red o capacidad recuperables y evitar tormentas sobre un servicio caído. Un error de schema, permiso o calidad determinista no mejora al repetir y consume tiempo de RTO. Repair conserva el contexto y parámetros del run original, por lo que es preferible a lanzar manualmente otro Job que pueda usar otra fecha. Antes de reparar se corrige la causa, se entiende qué outputs quedaron confirmados y se selecciona el mínimo subgrafo seguro. La publicación y los efectos externos necesitan claves que toleren repetición.
Nuevo intento automático de la misma tarea y contexto después de un fallo, sujeto a límites e intervalos configurados.
Recupera errores transitorios sin intervención, pero exige idempotencia y clasificación para no repetir fallos deterministas.Reanudación explícita de un run existente que reejecuta el subconjunto fallido o dependiente conservando su contexto original.
Reduce trabajo duplicado y mantiene la fecha y parámetros lógicos después de corregir una causa.Estrategia que aumenta el intervalo entre intentos, normalmente con aleatoriedad, para reducir presión sobre una dependencia degradada.
Evita tormentas coordinadas de reintentos y mejora probabilidad de recuperación de servicios con throttling.task_key: ingest_partner_api
max_retries: 3
min_retry_interval_millis: 60000
retry_on_timeout: true
timeout_seconds: 1800
python_wheel_task:
package_name: commerce
entry_point: ingest_partnerEl código debe escribir con una clave de petición o `MERGE`; tres reintentos de un append no idempotente pueden triplicar datos.
Puntos clave
- Retry es automática y cercana al fallo; repair es una decisión sobre una ejecución existente.
- No todos los errores son transitorios ni deben reintentarse.
- Idempotencia permite reejecutar una tarea sin duplicar o retroceder estado.
Evita
- Configurar retries ilimitados para un error de datos permanente y ocultar el incidente.
- Reparar downstream con parámetros distintos sin verificar que los outputs upstream exitosos siguen siendo compatibles.
Recuerdo activo