Módulo 4: Fan-out paralelo y agregación
Concurrencia real con ThreadPoolExecutor
Descripción
run_fanout_sequential (lecciones 03-04) corre cada sub-tarea una después de la otra — el orden es
fijo y determinista, pero el tiempo total es la suma de las dos. Esta lección construye la extensión
que la Regla dura del módulo anticipó desde la lección 01: run_fanout_parallel, que despacha cada
sub-tarea en su propio hilo, con concurrent.futures.ThreadPoolExecutor — el mismo módulo que
ya usa dispatch_parallel desde agent-fundamentals, aplicado ahora a agentes completos en vez de
a tool calls sueltas.
La parte que exige más disciplina es la que da nombre a la lección: con hilos reales, el orden en
que terminan no es determinista — depende del sistema operativo, de la carga de la máquina, de
factores que este código no controla ni debería intentar controlar. Esta lección lo demuestra
ejecutando el mismo fan-out varias veces y mostrando que el diccionario results se llena en un
orden distinto cada vez — y, aun así, la salida impresa, la que cualquier código consume, es
exactamente la misma en todas las corridas, porque nunca se construye a partir de ese orden de
llegada.
Conexión con el módulo
Esta lección extiende run_fanout_sequential de las lecciones 03-04 con run_fanout_parallel, sin
cambiar SubTask ni la forma de results. La lección 07 toma ambas versiones —secuencial y
paralela— y mide, con números, cuánto ahorra la concurrencia real en rondas de coordinación.
Analogía: dos mensajeros, cada uno con su propio camino
Vuelve a los dos mandados: comprar pan, revisar el correo. Hasta la lección 04, una sola persona hacía los dos, uno después del otro. Esta lección manda a dos mensajeros, cada uno por su cuenta, al mismo tiempo. Ninguno de los dos sabe ni le importa si el otro ya terminó — cada uno sigue su propio camino, a su propio ritmo, y puede llegar antes o después según el tráfico de ese día en particular. Lo que sí es fijo, sin importar quién llegue primero, es cómo se reporta el resultado al final: siempre "pan, después correo" en ese orden en el reporte — nunca "lo que haya llegado primero, primero".
Ejemplo trabajado: run_fanout_parallel, y la prueba de determinismo bajo hilos reales
import concurrent.futures
from dataclasses import dataclass
import reservo_tools as rt
def dispatch_parallel(tool_use_blocks, tools):
with concurrent.futures.ThreadPoolExecutor(max_workers=len(tool_use_blocks)) as pool:
futures = [pool.submit(tools[b["name"]], **b["input"]) for b in tool_use_blocks]
results = [f.result() for f in futures]
return [
{"type": "tool_result", "tool_use_id": b["id"], "content": str(r)}
for b, r in zip(tool_use_blocks, results)
]
def run_agent_parallel(question, model_script, tools, max_iterations=10):
messages = [{"role": "user", "content": question}]
for step in range(max_iterations):
turn = model_script[step]
messages.append({"role": "assistant", "content": turn["content"]})
if turn["stop_reason"] != "tool_use":
return turn, messages
tool_result_blocks = dispatch_parallel(turn["content"], tools)
messages.append({"role": "user", "content": tool_result_blocks})
raise RuntimeError(f"max_iterations alcanzado ({max_iterations})")
POLICY_DOCS = {
"cancellation-policy": (
"Las reservas se pueden cancelar sin cargo hasta 2 horas antes del "
"horario reservado. Cancelaciones dentro de esas 2 horas aplican "
"el cargo de no-presentación."
),
}
def search_docs(query):
q = query.lower()
if "cancela" in q:
return f"[cancellation-policy] {POLICY_DOCS['cancellation-policy']}"
return "No se encontró una política relevante para esa pregunta."
SPECIALISTS = {
"booking_agent": {
"tools": {
"list_rooms": rt.list_rooms, "get_quote": rt.get_quote,
"book_room": rt.book_room, "cancel_booking": rt.cancel_booking,
},
},
"policy_agent": {"tools": {"search_docs": search_docs}},
}
def run_specialist(name, task, model_script):
tools = SPECIALISTS[name]["tools"]
return run_agent_parallel(task, model_script, tools)
@dataclass
class SubTask:
agent: str
task: str
def run_fanout_sequential(subtasks, model_scripts):
"""El de las lecciones 03/04, sin cambios -- el caso base."""
ordered = sorted(subtasks, key=lambda s: s.agent)
results = {}
for sub in ordered:
final, history = run_specialist(sub.agent, sub.task, model_scripts[sub.agent])
results[sub.agent] = {
"task": sub.task,
"output": final["content"][0]["text"],
"history": history,
}
return results
def run_fanout_parallel(subtasks, model_scripts):
"""La extensión de esta lección: cada sub-tarea corre en su PROPIO hilo,
con ThreadPoolExecutor -- el mismo módulo que ya usa dispatch_parallel,
aplicado ahora a agentes completos en vez de a tool calls sueltas.
El orden en que los hilos TERMINAN no es determinista -- depende del
sistema operativo, no de este código. Por eso `results` se llena a
medida que cada hilo termina, PERO nunca se imprime ni se compone nada
directamente desde ese orden de llegada: `results` es un dict keyed
por nombre de agente, y quien lo consume SIEMPRE itera con
`sorted(results)` -- la misma disciplina de determinismo de la
lección 03, ahora bajo concurrencia real."""
results = {}
with concurrent.futures.ThreadPoolExecutor(max_workers=len(subtasks)) as pool:
future_to_subtask = {
pool.submit(run_specialist, sub.agent, sub.task, model_scripts[sub.agent]): sub
for sub in subtasks
}
for future in concurrent.futures.as_completed(future_to_subtask):
sub = future_to_subtask[future]
final, history = future.result()
results[sub.agent] = {
"task": sub.task,
"output": final["content"][0]["text"],
"history": history,
}
return results
COMPOUND_REQUEST = "Cotiza Focus pro 3h y dime la política de cancelación."
subtasks = [
SubTask(agent="booking_agent", task="Cotiza Focus pro 3h."),
SubTask(agent="policy_agent", task="¿Cuál es la política de cancelación?"),
]
model_scripts = {
"booking_agent": [
{"stop_reason": "tool_use", "content": [
{"type": "tool_use", "id": "toolu_01", "name": "get_quote",
"input": {"room": "Focus", "tier": "pro", "hours": 3}}]},
{"stop_reason": "end_turn", "content": [
{"type": "text", "text": "Focus pro 3h cuesta 6000 centavos."}]},
],
"policy_agent": [
{"stop_reason": "tool_use", "content": [
{"type": "tool_use", "id": "toolu_01", "name": "search_docs",
"input": {"query": "política de cancelación"}}]},
{"stop_reason": "end_turn", "content": [
{"type": "text", "text": (
"Puedes cancelar sin cargo hasta 2 horas antes del horario "
"reservado. Después de eso aplica el cargo de no-presentación."
)}]},
],
}
print("--- fan-out con concurrencia real (ThreadPoolExecutor) ---")
results_parallel = run_fanout_parallel(subtasks, model_scripts)
print("claves en el dict, según llegaron (informativo, NO se usa para imprimir):")
print(" ", list(results_parallel.keys()))
print()
print("salida SIEMPRE ordenada por nombre de agente, sin importar qué hilo terminó primero:")
for agent in sorted(results_parallel):
print(f"[{agent}] {results_parallel[agent]['output']}")
Qué esperar (corrida real; el orden de llegada de las claves puede variar entre ejecuciones):
--- fan-out con concurrencia real (ThreadPoolExecutor) ---
claves en el dict, según llegaron (informativo, NO se usa para imprimir):
['policy_agent', 'booking_agent']
salida SIEMPRE ordenada por nombre de agente, sin importar qué hilo terminó primero:
[booking_agent] Focus pro 3h cuesta 6000 centavos.
[policy_agent] Puedes cancelar sin cargo hasta 2 horas antes del horario reservado. Después de eso aplica el cargo de no-presentación.
Fíjate en el orden de las claves: en esta corrida, policy_agent llegó antes que
booking_agent en el diccionario — el hilo de policy_agent terminó primero. En otra corrida,
podría ser al revés. Ese orden de llegada se imprime solo como información, para que veas que
sí varía — pero la salida real, la que cualquier consumidor de este código usaría, viene de
sorted(results_parallel), que siempre produce booking_agent antes que policy_agent,
alfabéticamente, sin importar qué hilo ganó la carrera.
Probando el determinismo: la misma corrida, varias veces
print()
print("--- confirmando que sequential y parallel producen el MISMO resultado agregado ---")
results_sequential = run_fanout_sequential(subtasks, model_scripts)
same_output = all(
results_sequential[a]["output"] == results_parallel[a]["output"]
for a in results_sequential
)
same_keys = sorted(results_sequential) == sorted(results_parallel)
print("mismas claves (agentes):", same_keys)
print("mismo contenido agregado:", same_output)
print()
print("--- corriendo el fan-out paralelo varias veces: la salida IMPRESA nunca cambia ---")
for run_n in range(1, 4):
r = run_fanout_parallel(subtasks, model_scripts)
ordered_outputs = [r[a]["output"] for a in sorted(r)]
print(f"corrida {run_n}: {ordered_outputs == [results_parallel['booking_agent']['output'], results_parallel['policy_agent']['output']]}")
Qué esperar:
--- confirmando que sequential y parallel producen el MISMO resultado agregado ---
mismas claves (agentes): True
mismo contenido agregado: True
--- corriendo el fan-out paralelo varias veces: la salida IMPRESA nunca cambia ---
corrida 1: True
corrida 2: True
corrida 3: True
Tres corridas independientes de run_fanout_parallel, con hilos reales, cada una potencialmente con
un orden de finalización distinto por debajo — y las tres producen exactamente la misma salida
ordenada. Esa es la prueba, ejecutada varias veces (no una sola), de que la Regla dura de este
módulo se cumple: la salida agregada nunca depende de qué hilo terminó primero.
Por qué results se llena con as_completed, y se lee con sorted
Vale la pena ser preciso sobre dónde vive cada responsabilidad. concurrent.futures.as_completed
devuelve cada future en el orden en que termina — es la forma correcta y eficiente de recoger
resultados de hilos que corren a velocidades distintas, sin bloquear esperando al primero que se
lanzó si otro termina antes. Usar as_completed en vez de simplemente future.result() sobre cada
future en el orden en que se enviaron no es un error a corregir — es la forma idiomática de
ThreadPoolExecutor para no desperdiciar tiempo esperando innecesariamente.
El punto crítico es que ese orden de llegada nunca se propaga más allá del for que llena
results. Una vez que results[sub.agent] = {...} termina de ejecutarse para las dos sub-tareas,
el diccionario existe completo, y cualquier código que lo consuma después —imprimir, componer una
respuesta, contar costo— lee ese diccionario con sorted(), sin ninguna referencia a en qué orden
se llenó.
Errores comunes
-
Usar las claves de
resultsen el orden en que Python las itera por defecto, sinsorted(). Un diccionario de Python preserva el orden de inserción desde la versión 3.7 — pero ese orden de inserción, bajoThreadPoolExecutor, es el orden de finalización de los hilos, no ningún orden con significado. Iterarfor agent in resultssinsorted()reintroduciría exactamente el no-determinismo que esta lección elimina. -
Pensar que hace falta un
Locko algún mecanismo de sincronización manual para escribir enresults. No en este caso: cada hilo escribe en una clave distinta del diccionario (sub.agentes único por sub-tarea), así que no hay dos hilos escribiendo el mismo dato al mismo tiempo. Si dos sub-tareas compartieran el mismoagent, sí haría falta revisar esto con cuidado — un caso que este módulo no construye. -
Afirmar en un comentario o un log "el agente X terminó primero" como si fuera parte del comportamiento esperado del sistema. La Regla dura de este módulo lo prohíbe explícitamente — esa información existe solo para observar que el mecanismo es robusto, nunca como parte de la lógica ni de la respuesta final.
-
Confundir "concurrencia real" con "más rápido en este ejemplo". Con guiones concepto que no hacen ninguna llamada de red ni cómputo pesado, la diferencia de tiempo real entre secuencial y paralelo en esta guía es despreciable — el ahorro que importa es conceptual (rondas de coordinación, lección 07), no un cronómetro corriendo sobre este código de ejemplo.
-
Olvidar
max_workers.ThreadPoolExecutor(max_workers=len(subtasks))garantiza que cada sub-tarea tenga su propio hilo disponible de inmediato. Con un pool más chico, algunas sub-tareas esperarían a que se libere un hilo — el Ejercicio 3 de esta lección confirma que el resultado agregado sigue siendo idéntico de todos modos, aunque el paralelismo real sea menor.
Ejercicios
Ejercicio 1: Confirma tu propia ejecución, varias veces (Fácil)
Ejecuta run_fanout_parallel sobre subtasks cinco veces seguidas, en el mismo proceso. Confirma
que la salida ordenada ([results[a]["output"] for a in sorted(results)]) es idéntica en las cinco
corridas, aunque el orden de las claves en results.keys() pueda variar.
Ver solución
No hay una única "solución de código" para este ejercicio — es una verificación: si las cinco
corridas producen la misma salida ordenada (aunque results.keys() varíe entre ellas), confirmaste
con tus propios ojos la garantía de determinismo de esta lección.
Ejercicio 2: Fan-out real de tres agentes (Medio)
Extiende el ejemplo a tres sub-tareas —booking_agent cotizando Focus pro 2h, policy_agent
respondiendo sobre la política de cancelación, pricing_agent comparando Studio y Boardroom pro
3h— y despáchalas con run_fanout_parallel. Confirma que la salida queda ordenada
alfabéticamente por agente.
Ver solución
SPECIALISTS["pricing_agent"] = {"tools": {"get_quote": rt.get_quote}}
subtasks_3 = [
SubTask(agent="booking_agent", task="Cotiza Focus pro 2h."),
SubTask(agent="policy_agent", task="¿Cuál es la política de cancelación?"),
SubTask(agent="pricing_agent", task="Compara Studio y Boardroom, pro, 3h."),
]
model_scripts_3 = {
"booking_agent": [
{"stop_reason": "tool_use", "content": [
{"type": "tool_use", "id": "toolu_01", "name": "get_quote",
"input": {"room": "Focus", "tier": "pro", "hours": 2}}]},
{"stop_reason": "end_turn", "content": [{"type": "text", "text": "Focus pro 2h cuesta 4000 centavos."}]},
],
"policy_agent": [
{"stop_reason": "tool_use", "content": [
{"type": "tool_use", "id": "toolu_01", "name": "search_docs",
"input": {"query": "política de cancelación"}}]},
{"stop_reason": "end_turn", "content": [
{"type": "text", "text": "Puedes cancelar sin cargo hasta 2 horas antes del horario reservado."}]},
],
"pricing_agent": [
{"stop_reason": "tool_use", "content": [
{"type": "tool_use", "id": "toolu_01", "name": "get_quote",
"input": {"room": "Studio", "tier": "pro", "hours": 3}},
{"type": "tool_use", "id": "toolu_02", "name": "get_quote",
"input": {"room": "Boardroom", "tier": "pro", "hours": 3}},
]},
{"stop_reason": "end_turn", "content": [
{"type": "text", "text": "Studio pro 3h: 9600 centavos. Boardroom pro 3h: 19200 centavos. Studio es más barata."}]},
],
}
results_3 = run_fanout_parallel(subtasks_3, model_scripts_3)
for agent in sorted(results_3):
print(f"[{agent}] {results_3[agent]['output']}")
print("Focus pro 2h a mano:", 2500 * 2 * 80 // 100)
Salida esperada:
[booking_agent] Focus pro 2h cuesta 4000 centavos.
[policy_agent] Puedes cancelar sin cargo hasta 2 horas antes del horario reservado.
[pricing_agent] Studio pro 3h: 9600 centavos. Boardroom pro 3h: 19200 centavos. Studio es más barata.
Focus pro 2h a mano: 4000
Explicación: con tres sub-tareas, el mecanismo no cambia — max_workers=3 (uno por sub-tarea)
y la misma disciplina de sorted(results) para imprimir. 4000 = 2500 * 2 * 80 // 100.
Ejercicio 3: max_workers=1 (sin paralelismo real) da la misma salida (Difícil)
Ejecuta el fan-out de tres agentes del Ejercicio 2, pero con max_workers=1 — forzando a que los
tres hilos se turnen un único worker, en la práctica sin paralelismo real—. Confirma que la salida
agregada es idéntica a la del max_workers=3, y explica por qué eso confirma que el determinismo de
la salida no depende del grado de paralelismo real logrado.
Ver solución
def run_fanout_parallel_workers(subtasks, model_scripts, max_workers):
results = {}
with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as pool:
future_to_subtask = {
pool.submit(run_specialist, sub.agent, sub.task, model_scripts[sub.agent]): sub
for sub in subtasks
}
for future in concurrent.futures.as_completed(future_to_subtask):
sub = future_to_subtask[future]
final, history = future.result()
results[sub.agent] = {"task": sub.task, "output": final["content"][0]["text"], "history": history}
return results
results_seq_via_pool = run_fanout_parallel_workers(subtasks_3, model_scripts_3, max_workers=1)
same = all(results_3[a]["output"] == results_seq_via_pool[a]["output"] for a in results_3)
print("salida con max_workers=3 == salida con max_workers=1:", same)
for agent in sorted(results_seq_via_pool):
print(f"[{agent}] {results_seq_via_pool[agent]['output']}")
Salida esperada:
salida con max_workers=3 == salida con max_workers=1: True
[booking_agent] Focus pro 2h cuesta 4000 centavos.
[policy_agent] Puedes cancelar sin cargo hasta 2 horas antes del horario reservado.
[pricing_agent] Studio pro 3h: 9600 centavos. Boardroom pro 3h: 19200 centavos. Studio es más barata.
Explicación: con max_workers=1, el pool solo tiene un hilo disponible, así que las tres
sub-tareas terminan corriendo, en los hechos, una después de la otra dentro del pool —sin ningún
paralelismo real de por medio—. Aun así, results queda exactamente igual, porque el mecanismo de
agregación (sorted(results)) nunca dependió de cuánto paralelismo real hubo, solo de que las
claves del diccionario terminen completas antes de leerlas. Esta es la prueba más fuerte de esta
lección: el determinismo de la salida es una propiedad del diseño del código (agregar por clave,
leer ordenado), no un efecto colateral de cuántos hilos corrieron de verdad al mismo tiempo.
Resumen y siguiente paso
run_fanout_paralleldespacha cada sub-tarea en su propio hilo (ThreadPoolExecutor), recogiendo resultados conas_completeda medida que terminan.- El orden de llegada al diccionario
resultsvaría entre corridas — confirmado ejecutando el mismo fan-out varias veces— pero la salida impresa, siempre construida consorted(results), es idéntica en todas las corridas. - La Regla dura de este módulo se cumple: en ningún momento se afirma ni se usa "qué agente terminó primero" — solo el contenido agregado, reproducible, importa.
- El determinismo no depende del grado de paralelismo real logrado — confirmado con
max_workers=1, donde el resultado agregado es idéntico al demax_workers=3.
Siguiente lección: 07 — Midiendo el ahorro de latencia en rondas. Comparamos, con números
reales, cuántas rondas de coordinación ahorra run_fanout_parallel frente a
run_fanout_sequential — y confirmamos que el número de llamadas al modelo NO cambia entre los dos
caminos.
Recursos adicionales
- Python —
concurrent.futures.ThreadPoolExecutor— La clase central de esta lección, ya usada sin cambios pordispatch_paralleldesdeagent-fundamentalsM5. - Python —
concurrent.futures.as_completed— La función que recoge resultados en el orden en que terminan, la pieza que hace explícito el no-determinismo que esta lección neutraliza consorted(). - Anthropic — Multi-agent research system — Un sistema real donde varios sub-agentes corren en paralelo y sus resultados se agregan de forma reproducible, sin depender de cuál terminó primero.
- Python — Diccionarios (orden de inserción) — El comportamiento de Python 3.7+ que hace que
results.keys()refleje el orden de llegada de los hilos, y por qué esta lección nunca confía en ese orden.