Puedes leer sin crear un espacio. Créalo solo cuando quieras guardar.
Guardar progresoLecció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.
17 min aprox.
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.
{
"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.
¿Qué hace segura una reparación parcial de un For each?
Profundiza
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.Resumen
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.