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:
- Fan-out dinámico: crear un branch por cada item de la lista, sin importar cuántos hay
- Procesamiento independiente: cada item procesado por el mismo nodo, en paralelo
- 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:
| Fase | Qué hace | Analogía |
|---|---|---|
| Map | Distribuye cada item a una instancia del nodo procesador | "Reparte las tareas" |
| Process | Cada instancia procesa su item independientemente | "Cada uno hace su trabajo" |
| Reduce | Un 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:
- Crea una invocación del nodo
"process_result" - Le pasa
datacomo 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
| Aspecto | Branching estático | Send API |
|---|---|---|
| Número de branches | Fijo en compile time | Dinámico en runtime |
| Definición | add_conditional_edges → strings | Routing function → list[Send] |
| Datos por branch | Mismo estado compartido | Datos 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
search_papersretorna una lista de tamaño variable enraw_resultsroute_to_processorslee esa lista y crea unSendpor cada item — el fan-out dinámicosummarize_paperrecibe unProcessorState(no elMainStatecompleto) — solo los datos que necesita- El reducer
operator.addensummariesacumula los resultados de todos los branches compilees el deferred node — espera a que todos lossummarize_paperterminen 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
ProcessorStatedocumenta 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:
summarize_paperpuede ejecutarse N veces (por losSend)compileviene después desummarize_paper- Ergo,
compiledebe esperar a que las N instancias desummarize_paperterminen
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
| Criterio | Branching estático | Map-reduce con Send |
|---|---|---|
| Número de branches | Conocido al diseñar el grafo | Determinado por los datos en runtime |
| Ejemplo | "Busca en Wikipedia, arXiv y News" | "Resume cada paper de los resultados" |
| Datos por branch | Mismo estado compartido | Datos custom por branch (ProcessorState) |
| Nodos destino | Pueden ser diferentes nodos | Siempre el mismo nodo, diferentes datos |
| Complejidad | Menor (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 deSenden 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
- Map-Reduce in LangGraph — Tutorial oficial del patrón map-reduce con Send API
- Send API Reference — Referencia completa de la clase Send
- LangGraph Branching — Branching estático vs dinámico
- State Reducers — Cómo funcionan los reducers en LangGraph
- Fan-out / Fan-in patterns — Conceptos de fan-out dinámico
- Subgraphs vs Map-Reduce — Cuándo usar subgraphs vs map-reduce
- Python operator module — Referencia de operator.add y otros operadores usados como reducers
Módulo 7 — LangChain & LangGraph: From Chains to Agents