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

  1. Usar las claves de results en el orden en que Python las itera por defecto, sin sorted(). Un diccionario de Python preserva el orden de inserción desde la versión 3.7 — pero ese orden de inserción, bajo ThreadPoolExecutor, es el orden de finalización de los hilos, no ningún orden con significado. Iterar for agent in results sin sorted() reintroduciría exactamente el no-determinismo que esta lección elimina.

  2. Pensar que hace falta un Lock o algún mecanismo de sincronización manual para escribir en results. No en este caso: cada hilo escribe en una clave distinta del diccionario (sub.agent es único por sub-tarea), así que no hay dos hilos escribiendo el mismo dato al mismo tiempo. Si dos sub-tareas compartieran el mismo agent, sí haría falta revisar esto con cuidado — un caso que este módulo no construye.

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

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

  5. 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_parallel despacha cada sub-tarea en su propio hilo (ThreadPoolExecutor), recogiendo resultados con as_completed a medida que terminan.
  • El orden de llegada al diccionario results varía entre corridas — confirmado ejecutando el mismo fan-out varias veces— pero la salida impresa, siempre construida con sorted(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 de max_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

  1. Python — concurrent.futures.ThreadPoolExecutor — La clase central de esta lección, ya usada sin cambios por dispatch_parallel desde agent-fundamentals M5.
  2. 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 con sorted().
  3. 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.
  4. 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.