Módulo 4: Fan-out paralelo y agregación
Agregando resultados de forma determinista
Descripción
La lección 03 dejó el despacho resuelto —run_fanout_sequential corre cada sub-tarea en un orden
fijo— pero se detuvo en el resultado crudo: un dict con la salida de texto de cada especialista, sin
historial completo y sin ninguna respuesta unificada para el socio. Esta lección cierra ese ciclo:
imprime el historial completo de cada sub-tarea (igual que ya hizo el Módulo 3, lección 04, con el
pipeline), compone una respuesta final que combina las dos salidas en un solo mensaje (concepto,
igual que la síntesis del supervisor y del pipeline), y cuenta el costo completo de coordinación —
el split, las llamadas internas de cada especialista, y la síntesis.
El punto que sostiene todo lo demás sigue siendo el mismo de la lección 03, ahora aplicado a la
impresión y a la síntesis: todo lo que consume results itera sobre sorted(results), nunca
sobre el diccionario tal como quedó construido. Esta lección es la primera vez en el módulo donde
"agregar" significa algo más que "correr y guardar" — significa producir un mensaje único, coherente,
que un socio real podría leer.
Conexión con el módulo
Esta lección reusa, sin cambios, run_fanout_sequential de la lección 03. Lo nuevo es
compose_fanout_response —la síntesis, concepto— y el conteo de costo (FANOUT_SPLIT_CALLS,
COMPOSE_CALLS, llamadas internas). La lección 07 toma estos mismos números y los compara contra
la versión con concurrencia real de la lección 06, para medir el ahorro en rondas.
Analogía: el resumen que junta las dos respuestas en una sola carta
Retoma los dos mandados. Cuando la persona vuelve con el pan y con el correo revisado, no le entrega
al que preguntó dos reportes sueltos —"aquí está el pan" en un papel, "esto decía el correo" en
otro—. Junta las dos cosas en una sola respuesta: "Ya tienes el pan, y el correo no tenía nada
urgente." Esa unificación no agrega ningún mandado nuevo — es exactamente el mismo trabajo que ya
se hizo, presentado como una sola respuesta coherente. compose_fanout_response hace ese mismo
trabajo de unificación sobre las dos sub-tareas de Reservo.
Ejemplo trabajado: historial completo, síntesis y costo
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 la lección 03, sin cambios."""
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 count_model_calls(history):
return sum(1 for m in history if m["role"] == "assistant")
def count_tool_calls(history):
total = 0
for m in history:
if isinstance(m["content"], list):
total += sum(1 for b in m["content"] if b["type"] == "tool_use")
return total
def print_history(agent, history):
for i, m in enumerate(history):
role, content = m["role"], m["content"]
if isinstance(content, str):
print(f" [{agent}][{i}] {role:<9} pregunta: {content!r}")
continue
for block in content:
if block["type"] == "tool_use":
print(f" [{agent}][{i}] {role:<9} tool_use({block['name']}): {block['input']}")
elif block["type"] == "tool_result":
print(f" [{agent}][{i}] {role:<9} tool_result: {block['content']}")
elif block["type"] == "text":
print(f" [{agent}][{i}] {role:<9} texto final: {block['text']!r}")
def compose_fanout_response(results):
"""Paso final (concepto, claude-sonnet-5): combina la salida de cada
sub-tarea en una sola respuesta para el socio. Ningún dato nuevo --
solo redacción, igual que la síntesis del Módulo 3, lección 04."""
booking_out = results["booking_agent"]["output"]
policy_out = results["policy_agent"]["output"]
return f"{booking_out} Sobre la cancelación: {policy_out}"
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("--- petición compuesta ---")
print(repr(COMPOUND_REQUEST))
print()
print("--- fan-out secuencial, con historial completo por sub-tarea ---")
results = run_fanout_sequential(subtasks, model_scripts)
for agent in sorted(results):
print(f"=== {agent} ===")
print(f"tarea: {results[agent]['task']!r}")
print_history(agent, results[agent]["history"])
print()
print("--- respuesta agregada (concepto, síntesis) ---")
final_response = compose_fanout_response(results)
print(final_response)
print()
print("--- costo de coordinación ---")
FANOUT_SPLIT_CALLS = 1 # concepto: el supervisor reconoce las sub-tareas
COMPOSE_CALLS = 1 # concepto: la síntesis final para el socio
specialist_calls = sum(count_model_calls(results[a]["history"]) for a in results)
specialist_tools = sum(count_tool_calls(results[a]["history"]) for a in results)
total_calls = FANOUT_SPLIT_CALLS + specialist_calls + COMPOSE_CALLS
print(f"llamadas de split (concepto): {FANOUT_SPLIT_CALLS}")
print(f"llamadas internas (especialistas): {specialist_calls}")
print(f"llamadas de síntesis (concepto): {COMPOSE_CALLS}")
print(f"llamadas al modelo TOTAL: {total_calls}")
print(f"llamadas a tools: {specialist_tools}")
Qué esperar:
--- petición compuesta ---
'Cotiza Focus pro 3h y dime la política de cancelación.'
--- fan-out secuencial, con historial completo por sub-tarea ---
=== booking_agent ===
tarea: 'Cotiza Focus pro 3h.'
[booking_agent][0] user pregunta: 'Cotiza Focus pro 3h.'
[booking_agent][1] assistant tool_use(get_quote): {'room': 'Focus', 'tier': 'pro', 'hours': 3}
[booking_agent][2] user tool_result: {'price_cents': 6000}
[booking_agent][3] assistant texto final: 'Focus pro 3h cuesta 6000 centavos.'
=== policy_agent ===
tarea: '¿Cuál es la política de cancelación?'
[policy_agent][0] user pregunta: '¿Cuál es la política de cancelación?'
[policy_agent][1] assistant tool_use(search_docs): {'query': 'política de cancelación'}
[policy_agent][2] user tool_result: [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.
[policy_agent][3] assistant texto final: 'Puedes cancelar sin cargo hasta 2 horas antes del horario reservado. Después de eso aplica el cargo de no-presentación.'
--- respuesta agregada (concepto, síntesis) ---
Focus pro 3h cuesta 6000 centavos. Sobre la cancelación: Puedes cancelar sin cargo hasta 2 horas antes del horario reservado. Después de eso aplica el cargo de no-presentación.
--- costo de coordinación ---
llamadas de split (concepto): 1
llamadas internas (especialistas): 4
llamadas de síntesis (concepto): 1
llamadas al modelo TOTAL: 6
llamadas a tools: 2
Cada sub-tarea corrió su propio run_specialist, con su propio historial completo — booking_agent
gastó 2 llamadas al modelo (un tool_use, un texto final) y policy_agent gastó otras 2, idéntico
al costo interno que ya medías en los Módulos 2 y 3 para tareas de esta forma. Lo nuevo es la
respuesta agregada: combina "Focus pro 3h cuesta 6000 centavos" (de booking_agent) con "Puedes
cancelar sin cargo..." (de policy_agent) en un solo mensaje — ninguna de las dos sub-tareas, por sí
sola, tenía la información completa que el socio pidió.
Por qué el split y la síntesis cuestan una llamada cada uno, siempre
FANOUT_SPLIT_CALLS = 1 y COMPOSE_CALLS = 1 no son números que varíen según cuántas sub-tareas
haya — son costos fijos, uno por petición compuesta, sin importar si el fan-out reparte dos
sub-tareas o diez. El split necesita una decisión (concepto) para reconocer todas las
sub-tareas de una vez —no una decisión por sub-tarea, como si fuera un router del Módulo 2 rutando
una por una—; la síntesis necesita una decisión para combinar todos los resultados en un solo
mensaje. Esto contrasta con las llamadas internas de los especialistas, que sí crecen con el número
de sub-tareas y con la complejidad de cada una.
Errores comunes
-
Pensar que
compose_fanout_responsepuede correr antes de que TODAS las sub-tareas terminen. La función necesitaresults["booking_agent"]yresults["policy_agent"]completos — si alguna sub-tarea todavía no corrió,resultsno tiene esa clave, y la línea falla conKeyError. El Ejercicio 3 de esta lección lo confirma ejecutando. -
Sumar mal el costo total.
FANOUT_SPLIT_CALLS + specialist_calls + COMPOSE_CALLSno es lo mismo quespecialist_tools— el mismo error de contabilidad ya advertido en el Módulo 2, lección 06, y en el Módulo 3, lección 05. Esta lección los reporta por separado a propósito. -
Escribir
compose_fanout_responsede forma que dependa de CUÁLES claves existen enresults, en vez de asumir un conjunto fijo. La versión de este módulo asume quebooking_agentypolicy_agentsiempre están presentes para ESTA petición compuesta concreta — una versión más general, capaz de agregar un número variable de sub-tareas, es terreno del Módulo 7 (composición completa del sistema), no de esta lección. -
Olvidar que el historial completo (
print_history) es solo para trazabilidad, no para que la síntesis lo use.compose_fanout_responseleeresults[agent]["output"]—el texto final de cada especialista—, nunca el historial completo con lostool_use/tool_resultintermedios. -
Ejecutar esta lección después de otro ejemplo de la guía en el mismo intérprete. Aunque este ejemplo no llama a
book_room, otras lecciones de la guía sí — si compartes el proceso, alguna comparación de estado (reservo_tools.BOOKINGS) podría no coincidir con lo esperado.
Ejercicios
Ejercicio 1: Confirma tu propia ejecución (Fácil)
Ejecuta el ejemplo trabajado completo tú mismo y confirma, línea por línea, que tu salida coincide con el "Qué esperar" de arriba. Presta especial atención a los números del bloque de costo de coordinación.
Ver solución
No hay una única "solución de código" para este ejercicio — es una verificación: si tu salida coincide exactamente con el bloque "Qué esperar" del ejemplo trabajado, tu Reservo desechable arrancó limpio y el fan-out secuencial corrió sin desvíos.
Ejercicio 2: Boardroom pro 4h y el cargo por no-presentación (Medio)
Repite el flujo completo —split (a mano, sin concept_turn), despacho con run_fanout_sequential,
síntesis con compose_fanout_response— para la petición "Cotiza Boardroom pro 4h y dime si hay
cargo por no-presentarme.". Confirma el precio a mano.
Ver solución
subtasks_2 = [
SubTask(agent="booking_agent", task="Cotiza Boardroom pro 4h."),
SubTask(agent="policy_agent", task="¿Hay cargo si no me presento?"),
]
model_scripts_2 = {
"booking_agent": [
{"stop_reason": "tool_use", "content": [
{"type": "tool_use", "id": "toolu_01", "name": "get_quote",
"input": {"room": "Boardroom", "tier": "pro", "hours": 4}}]},
{"stop_reason": "end_turn", "content": [{"type": "text", "text": "Boardroom pro 4h cuesta 25600 centavos."}]},
],
"policy_agent": [
{"stop_reason": "tool_use", "content": [
{"type": "tool_use", "id": "toolu_01", "name": "search_docs",
"input": {"query": "cargo si no me presento"}}]},
{"stop_reason": "end_turn", "content": [
{"type": "text", "text": "Sí, el 50% del precio cotizado si no cancelas con 2 horas de anticipación."}]},
],
}
results_2 = run_fanout_sequential(subtasks_2, model_scripts_2)
final_response_2 = compose_fanout_response(results_2)
print("respuesta agregada:", final_response_2)
print("Boardroom pro 4h a mano:", 8000 * 4 * 80 // 100)
Salida esperada:
respuesta agregada: Boardroom pro 4h cuesta 25600 centavos. Sobre la cancelación: Sí, el 50% del precio cotizado si no cancelas con 2 horas de anticipación.
Boardroom pro 4h a mano: 25600
Explicación: 25600 = 8000 * 4 * 80 // 100, el mismo cálculo de siempre para pro.
compose_fanout_response no cambió — sigue leyendo results["booking_agent"]["output"] y
results["policy_agent"]["output"], sin importar qué sala, tier u horas se hayan pedido.
Ejercicio 3: ¿Qué pasa si compose_fanout_response recibe resultados incompletos? (Difícil)
Construye un results con solo la entrada de booking_agent (simulando que policy_agent
todavía no terminó) y ejecuta compose_fanout_response sobre él. ¿Qué excepción se produce, y por
qué confirma que la síntesis necesita que TODAS las sub-tareas hayan terminado?
Ver solución
partial_results = {"booking_agent": results_2["booking_agent"]} # falta policy_agent
try:
compose_fanout_response(partial_results)
except KeyError as e:
print(f"KeyError capturado: {e!r} -- compose necesita AMBAS claves, no solo una")
Salida esperada:
KeyError capturado: KeyError('policy_agent') -- compose necesita AMBAS claves, no solo una
Explicación: compose_fanout_response intenta leer results["policy_agent"]["output"] en su
segunda línea — si esa clave no existe todavía, Python lanza KeyError de inmediato, antes de
producir ningún texto parcial. Esto confirma, en código, algo importante sobre el patrón fan-out: la
síntesis final es un punto de sincronización — no puede correr hasta que todas las sub-tareas
que agrega hayan terminado, sin importar si el despacho fue secuencial (esta lección) o con
concurrencia real (lección 06). La lección 07 retoma esta misma idea para explicar por qué la
síntesis siempre cuenta como una ronda adicional, después de que la más lenta de las sub-tareas
termine.
Resumen y siguiente paso
- El caso base del fan-out queda completo: split (concepto) → despacho en orden fijo
(
run_fanout_sequential, lección 03) → agregación con historial completo y síntesis final (concepto) → conteo de costo. - El costo de coordinación de esta petición: 1 llamada de split, 4 llamadas internas (2 por especialista), 1 llamada de síntesis — 6 llamadas al modelo en total, 2 tool calls.
compose_fanout_responsees un punto de sincronización: necesita que TODAS las sub-tareas que combina hayan terminado — confirmado ejecutando el caso donde falta una.- El split y la síntesis cuestan una llamada cada uno, sin importar cuántas sub-tareas haya — solo las llamadas internas de los especialistas crecen con el tamaño del fan-out.
Siguiente lección: 05 — Fan-out dentro de un solo agente: la variante con pricing_agent.
Retomamos pricing_agent comparando tres salas —ya construido en el Módulo 2— para distinguirlo con
precisión del fan-out entre agentes que acabamos de construir.
Recursos adicionales
- Anthropic — Building effective agents — El patrón "parallelization" incluye, como paso final, "aggregator": combinar los resultados de las sub-tareas en una salida — exactamente lo que construye esta lección.
- Anthropic — Multi-agent research system — Un orquestador real sintetizando resultados de varios sub-agentes en una sola respuesta coherente, la misma disciplina de esta lección.
- Anthropic — Messages API reference — La forma exacta de
tool_use/tool_result/stop_reasonque cada sub-tarea de este fan-out respeta, sin cambios. - Python — Diccionarios — La estructura detrás de
results, keyed por nombre de agente, consumida siempre consorted().