Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.
Lección 4 de 5
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.
- Duración
- 17 min aprox.
- Objetivo
- For each ejecuta una tarea anidada por elemento con concurrencia limitada y exige que cada iteración sea aislable e idempotente.
- Siguiente paso
- Continuar con la siguiente lección
Ver detalles del módulo
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
04Diagnó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.
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