Saltar al contenido
Lakehouse LabLakehouse LabPreparación Databricks Data Engineer
Módulo 20 · Professional

Jobs avanzado

Contenido abierto

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.

Lectura pública
Al terminar podrás
  • Usar branching y for-each con límites
  • Aplicar retries y repairs correctamente
  • Transferir parámetros y task values
Ver fuentes y revisión

Metadatos editoriales

Última revisión
21 jul 2026
Nivel
Professional
Ruta relacionada
pipelines
Dominios blueprint
Debugging and Deploying · Lakeflow Jobs
Estado
Revisión editorial interna
Fuentes principales
Control the flow of tasks within Lakeflow Jobs · Databricks · Dynamic value references · Databricks
Reportar un error
01
Modelo mental

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
Reportar un error en esta lección

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.

Task key

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.
Camino crítico

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.
Dependencia de efecto

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.
YAMLParalelismo y convergencia explícitos
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 navegador
02
Implementación

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
Reportar un error en esta lección

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.

Job parameter

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

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.
Referencia dinámica

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.
PythonPublicación de evidencia para una condición
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 navegador
03
Operación

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
Reportar un error en esta lección

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.

If/else task

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

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

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.
JSONCondición basada en calidad
{
  "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 navegador
04
Diagnóstico

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
Reportar un error en esta lección

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.

Tarea anidada

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.
Concurrencia del bucle

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.
Aislamiento por elemento

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.
JSONProcesamiento regional acotado
{
  "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 navegador
05
Decisión de diseño

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
Reportar un error en esta lección

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.

Retry

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

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

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.
YAMLPolítica de retry limitada
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_partner

El 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

¿Por qué una tarea que soporta retries debe ser idempotente?

Borrador privado · solo en este navegador
5 lecciones pendientes

Vista de lectura · sin ejecución

Learn Databricks

Learn Databricks · commit 08c378c

README.md