Módulo 7: Flujos Avanzados

Map-Reduce y Deferred Nodes

Descripción de la cápsula

Tienes un sistema que busca en múltiples fuentes. En la cápsula de branching, definiste 3 branches fijos — uno por fuente. Funcionó perfecto. Pero ahora llega un problema nuevo: el usuario ingresa una query, y el número de resultados varía. A veces son 2, a veces 7, a veces 15. No puedes hardcodear branches para cada posible cantidad de resultados.

Este es el problema que resuelve map-reduce: crear branches dinámicamente en runtime, basándose en los datos — no en la estructura del grafo. LangGraph lo implementa con la Send API, que permite "enviar" datos a un nodo N veces, donde N se determina durante la ejecución. Después, un deferred node espera a que todos los branches terminen antes de continuar.

Si branching es "tengo 3 fuentes fijas y busco en paralelo", map-reduce es "no sé cuántos items tengo hasta que los veo, pero quiero procesar cada uno en paralelo y luego combinar los resultados."


El problema: colecciones de tamaño dinámico

Imagina este escenario concreto: tu Research Assistant busca papers y obtiene una lista de resultados. Necesitas resumir cada paper. Pero no sabes cuántos papers va a encontrar:

# Búsqueda 1: "transformer architecture" → 5 papers
# Búsqueda 2: "quantum error correction" → 2 papers
# Búsqueda 3: "CRISPR gene editing" → 8 papers

Con branching estático, necesitarías definir un branch por cada posible paper — imposible. Lo que necesitas es:

  1. Fan-out dinámico: crear un branch por cada item de la lista, sin importar cuántos hay
  2. Procesamiento independiente: cada item procesado por el mismo nodo, en paralelo
  3. Fan-in automático: todos los resultados recolectados en el estado cuando terminen

Este es el patrón map-reduce.


Map-reduce: el patrón

El patrón map-reduce tiene tres fases:

FaseQué haceAnalogía
MapDistribuye cada item a una instancia del nodo procesador"Reparte las tareas"
ProcessCada instancia procesa su item independientemente"Cada uno hace su trabajo"
ReduceUn nodo recolecta todos los resultados y los combina"Junta todo"

En programación funcional, es map() seguido de reduce(). En LangGraph, es la Send API seguida de un deferred node con un reducer en el estado.


La Send API: fan-out dinámico

La clave es Send de langgraph.types. En vez de que un conditional edge retorne un string (el nombre del siguiente nodo), retorna una lista de objetos Send. Cada Send crea un branch independiente hacia el mismo nodo, pero con datos diferentes:

from langgraph.types import Send

def route_to_processors(state):
    """Crea un branch dinámico por cada resultado."""
    return [
        Send("process_result", {"item": item, "index": i})
        for i, item in enumerate(state["raw_results"])
    ]

Cada Send("process_result", data) hace dos cosas:

  1. Crea una invocación del nodo "process_result"
  2. Le pasa data como el estado de entrada para esa invocación

Si state["raw_results"] tiene 3 items, se crean 3 invocaciones paralelas. Si tiene 10, se crean 10. El número se determina en runtime.

Send vs branching estático

AspectoBranching estáticoSend API
Número de branchesFijo en compile timeDinámico en runtime
Definiciónadd_conditional_edges → stringsRouting function → list[Send]
Datos por branchMismo estado compartidoDatos custom por branch
Caso de uso"Busca en estas 3 fuentes fijas""Procesa cada item de esta lista"

Ejemplo completo: procesar N resultados de búsqueda

Vamos a construir un grafo que busca papers, y luego resume cada paper dinámicamente:

from dotenv import load_dotenv
load_dotenv()

import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send
from langchain.chat_models import init_chat_model
from IPython.display import Image, display

class MainState(TypedDict):
    query: str
    raw_results: list[str]
    summaries: Annotated[list[str], operator.add]

class ProcessorState(TypedDict):
    item: str
    index: int

def search_papers(state: MainState) -> dict:
    """Simula búsqueda que retorna N resultados."""
    query = state["query"]
    papers = [
        f"Paper 1: '{query}' - Propone un nuevo framework basado en attention mechanisms.",
        f"Paper 2: '{query}' - Analiza las limitaciones de los enfoques actuales.",
        f"Paper 3: '{query}' - Presenta resultados experimentales con 15% de mejora.",
        f"Paper 4: '{query}' - Review del estado del arte con 200+ referencias.",
    ]
    return {"raw_results": papers}

def route_to_processors(state: MainState) -> list[Send]:
    """Crea un Send por cada paper encontrado."""
    return [
        Send("summarize_paper", {"item": paper, "index": i})
        for i, paper in enumerate(state["raw_results"])
    ]

def summarize_paper(state: ProcessorState) -> dict:
    """Procesa un paper individual — se ejecuta N veces en paralelo."""
    model = init_chat_model("openai:gpt-4.1-mini")
    response = model.invoke(
        f"Resume este paper en una oración técnica:\n\n{state['item']}"
    )
    return {"summaries": [f"[{state['index']}] {response.content}"]}

def compile_results(state: MainState) -> dict:
    """Combina todos los resúmenes en un resultado final."""
    model = init_chat_model("openai:gpt-4.1-mini")
    context = "\n".join(state["summaries"])
    response = model.invoke(
        f"Genera un resumen ejecutivo de estos papers sobre '{state['query']}':\n\n{context}"
    )
    return {"summaries": [f"\n--- RESUMEN EJECUTIVO ---\n{response.content}"]}

graph_builder = StateGraph(MainState)
graph_builder.add_node("search", search_papers)
graph_builder.add_node("summarize_paper", summarize_paper)
graph_builder.add_node("compile", compile_results)

graph_builder.add_edge(START, "search")
graph_builder.add_conditional_edges("search", route_to_processors)
graph_builder.add_edge("summarize_paper", "compile")
graph_builder.add_edge("compile", END)

graph = graph_builder.compile()
display(Image(graph.get_graph().draw_mermaid_png()))

result = graph.invoke({"query": "transformer architecture", "raw_results": [], "summaries": []})
print(f"Papers encontrados: {len(result['raw_results'])}")
print(f"Resúmenes generados: {len(result['summaries']) - 1}")
for s in result["summaries"]:
    print(s)
# Output esperado:
# Papers encontrados: 4
# Resúmenes generados: 4
# [0] Este paper propone un nuevo framework basado en mecanismos de atención...
# [1] El estudio analiza las limitaciones de los enfoques actuales en...
# [2] Se presentan resultados experimentales que demuestran una mejora del 15%...
# [3] Una revisión exhaustiva del estado del arte con más de 200 referencias...
#
# --- RESUMEN EJECUTIVO ---
# La investigación sobre transformer architecture muestra avances en...

Anatomía del ejemplo

  1. search_papers retorna una lista de tamaño variable en raw_results
  2. route_to_processors lee esa lista y crea un Send por cada item — el fan-out dinámico
  3. summarize_paper recibe un ProcessorState (no el MainState completo) — solo los datos que necesita
  4. El reducer operator.add en summaries acumula los resultados de todos los branches
  5. compile es el deferred node — espera a que todos los summarize_paper terminen antes de ejecutar

El estado del procesador: separar datos por branch

Observa que summarize_paper usa ProcessorState, no MainState:

class MainState(TypedDict):
    query: str
    raw_results: list[str]
    summaries: Annotated[list[str], operator.add]

class ProcessorState(TypedDict):
    item: str
    index: int

Cada Send("summarize_paper", {"item": paper, "index": i}) crea una invocación con solo los datos que ese branch necesita. El nodo summarize_paper no ve la query ni los demás papers — solo su item y su índice. Esto es intencional:

  • Aislamiento: cada branch es independiente, no puede interferir con otros
  • Eficiencia: no copia el estado completo N veces
  • Claridad: el tipo ProcessorState documenta qué datos recibe cada branch

El retorno {"summaries": [...]} se escribe de vuelta al MainState porque summaries tiene el reducer operator.add. Cada branch agrega su resultado a la lista compartida.


Deferred nodes: el fan-in automático

El nodo compile en el ejemplo es un deferred node — un nodo que espera a que todos los branches upstream terminen antes de ejecutarse. No necesitas configurar nada especial: LangGraph lo infiere de la estructura del grafo.

Cuando defines graph_builder.add_edge("summarize_paper", "compile"), LangGraph sabe que:

  1. summarize_paper puede ejecutarse N veces (por los Send)
  2. compile viene después de summarize_paper
  3. Ergo, compile debe esperar a que las N instancias de summarize_paper terminen

El deferred node tiene acceso al estado completo con todos los resultados acumulados por el reducer. Cuando compile se ejecuta, state["summaries"] contiene los resúmenes de todos los papers — no importa si fueron 2 o 20.

Cómo sabe LangGraph que "todos terminaron"

LangGraph mantiene un contador interno de branches activos. Cada Send incrementa el contador. Cada finalización de summarize_paper lo decrementa. Cuando llega a 0, el deferred node se activa. Tú no manejas este contador — es automático.


Reducers: la pieza clave del fan-in

Sin un reducer, el fan-in no funciona. Si summaries fuera list[str] sin Annotated[..., operator.add], cada branch sobrescribiría el resultado del anterior:

# ❌ Sin reducer: solo el último branch sobrevive
class BadState(TypedDict):
    summaries: list[str]

# ✅ Con reducer: todos los branches acumulan
class GoodState(TypedDict):
    summaries: Annotated[list[str], operator.add]

El reducer operator.add concatena listas. Si branch 1 retorna ["resumen A"] y branch 2 retorna ["resumen B"], el estado termina con ["resumen A", "resumen B"].

Otros reducers útiles para map-reduce:

import operator

# Acumular listas
results: Annotated[list[str], operator.add]

# Contar éxitos
success_count: Annotated[int, operator.add]

# Reducer custom: mantener solo los mejores resultados
def keep_top_scores(current: list[dict], new: list[dict]) -> list[dict]:
    combined = current + new
    return sorted(combined, key=lambda x: x["score"], reverse=True)[:5]

top_results: Annotated[list[dict], keep_top_scores]

Cuándo usar map-reduce vs branching estático

CriterioBranching estáticoMap-reduce con Send
Número de branchesConocido al diseñar el grafoDeterminado por los datos en runtime
Ejemplo"Busca en Wikipedia, arXiv y News""Resume cada paper de los resultados"
Datos por branchMismo estado compartidoDatos custom por branch (ProcessorState)
Nodos destinoPueden ser diferentes nodosSiempre el mismo nodo, diferentes datos
ComplejidadMenor (edges explícitos)Mayor (Send + ProcessorState + reducer)

Regla: si sabes los branches al escribir el código, usa branching estático. Si los branches dependen de los datos, usa Send.


Patterns de consenso: combinar resultados de branches paralelos

Cuando N branches retornan resultados, el deferred node necesita una estrategia para combinarlos. Estos son los patterns más comunes:

Pattern 1: Concatenación simple

Todos los resultados se juntan en una lista. Útil cuando cada resultado es independiente y valioso:

def compile_all(state: MainState) -> dict:
    return {"final_output": "\n".join(state["summaries"])}

Pattern 2: Votación por mayoría

Múltiples branches analizan el mismo input y votan. Útil para clasificación o verificación:

from dotenv import load_dotenv
load_dotenv()

import operator
from typing import TypedDict, Annotated
from collections import Counter
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send
from langchain.chat_models import init_chat_model

class VoteState(TypedDict):
    text: str
    votes: Annotated[list[str], operator.add]

class VoterInput(TypedDict):
    text: str
    voter_id: int

def create_voters(state: VoteState) -> list[Send]:
    return [Send("vote", {"text": state["text"], "voter_id": i}) for i in range(3)]

def vote(state: VoterInput) -> dict:
    model = init_chat_model("openai:gpt-4.1-mini")
    response = model.invoke(
        f"Clasifica este texto como 'positive', 'negative', o 'neutral'. "
        f"Responde SOLO con la clasificación.\n\n{state['text']}"
    )
    return {"votes": [response.content.strip().lower()]}

def tally_votes(state: VoteState) -> dict:
    counts = Counter(state["votes"])
    winner = counts.most_common(1)[0][0]
    return {"votes": [f"RESULTADO: {winner} ({dict(counts)})"]}

graph_builder = StateGraph(VoteState)
graph_builder.add_node("vote", vote)
graph_builder.add_node("tally", tally_votes)

graph_builder.add_conditional_edges(START, create_voters)
graph_builder.add_edge("vote", "tally")
graph_builder.add_edge("tally", END)

graph = graph_builder.compile()

result = graph.invoke({
    "text": "El producto es bueno pero el servicio post-venta es terrible.",
    "votes": [],
})
print(f"Votos: {result['votes']}")
# Output esperado:
# Votos: ['negative', 'negative', 'neutral', 'RESULTADO: negative ({negative: 2, neutral: 1})']

Pattern 3: Promedio ponderado

Cada branch retorna un score. El deferred node calcula el promedio:

import operator
from typing import TypedDict, Annotated

class ScoreState(TypedDict):
    scores: Annotated[list[float], operator.add]

def aggregate_scores(state: ScoreState) -> dict:
    avg = sum(state["scores"]) / len(state["scores"]) if state["scores"] else 0
    return {"scores": [avg]}

Pattern 4: Filtrado por calidad

Cada branch retorna un resultado con un score de confianza. El deferred node filtra los de baja calidad:

import operator
from typing import TypedDict, Annotated

class FilterState(TypedDict):
    results: Annotated[list[dict], operator.add]

def filter_high_quality(state: FilterState) -> dict:
    good = [r for r in state["results"] if r.get("confidence", 0) >= 0.7]
    return {"results": good if good else state["results"][:1]}

Map-reduce con la Functional API

El mismo patrón es posible con la Functional API usando Futures para paralelismo:

from dotenv import load_dotenv
load_dotenv()

from langgraph.func import entrypoint, task
from langchain.chat_models import init_chat_model

model = init_chat_model("openai:gpt-4.1-mini")

@task
def search(query: str) -> list:
    return [
        f"Paper 1 sobre {query}: nuevo framework propuesto.",
        f"Paper 2 sobre {query}: análisis de limitaciones.",
        f"Paper 3 sobre {query}: mejora experimental del 15%.",
    ]

@task
def summarize_one(paper: str) -> str:
    return model.invoke(f"Resume en una oración: {paper}").content

@task
def compile_summaries(query: str, summaries: list) -> str:
    context = "\n".join(f"- {s}" for s in summaries)
    return model.invoke(
        f"Resumen ejecutivo sobre '{query}':\n\n{context}"
    ).content

@entrypoint()
def map_reduce_functional(query: str) -> dict:
    papers = search(query).result()

    summary_futures = [summarize_one(p) for p in papers]
    summaries = [f.result() for f in summary_futures]

    executive = compile_summaries(query, summaries).result()

    return {
        "papers_found": len(papers),
        "summaries": summaries,
        "executive_summary": executive,
    }

result = map_reduce_functional.invoke("retrieval augmented generation")
print(f"Papers: {result['papers_found']}")
for s in result["summaries"]:
    print(f"  - {s[:80]}...")
print(f"\nResumen ejecutivo: {result['executive_summary'][:150]}...")
# Output esperado:
# Papers: 3
#   - Este paper propone un nuevo framework para RAG basado en...
#   - El estudio analiza las principales limitaciones de los enfoques...
#   - Los resultados experimentales muestran una mejora del 15% en...
# Resumen ejecutivo: La investigación sobre RAG muestra avances significativos...

La diferencia clave: con la Functional API no tienes Send — usas list comprehensions y Futures. Esto funciona bien cuando el procesamiento es simple. Si necesitas subgrafos, estado tipado, o visualización, la Graph API con Send es más potente.


Ejemplo avanzado: map-reduce con procesamiento condicional

No todos los items necesitan el mismo procesamiento. Puedes usar Send para dirigir items a diferentes nodos según sus características:

from dotenv import load_dotenv
load_dotenv()

import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send
from langchain.chat_models import init_chat_model

class MainState(TypedDict):
    items: list[dict]
    results: Annotated[list[str], operator.add]

class ItemState(TypedDict):
    item: dict

def classify_and_route(state: MainState) -> list[Send]:
    sends = []
    for item in state["items"]:
        if item.get("type") == "short":
            sends.append(Send("process_short", {"item": item}))
        else:
            sends.append(Send("process_long", {"item": item}))
    return sends

def process_short(state: ItemState) -> dict:
    text = state["item"]["text"]
    return {"results": [f"[SHORT] Procesado rápido: {text[:50]}"]}

def process_long(state: ItemState) -> dict:
    model = init_chat_model("openai:gpt-4.1-mini")
    text = state["item"]["text"]
    response = model.invoke(f"Resume este texto largo en una oración:\n\n{text}")
    return {"results": [f"[LONG] {response.content}"]}

def merge_results(state: MainState) -> dict:
    return {}

graph_builder = StateGraph(MainState)
graph_builder.add_node("process_short", process_short)
graph_builder.add_node("process_long", process_long)
graph_builder.add_node("merge", merge_results)

graph_builder.add_conditional_edges(START, classify_and_route)
graph_builder.add_edge("process_short", "merge")
graph_builder.add_edge("process_long", "merge")
graph_builder.add_edge("merge", END)

graph = graph_builder.compile()

result = graph.invoke({
    "items": [
        {"type": "short", "text": "Python es un lenguaje de programación."},
        {"type": "long", "text": "Los transformers son una arquitectura de deep learning que usa mecanismos de atención para procesar secuencias de datos en paralelo, eliminando la necesidad de procesamiento recurrente."},
        {"type": "short", "text": "LangGraph maneja grafos."},
    ],
    "results": [],
})
for r in result["results"]:
    print(r)
# Output esperado:
# [SHORT] Procesado rápido: Python es un lenguaje de programación.
# [LONG] Los transformers son una arquitectura basada en atención que...
# [SHORT] Procesado rápido: LangGraph maneja grafos.

Los Send pueden apuntar a diferentes nodos — no todos los items tienen que ir al mismo procesador. La función de routing decide qué nodo procesa cada item.


Troubleshooting

Problema 1: "Los resultados del map están vacíos"

Síntoma: El deferred node recibe una lista vacía en el campo acumulado.

Causa: El nodo procesador no retorna datos en el campo correcto, o el campo no tiene reducer.

Solución: Verifica dos cosas:

# 1. El campo del estado tiene reducer
class State(TypedDict):
    results: Annotated[list[str], operator.add]  # ✅ Con reducer

# 2. El nodo procesador retorna en el campo correcto
def process(state: ProcessorState) -> dict:
    return {"results": ["mi resultado"]}  # ✅ Key coincide

Problema 2: "Send retorna error: node not found"

Síntoma: ValueError: Node 'process_result' not found in graph.

Causa: El nombre del nodo en Send("process_result", data) no coincide con el nombre registrado en add_node.

Solución: Verifica que el string en Send coincida exactamente:

graph_builder.add_node("process_result", my_function)  # Registro
Send("process_result", data)                            # Uso — debe coincidir

Problema 3: "Solo recibo el resultado del último branch"

Síntoma: El deferred node solo tiene un resultado en vez de N.

Causa: El campo del estado no tiene operator.add como reducer — cada branch sobrescribe al anterior.

Solución:

# ❌ Sin reducer — último branch gana
results: list[str]

# ✅ Con reducer — todos los branches acumulan
results: Annotated[list[str], operator.add]

Problema 4: "El orden de los resultados es impredecible"

Síntoma: Los resultados del map-reduce llegan en orden diferente cada vez.

Causa: Los branches se ejecutan en paralelo — el orden de finalización no es determinista.

Solución: Incluye un índice en los datos del Send y ordena en el deferred node:

def route(state):
    return [
        Send("process", {"item": item, "index": i})
        for i, item in enumerate(state["items"])
    ]

def merge(state):
    sorted_results = sorted(state["results"], key=lambda r: r["index"])
    return {"results": sorted_results}

Problema 5: "Un branch falla y el deferred node nunca se ejecuta"

Síntoma: Si uno de los N branches lanza una excepción, todo el grafo falla.

Causa: Por defecto, una excepción en cualquier branch cancela la ejecución.

Solución: Maneja errores dentro del nodo procesador:

def process(state: ProcessorState) -> dict:
    try:
        result = expensive_operation(state["item"])
        return {"results": [{"status": "ok", "data": result}]}
    except Exception as e:
        return {"results": [{"status": "error", "error": str(e)}]}

Ejercicios

Ejercicio 1: Map-reduce básico con Send (Fácil)

Crea un grafo que reciba una lista de ciudades en el estado, use Send para crear un branch por cada ciudad que obtenga su "clima" (simulado), y un deferred node que combine todos los climas en un reporte.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send
from IPython.display import Image, display

class MainState(TypedDict):
    cities: list[str]
    weather_reports: Annotated[list[str], operator.add]

class CityState(TypedDict):
    city: str

MOCK_WEATHER = {
    "madrid": "Soleado, 28°C",
    "london": "Nublado, 14°C",
    "tokyo": "Lluvia ligera, 22°C",
    "new york": "Parcialmente nublado, 25°C",
}

def route_to_cities(state: MainState) -> list[Send]:
    return [Send("get_weather", {"city": city}) for city in state["cities"]]

def get_weather(state: CityState) -> dict:
    city = state["city"]
    weather = MOCK_WEATHER.get(city.lower(), "Datos no disponibles")
    return {"weather_reports": [f"{city}: {weather}"]}

def compile_report(state: MainState) -> dict:
    return {}

graph_builder = StateGraph(MainState)
graph_builder.add_node("get_weather", get_weather)
graph_builder.add_node("compile", compile_report)

graph_builder.add_conditional_edges(START, route_to_cities)
graph_builder.add_edge("get_weather", "compile")
graph_builder.add_edge("compile", END)

graph = graph_builder.compile()
display(Image(graph.get_graph().draw_mermaid_png()))

result = graph.invoke({
    "cities": ["Madrid", "London", "Tokyo"],
    "weather_reports": [],
})
print("Reporte del clima:")
for report in result["weather_reports"]:
    print(f"  {report}")
# Output esperado:
# Reporte del clima:
#   Madrid: Soleado, 28°C
#   London: Nublado, 14°C
#   Tokyo: Lluvia ligera, 22°C

El grafo crea dinámicamente un branch por cada ciudad. Si agregas "New York" a la lista, automáticamente se crea un cuarto branch.

Ejercicio 2: Map-reduce con procesamiento LLM (Fácil)

Adapta el ejercicio anterior para que, en vez de buscar clima, cada branch use un LLM para generar un dato curioso sobre la ciudad. El deferred node debe compilar todos los datos en un párrafo usando el LLM.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send
from langchain.chat_models import init_chat_model

class MainState(TypedDict):
    cities: list[str]
    fun_facts: Annotated[list[str], operator.add]
    compiled: str

class CityState(TypedDict):
    city: str

def route_cities(state: MainState) -> list[Send]:
    return [Send("get_fact", {"city": c}) for c in state["cities"]]

def get_fact(state: CityState) -> dict:
    model = init_chat_model("openai:gpt-4.1-mini")
    response = model.invoke(
        f"Genera UN dato curioso sobre {state['city']} en una oración."
    )
    return {"fun_facts": [f"{state['city']}: {response.content}"]}

def compile_facts(state: MainState) -> dict:
    model = init_chat_model("openai:gpt-4.1-mini")
    context = "\n".join(state["fun_facts"])
    response = model.invoke(
        f"Combina estos datos curiosos en un párrafo entretenido:\n\n{context}"
    )
    return {"compiled": response.content}

graph_builder = StateGraph(MainState)
graph_builder.add_node("get_fact", get_fact)
graph_builder.add_node("compile", compile_facts)

graph_builder.add_conditional_edges(START, route_cities)
graph_builder.add_edge("get_fact", "compile")
graph_builder.add_edge("compile", END)

graph = graph_builder.compile()

result = graph.invoke({
    "cities": ["París", "Tokio", "Buenos Aires", "El Cairo"],
    "fun_facts": [],
    "compiled": "",
})
print(f"Ciudades procesadas: {len(result['fun_facts'])}")
for fact in result["fun_facts"]:
    print(f"  {fact[:80]}...")
print(f"\nPárrafo compilado: {result['compiled'][:200]}...")
# Output esperado:
# Ciudades procesadas: 4
#   París: La Torre Eiffel fue construida como estructura temporal para...
#   Tokio: Tokio tiene más restaurantes con estrellas Michelin que...
#   Buenos Aires: Buenos Aires tiene la avenida más ancha del mundo...
#   El Cairo: Las pirámides de Giza son las únicas maravillas del...
# Párrafo compilado: Alrededor del mundo hay datos fascinantes...

Ejercicio 3: Votación con 5 jueces (Medio)

Implementa un sistema de clasificación de sentimiento donde 5 "jueces" LLM clasifican un texto independientemente. Usa Send para crear los 5 jueces. El deferred node cuenta votos y decide por mayoría. Muestra el desglose de votos.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import operator
from typing import TypedDict, Annotated
from collections import Counter
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send
from langchain.chat_models import init_chat_model

class MainState(TypedDict):
    text: str
    num_judges: int
    votes: Annotated[list[dict], operator.add]
    verdict: str

class JudgeState(TypedDict):
    text: str
    judge_id: int

def create_judges(state: MainState) -> list[Send]:
    return [
        Send("judge", {"text": state["text"], "judge_id": i})
        for i in range(state["num_judges"])
    ]

def judge(state: JudgeState) -> dict:
    model = init_chat_model("openai:gpt-4.1-mini")
    response = model.invoke(
        f"Eres el juez #{state['judge_id']}. Clasifica el sentimiento de este texto "
        f"como EXACTAMENTE una de estas opciones: positive, negative, neutral.\n"
        f"Responde SOLO con la clasificación, nada más.\n\n"
        f"Texto: {state['text']}"
    )
    classification = response.content.strip().lower()
    valid = {"positive", "negative", "neutral"}
    if classification not in valid:
        classification = "neutral"
    return {"votes": [{"judge": state["judge_id"], "vote": classification}]}

def tally(state: MainState) -> dict:
    vote_values = [v["vote"] for v in state["votes"]]
    counts = Counter(vote_values)
    winner, count = counts.most_common(1)[0]
    total = len(vote_values)
    breakdown = ", ".join(f"{k}: {v}/{total}" for k, v in counts.most_common())
    return {"verdict": f"{winner} (consenso: {breakdown})"}

graph_builder = StateGraph(MainState)
graph_builder.add_node("judge", judge)
graph_builder.add_node("tally", tally)

graph_builder.add_conditional_edges(START, create_judges)
graph_builder.add_edge("judge", "tally")
graph_builder.add_edge("tally", END)

graph = graph_builder.compile()

result = graph.invoke({
    "text": "El producto llegó rápido y funciona bien, aunque el empaque estaba dañado.",
    "num_judges": 5,
    "votes": [],
    "verdict": "",
})
print(f"Votos individuales:")
for v in result["votes"]:
    print(f"  Juez #{v['judge']}: {v['vote']}")
print(f"\nVeredicto: {result['verdict']}")
# Output esperado:
# Votos individuales:
#   Juez #0: positive
#   Juez #1: positive
#   Juez #2: neutral
#   Juez #3: positive
#   Juez #4: positive
# Veredicto: positive (consenso: positive: 4/5, neutral: 1/5)

El número de jueces es configurable — num_judges determina cuántos Send se crean. Podrías usar 3 para velocidad o 7 para más confianza.

Ejercicio 4: Map-reduce con manejo de errores (Medio)

Crea un grafo map-reduce donde algunos items fallan al procesarse (simula errores para items específicos). El deferred node debe separar éxitos de errores, y generar un reporte que incluya ambos. El sistema no debe crashear por un item fallido.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send

class MainState(TypedDict):
    urls: list[str]
    results: Annotated[list[dict], operator.add]

class FetchState(TypedDict):
    url: str
    index: int

MOCK_RESPONSES = {
    "https://api.example.com/data1": "Datos del endpoint 1 obtenidos correctamente.",
    "https://api.example.com/data2": None,
    "https://api.example.com/data3": "Datos del endpoint 3 obtenidos correctamente.",
    "https://api.example.com/data4": None,
    "https://api.example.com/data5": "Datos del endpoint 5 obtenidos correctamente.",
}

def route_to_fetchers(state: MainState) -> list[Send]:
    return [
        Send("fetch_url", {"url": url, "index": i})
        for i, url in enumerate(state["urls"])
    ]

def fetch_url(state: FetchState) -> dict:
    url = state["url"]
    response = MOCK_RESPONSES.get(url)
    if response is None:
        return {"results": [{
            "index": state["index"],
            "url": url,
            "status": "error",
            "data": None,
            "error": f"Connection timeout for {url}",
        }]}
    return {"results": [{
        "index": state["index"],
        "url": url,
        "status": "ok",
        "data": response,
        "error": None,
    }]}

def compile_report(state: MainState) -> dict:
    sorted_results = sorted(state["results"], key=lambda r: r["index"])
    successes = [r for r in sorted_results if r["status"] == "ok"]
    failures = [r for r in sorted_results if r["status"] == "error"]
    print(f"Éxitos: {len(successes)}/{len(sorted_results)}")
    print(f"Errores: {len(failures)}/{len(sorted_results)}")
    for f in failures:
        print(f"  ⚠️ {f['url']}: {f['error']}")
    for s in successes:
        print(f"  ✅ {s['url']}: {s['data'][:50]}")
    return {}

graph_builder = StateGraph(MainState)
graph_builder.add_node("fetch_url", fetch_url)
graph_builder.add_node("compile", compile_report)

graph_builder.add_conditional_edges(START, route_to_fetchers)
graph_builder.add_edge("fetch_url", "compile")
graph_builder.add_edge("compile", END)

graph = graph_builder.compile()

result = graph.invoke({
    "urls": list(MOCK_RESPONSES.keys()),
    "results": [],
})
# Output esperado:
# Éxitos: 3/5
# Errores: 2/5
#   ⚠️ https://api.example.com/data2: Connection timeout for https://api.example.com/data2
#   ⚠️ https://api.example.com/data4: Connection timeout for https://api.example.com/data4
#   ✅ https://api.example.com/data1: Datos del endpoint 1 obtenidos correctamente.
#   ✅ https://api.example.com/data3: Datos del endpoint 3 obtenidos correctamente.
#   ✅ https://api.example.com/data5: Datos del endpoint 5 obtenidos correctamente.

El truco: en vez de lanzar excepciones, el nodo procesador retorna un dict con status: "error". Esto evita que el grafo falle y permite que el deferred node genere un reporte completo.

Ejercicio 5: Map-reduce con routing condicional por tipo (Avanzado)

Implementa un grafo que reciba una lista de documentos mixtos (emails, tweets, artículos). Usa Send para dirigir cada documento a un nodo procesador diferente según su tipo: "process_email", "process_tweet", "process_article". Cada procesador extrae información diferente. El deferred node combina todo.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send
from langchain.chat_models import init_chat_model

class MainState(TypedDict):
    documents: list[dict]
    extractions: Annotated[list[dict], operator.add]
    summary: str

class DocState(TypedDict):
    doc: dict

def route_by_type(state: MainState) -> list[Send]:
    sends = []
    for doc in state["documents"]:
        doc_type = doc.get("type", "article")
        node_name = f"process_{doc_type}"
        sends.append(Send(node_name, {"doc": doc}))
    return sends

def process_email(state: DocState) -> dict:
    doc = state["doc"]
    return {"extractions": [{
        "type": "email",
        "from": doc.get("from", "unknown"),
        "subject": doc.get("subject", ""),
        "urgency": "high" if "urgente" in doc.get("body", "").lower() else "normal",
    }]}

def process_tweet(state: DocState) -> dict:
    doc = state["doc"]
    body = doc.get("body", "")
    hashtags = [w for w in body.split() if w.startswith("#")]
    return {"extractions": [{
        "type": "tweet",
        "author": doc.get("author", "unknown"),
        "hashtags": hashtags,
        "char_count": len(body),
    }]}

def process_article(state: DocState) -> dict:
    model = init_chat_model("openai:gpt-4.1-mini")
    response = model.invoke(
        f"Resume este artículo en una oración:\n\n{state['doc'].get('body', '')}"
    )
    return {"extractions": [{
        "type": "article",
        "title": state["doc"].get("title", ""),
        "summary": response.content,
    }]}

def compile_all(state: MainState) -> dict:
    model = init_chat_model("openai:gpt-4.1-mini")
    context = "\n".join(str(e) for e in state["extractions"])
    response = model.invoke(
        f"Genera un resumen ejecutivo de estos {len(state['extractions'])} documentos procesados:\n\n{context}"
    )
    return {"summary": response.content}

graph_builder = StateGraph(MainState)
graph_builder.add_node("process_email", process_email)
graph_builder.add_node("process_tweet", process_tweet)
graph_builder.add_node("process_article", process_article)
graph_builder.add_node("compile", compile_all)

graph_builder.add_conditional_edges(START, route_by_type)
graph_builder.add_edge("process_email", "compile")
graph_builder.add_edge("process_tweet", "compile")
graph_builder.add_edge("process_article", "compile")
graph_builder.add_edge("compile", END)

graph = graph_builder.compile()

result = graph.invoke({
    "documents": [
        {"type": "email", "from": "jefe@empresa.com", "subject": "Reporte Q4", "body": "Necesito el reporte urgente para la junta."},
        {"type": "tweet", "author": "@techguru", "body": "LangGraph es increíble para workflows de IA #AI #LangGraph #Python"},
        {"type": "article", "title": "RAG en 2026", "body": "La técnica RAG ha evolucionado significativamente, integrando búsqueda semántica con modelos generativos para respuestas más precisas."},
        {"type": "email", "from": "equipo@dev.com", "subject": "Deploy v2.1", "body": "El deploy fue exitoso, sin incidentes."},
    ],
    "extractions": [],
    "summary": "",
})
print(f"Documentos procesados: {len(result['extractions'])}")
for ext in result["extractions"]:
    print(f"  [{ext['type']}] {ext}")
print(f"\nResumen: {result['summary'][:200]}...")
# Output esperado:
# Documentos procesados: 4
#   [email] {'type': 'email', 'from': 'jefe@empresa.com', 'subject': 'Reporte Q4', 'urgency': 'high'}
#   [tweet] {'type': 'tweet', 'author': '@techguru', 'hashtags': ['#AI', '#LangGraph', '#Python'], 'char_count': 67}
#   [article] {'type': 'article', 'title': 'RAG en 2026', 'summary': 'RAG ha evolucionado integrando...'}
#   [email] {'type': 'email', 'from': 'equipo@dev.com', 'subject': 'Deploy v2.1', 'urgency': 'normal'}
# Resumen: Se procesaron 4 documentos: 2 emails (1 urgente sobre reporte Q4...

Este ejercicio combina map-reduce con routing condicional — cada tipo de documento va a un procesador especializado, pero todos convergen en el mismo deferred node.

Ejercicio 6: Map-reduce anidado (Avanzado)

Implementa un sistema de dos niveles: el primer map-reduce procesa una lista de temas (genera sub-preguntas por tema), y el segundo nivel procesa cada sub-pregunta. Usa Send en ambos niveles. El resultado final es un diccionario con tema → [respuestas].

Ver solución
from dotenv import load_dotenv
load_dotenv()

import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send
from langchain.chat_models import init_chat_model

class MainState(TypedDict):
    topics: list[str]
    sub_questions: Annotated[list[dict], operator.add]
    answers: Annotated[list[dict], operator.add]

class TopicState(TypedDict):
    topic: str

class QuestionState(TypedDict):
    topic: str
    question: str

model = init_chat_model("openai:gpt-4.1-mini")

def route_topics(state: MainState) -> list[Send]:
    return [Send("generate_questions", {"topic": t}) for t in state["topics"]]

def generate_questions(state: TopicState) -> dict:
    response = model.invoke(
        f"Genera exactamente 2 preguntas técnicas sobre '{state['topic']}'. "
        f"Responde SOLO con las preguntas, una por línea."
    )
    questions = [q.strip() for q in response.content.strip().split("\n") if q.strip()][:2]
    return {"sub_questions": [
        {"topic": state["topic"], "question": q} for q in questions
    ]}

def route_questions(state: MainState) -> list[Send]:
    return [
        Send("answer_question", {"topic": sq["topic"], "question": sq["question"]})
        for sq in state["sub_questions"]
    ]

def answer_question(state: QuestionState) -> dict:
    response = model.invoke(
        f"Responde esta pregunta técnica en 2 oraciones:\n{state['question']}"
    )
    return {"answers": [{
        "topic": state["topic"],
        "question": state["question"],
        "answer": response.content,
    }]}

def compile_final(state: MainState) -> dict:
    return {}

graph_builder = StateGraph(MainState)
graph_builder.add_node("generate_questions", generate_questions)
graph_builder.add_node("collect_questions", lambda state: {})
graph_builder.add_node("answer_question", answer_question)
graph_builder.add_node("compile", compile_final)

graph_builder.add_conditional_edges(START, route_topics)
graph_builder.add_edge("generate_questions", "collect_questions")
graph_builder.add_conditional_edges("collect_questions", route_questions)
graph_builder.add_edge("answer_question", "compile")
graph_builder.add_edge("compile", END)

graph = graph_builder.compile()

result = graph.invoke({
    "topics": ["RAG", "Fine-tuning"],
    "sub_questions": [],
    "answers": [],
})

print(f"Temas: {len(result['topics'])}")
print(f"Sub-preguntas generadas: {len(result['sub_questions'])}")
print(f"Respuestas: {len(result['answers'])}")
for topic in result["topics"]:
    topic_answers = [a for a in result["answers"] if a["topic"] == topic]
    print(f"\n--- {topic} ---")
    for a in topic_answers:
        print(f"  Q: {a['question'][:60]}...")
        print(f"  A: {a['answer'][:80]}...")
# Output esperado:
# Temas: 2
# Sub-preguntas generadas: 4
# Respuestas: 4
#
# --- RAG ---
#   Q: ¿Cuáles son las principales estrategias de chunking para RAG...
#   A: Las estrategias de chunking incluyen fixed-size, semantic, y recursive...
#   Q: ¿Cómo se evalúa la calidad de un sistema RAG...
#   A: La calidad de un sistema RAG se evalúa con métricas como faithfulness...
#
# --- Fine-tuning ---
#   Q: ¿Cuándo es preferible fine-tuning sobre few-shot prompting...
#   A: Fine-tuning es preferible cuando tienes suficientes datos de entrenamiento...
#   Q: ¿Qué técnicas de fine-tuning eficiente existen para LLMs...
#   A: Las técnicas de PEFT como LoRA y QLoRA permiten fine-tuning eficiente...

Este es map-reduce de dos niveles: primero se expanden temas en sub-preguntas (map), se recolectan (reduce), luego se expanden sub-preguntas en respuestas (map) y se recolectan (reduce). El nodo intermedio collect_questions actúa como deferred node del primer nivel y punto de partida del segundo.


Resumen

En esta cápsula aprendiste:

  • El problema: procesar colecciones de tamaño variable — no puedes hardcodear N branches cuando N cambia en runtime
  • Send API (Send("node", data)) crea branches dinámicamente en runtime — la routing function retorna una lista de Send en vez de un string
  • ProcessorState separado permite que cada branch reciba solo los datos que necesita, sin copiar el estado completo
  • Reducers (Annotated[list, operator.add]) son obligatorios para que los resultados de múltiples branches se acumulen en vez de sobrescribirse
  • Deferred nodes esperan automáticamente a que todos los branches upstream terminen — LangGraph lleva el contador internamente
  • Patterns de consenso: concatenación, votación por mayoría, promedio ponderado, filtrado por calidad — cada uno resuelve un tipo diferente de combinación de resultados
  • Send puede apuntar a diferentes nodos — no todos los items tienen que ir al mismo procesador
  • Functional API puede lograr map-reduce con Futures y list comprehensions, pero la Graph API con Send es más potente para casos complejos

Próxima cápsula: Error Handling Avanzado — qué pasa cuando tus branches fallan, y cómo construir sistemas que se degradan con gracia en vez de crashear.


Recursos adicionales

  1. Map-Reduce in LangGraph — Tutorial oficial del patrón map-reduce con Send API
  2. Send API Reference — Referencia completa de la clase Send
  3. LangGraph Branching — Branching estático vs dinámico
  4. State Reducers — Cómo funcionan los reducers en LangGraph
  5. Fan-out / Fan-in patterns — Conceptos de fan-out dinámico
  6. Subgraphs vs Map-Reduce — Cuándo usar subgraphs vs map-reduce
  7. Python operator module — Referencia de operator.add y otros operadores usados como reducers

Módulo 7 — LangChain & LangGraph: From Chains to Agents