Módulo 7: Flujos Avanzados
Branching y Merging
Descripción de la cápsula
Tu Research Agent busca información en 3 fuentes: web, académica y noticias. Ahora mismo, las búsquedas van una detrás de otra: 3 segundos + 3 segundos + 3 segundos = 9 segundos. Pero las búsquedas no dependen entre sí — ninguna necesita el resultado de otra para empezar. ¿Y si las ejecutaras al mismo tiempo? En paralelo: max(3s, 3s, 3s) = 3 segundos. Tres veces más rápido, sin cambiar la lógica de búsqueda.
Este patrón se llama fan-out/fan-in: un nodo dispara múltiples nodos que corren en paralelo (fan-out), y luego un nodo de merge recoge todos los resultados (fan-in). Es el patrón que convierte latencia secuencial en latencia paralela — y los usuarios lo notan.
Branching sin merge es un grafo incompleto. Si disparas 3 búsquedas pero nunca juntas los resultados, no tienes un workflow — tienes 3 workflows inconexos. Esta cápsula cubre el ciclo completo: divergir → ejecutar en paralelo → convergir → procesar resultados combinados.
El problema: búsquedas secuenciales
Observa este grafo que busca en 3 fuentes de forma secuencial:
from dotenv import load_dotenv
load_dotenv()
from typing import TypedDict, Annotated
import operator
import time
from langgraph.graph import StateGraph, START, END
class ResearchState(TypedDict):
query: str
results: Annotated[list[str], operator.add]
total_time: float
def search_web(state: ResearchState) -> dict:
time.sleep(1)
return {"results": [f"[Web] Resultados sobre '{state['query']}': Python 3.12 mejora tipado y rendimiento."]}
def search_academic(state: ResearchState) -> dict:
time.sleep(1)
return {"results": [f"[Academic] Paper sobre '{state['query']}': Análisis comparativo de type checkers."]}
def search_news(state: ResearchState) -> dict:
time.sleep(1)
return {"results": [f"[News] Noticias sobre '{state['query']}': Python domina en AI según encuesta 2026."]}
graph = StateGraph(ResearchState)
graph.add_node("search_web", search_web)
graph.add_node("search_academic", search_academic)
graph.add_node("search_news", search_news)
graph.add_edge(START, "search_web")
graph.add_edge("search_web", "search_academic")
graph.add_edge("search_academic", "search_news")
graph.add_edge("search_news", END)
app = graph.compile()
start = time.time()
result = app.invoke({"query": "Python 3.12", "results": [], "total_time": 0.0})
elapsed = time.time() - start
print(f"Tiempo total: {elapsed:.1f}s")
print(f"Resultados: {len(result['results'])}")
for r in result["results"]:
print(f" {r[:70]}...")
# Output esperado:
# Tiempo total: 3.0s
# Resultados: 3
# [Web] Resultados sobre 'Python 3.12': Python 3.12 mejora tipado y rend...
# [Academic] Paper sobre 'Python 3.12': Análisis comparativo de type chec...
# [News] Noticias sobre 'Python 3.12': Python domina en AI según encuest...
3 segundos. Cada búsqueda espera a que la anterior termine. El flujo es search_web → search_academic → search_news. Pero ninguna búsqueda usa el resultado de otra — son completamente independientes.
Fan-out: un nodo dispara múltiples nodos en paralelo
Fan-out es cuando un nodo (o el punto de entrada) tiene edges hacia múltiples nodos. LangGraph detecta estos branches y ejecuta los nodos destino concurrentemente:
graph.add_edge(START, "search_web")
graph.add_edge(START, "search_academic")
graph.add_edge(START, "search_news")
Estas 3 líneas dicen: "cuando el grafo inicie, lanza search_web, search_academic y search_news al mismo tiempo." No necesitas threading, asyncio ni ninguna librería de concurrencia — LangGraph lo maneja internamente.
Cómo funciona bajo el capó
LangGraph usa un scheduler interno. Cuando un nodo tiene múltiples edges de salida que van a nodos diferentes, o cuando múltiples edges apuntan al mismo punto de origen (como START), el scheduler:
- Identifica qué nodos pueden ejecutarse sin dependencias pendientes
- Los lanza en paralelo (usando threads internamente)
- Espera a que todos terminen antes de avanzar a nodos que dependan de ellos
No controlas el threading directamente. Solo declaras la estructura del grafo y LangGraph optimiza la ejecución.
Fan-in: múltiples nodos convergen en uno solo
Fan-in es el complemento del fan-out: múltiples nodos paralelos dirigen sus edges hacia un único nodo de merge:
graph.add_edge("search_web", "merge_results")
graph.add_edge("search_academic", "merge_results")
graph.add_edge("search_news", "merge_results")
El nodo merge_results se ejecuta solo cuando TODOS los nodos upstream han completado. No necesitas sincronización manual — LangGraph garantiza que el merge recibe todos los resultados.
La regla del merge
El nodo de merge se ejecuta una sola vez, cuando el último nodo paralelo termina. Si search_web tarda 1s, search_academic tarda 2s, y search_news tarda 3s, el merge se ejecuta a los 3 segundos — esperó al más lento. Pero el tiempo total es 3s, no 1+2+3=6s.
Estado diseñado para branching: Annotated con operator.add
Cuando múltiples nodos escriben al mismo campo del estado, necesitas un reducer que combine los valores en vez de sobrescribirlos. Sin reducer, el último nodo en terminar sobreescribiría los resultados de los anteriores:
from typing import TypedDict, Annotated
import operator
class ResearchState(TypedDict):
query: str
results: Annotated[list[str], operator.add]
Annotated[list[str], operator.add] dice: "cuando un nodo retorne {"results": [...]}, concatena la lista nueva con la existente en vez de reemplazarla." Así cada nodo paralelo agrega sus resultados y al final tienes todos.
| Sin reducer | Con operator.add |
|---|---|
search_web escribe ["web result"] | search_web agrega ["web result"] |
search_academic sobreescribe con ["academic result"] | search_academic agrega ["academic result"] |
Resultado final: ["academic result"] (se perdió web) | Resultado final: ["web result", "academic result"] |
Este diseño ya lo conoces del Módulo 5 (estado tipado con TypedDict y Annotated). Aquí lo aplicas en un contexto donde es crítico: sin el reducer, branching pierde datos.
Ejemplo completo: Research Agent con búsqueda paralela
Ahora juntamos fan-out + fan-in + estado con reducer:
from dotenv import load_dotenv
load_dotenv()
from typing import TypedDict, Annotated
import operator
import time
from langgraph.graph import StateGraph, START, END
from langchain.chat_models import init_chat_model
from IPython.display import Image, display
class ResearchState(TypedDict):
query: str
results: Annotated[list[dict], operator.add]
def search_web(state: ResearchState) -> dict:
time.sleep(1)
return {"results": [
{"source": "web", "content": f"Resultados web sobre '{state['query']}': Python 3.12 introduce mejoras en tipado, rendimiento de comprehensions y f-strings mejorados."}
]}
def search_academic(state: ResearchState) -> dict:
time.sleep(1)
return {"results": [
{"source": "academic", "content": f"Paper sobre '{state['query']}': Estudio comparativo de type checkers estáticos muestra mejoras del 15% en detección de bugs con Python 3.12."}
]}
def search_news(state: ResearchState) -> dict:
time.sleep(1)
return {"results": [
{"source": "news", "content": f"Noticias sobre '{state['query']}': Python mantiene posición #1 en índice TIOBE 2026, impulsado por adopción en AI/ML."}
]}
def merge_results(state: ResearchState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
sources_text = "\n\n".join(
f"[{r['source'].upper()}]: {r['content']}" for r in state["results"]
)
response = model.invoke(
f"Sintetiza estas fuentes en un resumen ejecutivo de 3-4 oraciones sobre '{state['query']}':\n\n{sources_text}"
)
return {"results": [{"source": "synthesis", "content": response.content}]}
graph = StateGraph(ResearchState)
graph.add_node("search_web", search_web)
graph.add_node("search_academic", search_academic)
graph.add_node("search_news", search_news)
graph.add_node("merge_results", merge_results)
graph.add_edge(START, "search_web")
graph.add_edge(START, "search_academic")
graph.add_edge(START, "search_news")
graph.add_edge("search_web", "merge_results")
graph.add_edge("search_academic", "merge_results")
graph.add_edge("search_news", "merge_results")
graph.add_edge("merge_results", END)
app = graph.compile()
display(Image(app.get_graph().draw_mermaid_png()))
start = time.time()
result = app.invoke({"query": "Python 3.12", "results": []})
elapsed = time.time() - start
print(f"Tiempo total: {elapsed:.1f}s (vs ~3.0s secuencial)")
print(f"Resultados recolectados: {len(result['results'])}")
for r in result["results"]:
print(f" [{r['source']}] {r['content'][:70]}...")
# Output esperado:
# Tiempo total: ~1.5s (vs ~3.0s secuencial)
# Resultados recolectados: 4
# [web] Resultados web sobre 'Python 3.12': Python 3.12 introduce mejor...
# [academic] Paper sobre 'Python 3.12': Estudio comparativo de type chec...
# [news] Noticias sobre 'Python 3.12': Python mantiene posición #1 en í...
# [synthesis] Python 3.12 ha consolidado mejoras significativas en tipado...
Anatomía del patrón
- Fan-out: 3 edges desde
STARThacia los 3 nodos de búsqueda - Ejecución paralela: LangGraph ejecuta los 3 nodos concurrentemente
- Fan-in: 3 edges desde los nodos de búsqueda hacia
merge_results - Reducer:
Annotated[list[dict], operator.add]acumula resultados de todos los nodos - Merge:
merge_resultsrecibe el estado con las 3 búsquedas y sintetiza
El diagrama que genera draw_mermaid_png() muestra claramente la forma de diamante: un punto de divergencia (START), 3 caminos paralelos, y un punto de convergencia (merge_results).
Fan-out con nodo planificador
En el ejemplo anterior, el fan-out sale directamente de START. Pero en un workflow real, a menudo hay un nodo planificador que decide qué buscar antes de lanzar las búsquedas:
from dotenv import load_dotenv
load_dotenv()
from typing import TypedDict, Annotated
import operator
import time
from langgraph.graph import StateGraph, START, END
from IPython.display import Image, display
class PlanState(TypedDict):
query: str
plan: str
results: Annotated[list[str], operator.add]
def planner(state: PlanState) -> dict:
return {"plan": f"Buscar '{state['query']}' en web, academic y news simultáneamente"}
def search_web(state: PlanState) -> dict:
time.sleep(1)
return {"results": [f"[Web] Info sobre '{state['query']}'"]}
def search_academic(state: PlanState) -> dict:
time.sleep(1)
return {"results": [f"[Academic] Papers sobre '{state['query']}'"]}
def search_news(state: PlanState) -> dict:
time.sleep(1)
return {"results": [f"[News] Noticias sobre '{state['query']}'"]}
def synthesizer(state: PlanState) -> dict:
combined = " | ".join(state["results"])
return {"results": [f"[Síntesis] {combined}"]}
graph = StateGraph(PlanState)
graph.add_node("planner", planner)
graph.add_node("search_web", search_web)
graph.add_node("search_academic", search_academic)
graph.add_node("search_news", search_news)
graph.add_node("synthesizer", synthesizer)
graph.add_edge(START, "planner")
graph.add_edge("planner", "search_web")
graph.add_edge("planner", "search_academic")
graph.add_edge("planner", "search_news")
graph.add_edge("search_web", "synthesizer")
graph.add_edge("search_academic", "synthesizer")
graph.add_edge("search_news", "synthesizer")
graph.add_edge("synthesizer", END)
app = graph.compile()
display(Image(app.get_graph().draw_mermaid_png()))
start = time.time()
result = app.invoke({"query": "LangGraph branching", "plan": "", "results": []})
elapsed = time.time() - start
print(f"Tiempo: {elapsed:.1f}s")
print(f"Plan: {result['plan']}")
print(f"Resultados: {len(result['results'])}")
# Output esperado:
# Tiempo: ~1.1s
# Plan: Buscar 'LangGraph branching' en web, academic y news simultáneamente
# Resultados: 4
El flujo es: START → planner → [search_web, search_academic, search_news] → synthesizer → END. El planificador ejecuta primero, luego las 3 búsquedas en paralelo, y finalmente el sintetizador.
Estrategias de merge
El nodo de merge recibe todos los resultados. Pero "juntar resultados" puede significar cosas diferentes según tu caso de uso:
Concatenación simple
La estrategia más básica. Junta todos los resultados en una lista:
def merge_concatenate(state: ResearchState) -> dict:
all_content = "\n\n".join(r["content"] for r in state["results"])
return {"results": [{"source": "merged", "content": all_content}]}
Deduplicación
Cuando múltiples fuentes pueden encontrar la misma información:
def merge_deduplicate(state: ResearchState) -> dict:
seen = set()
unique_results = []
for r in state["results"]:
content_key = r["content"][:100]
if content_key not in seen:
seen.add(content_key)
unique_results.append(r)
return {"results": unique_results}
Ranking por confianza
Cuando cada fuente tiene un score de relevancia:
def merge_ranked(state: ResearchState) -> dict:
source_priority = {"academic": 3, "web": 2, "news": 1}
sorted_results = sorted(
state["results"],
key=lambda r: source_priority.get(r["source"], 0),
reverse=True,
)
return {"results": sorted_results}
Síntesis con LLM
La estrategia más sofisticada — usa un LLM para combinar y sintetizar:
from langchain.chat_models import init_chat_model
def merge_synthesize(state: ResearchState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
context = "\n\n".join(f"[{r['source']}]: {r['content']}" for r in state["results"])
response = model.invoke(
f"Sintetiza estas fuentes en un resumen coherente:\n\n{context}"
)
return {"results": [{"source": "synthesis", "content": response.content}]}
| Estrategia | Cuándo usarla | Costo |
|---|---|---|
| Concatenación | Resultados ya están limpios, solo necesitas juntarlos | Mínimo |
| Deduplicación | Fuentes con overlap (web + news pueden tener lo mismo) | Bajo |
| Ranking | Necesitas priorizar fuentes más confiables | Bajo |
| Síntesis LLM | Necesitas un resumen coherente de fuentes diversas | Alto (llamada al LLM) |
Fan-out/fan-in con Functional API: Futures en paralelo
En el Módulo 6, Pattern 3, viste cómo la Functional API ejecuta tasks en paralelo con Futures. Ese mismo patrón es el equivalente del fan-out/fan-in de la Graph API:
from dotenv import load_dotenv
load_dotenv()
import time
from langgraph.func import entrypoint, task
@task
def search_web_fn(query: str) -> dict:
time.sleep(1)
return {"source": "web", "content": f"Web results for '{query}'"}
@task
def search_academic_fn(query: str) -> dict:
time.sleep(1)
return {"source": "academic", "content": f"Academic papers on '{query}'"}
@task
def search_news_fn(query: str) -> dict:
time.sleep(1)
return {"source": "news", "content": f"News about '{query}'"}
@task
def synthesize_fn(results: list) -> str:
return f"Síntesis de {len(results)} fuentes: " + " | ".join(r["content"] for r in results)
@entrypoint()
def parallel_research(query: str) -> dict:
web_future = search_web_fn(query)
academic_future = search_academic_fn(query)
news_future = search_news_fn(query)
results = [
web_future.result(),
academic_future.result(),
news_future.result(),
]
summary = synthesize_fn(results).result()
return {"results": results, "summary": summary}
start = time.time()
result = parallel_research.invoke("transformer architectures")
elapsed = time.time() - start
print(f"Tiempo: {elapsed:.1f}s")
print(f"Fuentes: {len(result['results'])}")
print(f"Síntesis: {result['summary'][:80]}...")
# Output esperado:
# Tiempo: ~1.1s
# Fuentes: 3
# Síntesis: Síntesis de 3 fuentes: Web results for 'transformer architectures'...
La mecánica es idéntica
# Fan-out: lanzar sin esperar
web_future = search_web_fn(query) # Lanza (no espera)
academic_future = search_academic_fn(query) # Lanza (no espera)
news_future = search_news_fn(query) # Lanza (no espera)
# Fan-in: recoger todos los resultados
results = [
web_future.result(), # Espera a web
academic_future.result(), # Espera a academic
news_future.result(), # Espera a news
]
El "fan-out" es lanzar múltiples @task sin .result(). El "fan-in" es recoger todos los .result() después. Es exactamente lo que viste en el Módulo 6, pero ahora con el vocabulario correcto: fan-out y fan-in.
Comparación: Graph API branching vs Functional API Futures
| Aspecto | Graph API (branching) | Functional API (Futures) |
|---|---|---|
| Cómo declaras paralelismo | Edges de un nodo a múltiples nodos | Lanzar @task sin .result() |
| Cómo declaras merge | Edges de múltiples nodos a uno | Recoger .result() en una lista |
| Estado compartido | TypedDict con reducer (operator.add) | Variables locales del @entrypoint |
| Visualización | draw_mermaid_png() muestra el diamante | No hay visualización nativa |
| Scheduler | Automático (LangGraph detecta el paralelismo) | Automático (Futures se ejecutan en background) |
| Checkpointing | Por nodo (cada nodo guarda su resultado) | Por @task (cada task guarda su resultado) |
| Manejo de errores | Nodo de merge recibe lo que haya completado | try/except al llamar .result() |
| Escalabilidad | Agregar nodos = agregar edges | Agregar tasks = agregar líneas |
¿Cuándo usar cada uno?
- ✅ Graph API cuando necesitas visualización del flujo, el equipo no-técnico necesita entender el workflow, o tienes múltiples estrategias de merge
- ✅ Functional API cuando el paralelismo es simple (lanzar N tasks, recoger resultados), no necesitas visualización, o prefieres escribir Python nativo
Cuándo branching ayuda vs cuándo es overkill
Branching no es gratis. Tiene overhead de scheduling, sincronización y complejidad de estado. Úsalo cuando el beneficio supera el costo:
| Escenario | ¿Branching? | Por qué |
|---|---|---|
| 3 búsquedas de 3s cada una | ✅ Sí | 9s → 3s, 3x más rápido |
| 3 operaciones de 50ms | ❌ No | 150ms → 50ms, pero el overhead de threading puede ser mayor que el ahorro |
| 2 llamadas al LLM independientes | ✅ Sí | Llamadas LLM son I/O-bound, paralelo ayuda mucho |
| 1 operación CPU-intensive | ❌ No | Solo 1 operación, nada que paralelizar |
| 10+ búsquedas con rate limit | ⚠️ Depende | Paralelo es más rápido, pero puedes exceder rate limits |
La regla de oro: branching ayuda cuando tienes múltiples operaciones I/O-bound que no dependen entre sí. Llamadas a APIs, búsquedas en bases de datos, invocaciones a LLMs — todas son I/O-bound y se benefician de paralelismo.
Branching con retry: combinando patrones
En la cápsula 02 aprendiste retry con backoff exponencial para manejar APIs que fallan. Ahora combina retry con branching — cada búsqueda paralela puede tener su propio retry independiente:
from dotenv import load_dotenv
load_dotenv()
from typing import TypedDict, Annotated
import operator
import time
import random
from langgraph.graph import StateGraph, START, END
class RetryState(TypedDict):
query: str
results: Annotated[list[dict], operator.add]
def make_search_with_retry(source_name: str, fail_rate: float = 0.5):
def search(state: RetryState) -> dict:
max_retries = 3
for attempt in range(max_retries):
try:
if random.random() < fail_rate and attempt < max_retries - 1:
raise ConnectionError(f"{source_name} timeout")
time.sleep(0.5)
return {"results": [{
"source": source_name,
"content": f"Resultados de {source_name} para '{state['query']}'",
"attempts": attempt + 1,
}]}
except ConnectionError:
wait = (2 ** attempt) * 0.1 + random.uniform(0, 0.05)
time.sleep(wait)
return {"results": [{
"source": source_name,
"content": f"[FALLBACK] {source_name} no disponible después de {max_retries} intentos",
"attempts": max_retries,
}]}
return search
def merge_with_status(state: RetryState) -> dict:
successful = [r for r in state["results"] if not r["content"].startswith("[FALLBACK]")]
failed = [r for r in state["results"] if r["content"].startswith("[FALLBACK]")]
summary = f"Fuentes exitosas: {len(successful)}, fallidas: {len(failed)}"
return {"results": [{"source": "merge", "content": summary, "attempts": 0}]}
graph = StateGraph(RetryState)
graph.add_node("search_web", make_search_with_retry("web", fail_rate=0.3))
graph.add_node("search_academic", make_search_with_retry("academic", fail_rate=0.3))
graph.add_node("search_news", make_search_with_retry("news", fail_rate=0.3))
graph.add_node("merge", merge_with_status)
graph.add_edge(START, "search_web")
graph.add_edge(START, "search_academic")
graph.add_edge(START, "search_news")
graph.add_edge("search_web", "merge")
graph.add_edge("search_academic", "merge")
graph.add_edge("search_news", "merge")
graph.add_edge("merge", END)
app = graph.compile()
result = app.invoke({"query": "LangGraph patterns", "results": []})
for r in result["results"]:
print(f" [{r['source']}] (intentos: {r['attempts']}) {r['content'][:60]}...")
# Output esperado (varía por aleatoridad):
# [web] (intentos: 2) Resultados de web para 'LangGraph patterns'...
# [academic] (intentos: 1) Resultados de academic para 'LangGraph patterns'...
# [news] (intentos: 3) [FALLBACK] news no disponible después de 3 intentos...
# [merge] (intentos: 0) Fuentes exitosas: 2, fallidas: 1...
Cada nodo de búsqueda maneja sus propios reintentos internamente. El merge recibe todos los resultados (exitosos y fallback) y puede decidir cómo proceder. Los retries ocurren en paralelo — si search_web necesita 2 intentos y search_news necesita 3, ambos reintentan al mismo tiempo.
Branching condicional: fan-out selectivo
No siempre quieres buscar en todas las fuentes. A veces un planificador decide qué fuentes consultar. Para eso combinas conditional edges con branching:
from dotenv import load_dotenv
load_dotenv()
from typing import TypedDict, Annotated
import operator
from langgraph.graph import StateGraph, START, END
from IPython.display import Image, display
class SelectiveState(TypedDict):
query: str
query_type: str
results: Annotated[list[str], operator.add]
def classifier(state: SelectiveState) -> dict:
q = state["query"].lower()
if any(w in q for w in ["paper", "estudio", "investigación"]):
return {"query_type": "academic"}
elif any(w in q for w in ["noticia", "hoy", "reciente"]):
return {"query_type": "news"}
return {"query_type": "general"}
def route_searches(state: SelectiveState) -> list[str]:
qt = state["query_type"]
if qt == "academic":
return ["search_web", "search_academic"]
elif qt == "news":
return ["search_web", "search_news"]
return ["search_web", "search_academic", "search_news"]
def search_web(state: SelectiveState) -> dict:
return {"results": [f"[Web] Resultados para '{state['query']}'"]}
def search_academic(state: SelectiveState) -> dict:
return {"results": [f"[Academic] Papers sobre '{state['query']}'"]}
def search_news(state: SelectiveState) -> dict:
return {"results": [f"[News] Noticias sobre '{state['query']}'"]}
def merge(state: SelectiveState) -> dict:
return {"results": [f"[Merge] Combinados {len(state['results'])} resultados"]}
graph = StateGraph(SelectiveState)
graph.add_node("classifier", classifier)
graph.add_node("search_web", search_web)
graph.add_node("search_academic", search_academic)
graph.add_node("search_news", search_news)
graph.add_node("merge", merge)
graph.add_edge(START, "classifier")
graph.add_conditional_edges("classifier", route_searches, ["search_web", "search_academic", "search_news"])
graph.add_edge("search_web", "merge")
graph.add_edge("search_academic", "merge")
graph.add_edge("search_news", "merge")
graph.add_edge("merge", END)
app = graph.compile()
display(Image(app.get_graph().draw_mermaid_png()))
for query in ["Busco papers sobre transformers", "¿Qué noticias hay hoy?", "Explica qué es RAG"]:
result = app.invoke({"query": query, "query_type": "", "results": []})
sources = [r for r in result["results"] if not r.startswith("[Merge]")]
print(f"'{query}' → {len(sources)} fuentes: {sources}")
# Output esperado:
# 'Busco papers sobre transformers' → 2 fuentes: ['[Web]...', '[Academic]...']
# '¿Qué noticias hay hoy?' → 2 fuentes: ['[Web]...', '[News]...']
# 'Explica qué es RAG' → 3 fuentes: ['[Web]...', '[Academic]...', '[News]...']
La función de routing retorna una lista de nodos en vez de un string. LangGraph ejecuta en paralelo solo los nodos de la lista. Los nodos no seleccionados no se ejecutan, pero sus edges al merge siguen existiendo — el merge espera solo a los nodos que sí ejecutaron.
Troubleshooting
Problema 1: "El nodo de merge se ejecuta antes de que todos los nodos paralelos terminen"
Síntoma: El merge recibe resultados incompletos — faltan fuentes.
Causa: Un nodo paralelo no tiene edge hacia el merge, o el nodo paralelo no retorna estado correctamente.
Solución: Verifica que cada nodo paralelo tenga un edge hacia el merge. Usa draw_mermaid_png() para confirmar visualmente:
display(Image(app.get_graph().draw_mermaid_png()))
Si un nodo aparece desconectado del merge en el diagrama, falta un add_edge.
Problema 2: "Los resultados de nodos paralelos se sobreescriben"
Síntoma: El estado final solo tiene el resultado de un nodo, no de todos.
Causa: El campo del estado no tiene un reducer. Sin Annotated[list, operator.add], el último nodo en escribir sobreescribe a los anteriores.
Solución: Agrega el reducer al campo que acumula resultados:
# ❌ Sin reducer (el último gana)
class State(TypedDict):
results: list[str]
# ✅ Con reducer (todos acumulan)
class State(TypedDict):
results: Annotated[list[str], operator.add]
Problema 3: "El orden de resultados es impredecible"
Síntoma: Los resultados llegan en orden diferente cada ejecución.
Causa: Los nodos paralelos terminan en orden no determinístico. El primero en terminar agrega primero al estado.
Solución: Si el orden importa, agrega metadata (timestamp, source name) y ordena en el merge:
def merge_ordered(state: ResearchState) -> dict:
source_order = {"web": 0, "academic": 1, "news": 2}
sorted_results = sorted(state["results"], key=lambda r: source_order.get(r["source"], 99))
return {"results": sorted_results}
Problema 4: "Un nodo paralelo falla y el grafo completo se cae"
Síntoma: Una excepción en un nodo de búsqueda detiene todo el workflow.
Causa: Las excepciones no capturadas en nodos se propagan al scheduler. En branching, un fallo en cualquier rama cancela todas.
Solución: Captura excepciones dentro de cada nodo y retorna un resultado de fallback:
def search_web(state: ResearchState) -> dict:
try:
result = do_actual_search(state["query"])
return {"results": [{"source": "web", "content": result, "status": "ok"}]}
except Exception as e:
return {"results": [{"source": "web", "content": str(e), "status": "error"}]}
Problema 5: "No veo paralelismo — el tiempo es igual al secuencial"
Síntoma: 3 nodos de 1s tardan 3s en total, no ~1s.
Causa: Los edges están encadenados en vez de hacer fan-out:
# ❌ Secuencial (encadenado)
graph.add_edge(START, "search_web")
graph.add_edge("search_web", "search_academic")
graph.add_edge("search_academic", "search_news")
# ✅ Paralelo (fan-out desde el mismo nodo)
graph.add_edge(START, "search_web")
graph.add_edge(START, "search_academic")
graph.add_edge(START, "search_news")
Revisa que todos los edges de fan-out salgan del mismo nodo.
Ejercicios
Ejercicio 1: Fan-out básico (Fácil)
Crea un grafo con 2 nodos paralelos ("uppercase" y "word_count") que procesen el mismo texto. uppercase convierte a mayúsculas, word_count cuenta palabras. Ambos escriben a un campo con reducer. Un nodo "report" al final combina los resultados. Visualiza.
Ver solución
from dotenv import load_dotenv
load_dotenv()
from typing import TypedDict, Annotated
import operator
from langgraph.graph import StateGraph, START, END
from IPython.display import Image, display
class State(TypedDict):
text: str
analyses: Annotated[list[dict], operator.add]
def uppercase(state: State) -> dict:
return {"analyses": [{"type": "uppercase", "result": state["text"].upper()}]}
def word_count(state: State) -> dict:
count = len(state["text"].split())
return {"analyses": [{"type": "word_count", "result": str(count)}]}
def report(state: State) -> dict:
results = {a["type"]: a["result"] for a in state["analyses"]}
summary = f"Texto: {results.get('uppercase', 'N/A')} | Palabras: {results.get('word_count', 'N/A')}"
return {"analyses": [{"type": "report", "result": summary}]}
graph = StateGraph(State)
graph.add_node("uppercase", uppercase)
graph.add_node("word_count", word_count)
graph.add_node("report", report)
graph.add_edge(START, "uppercase")
graph.add_edge(START, "word_count")
graph.add_edge("uppercase", "report")
graph.add_edge("word_count", "report")
graph.add_edge("report", END)
app = graph.compile()
display(Image(app.get_graph().draw_mermaid_png()))
result = app.invoke({"text": "hola mundo desde langgraph", "analyses": []})
for a in result["analyses"]:
print(f" [{a['type']}] {a['result']}")
# Output esperado:
# [uppercase] HOLA MUNDO DESDE LANGGRAPH
# [word_count] 4
# [report] Texto: HOLA MUNDO DESDE LANGGRAPH | Palabras: 4
El diagrama muestra el fan-out desde START a 2 nodos, y el fan-in hacia report.
Ejercicio 2: Merge con deduplicación (Fácil)
Crea 3 nodos paralelos de búsqueda que retornen resultados con overlap (2 de los 3 nodos retornan el mismo resultado). El nodo de merge debe deduplicar por contenido antes de reportar. Usa operator.add en el estado.
Ver solución
from dotenv import load_dotenv
load_dotenv()
from typing import TypedDict, Annotated
import operator
from langgraph.graph import StateGraph, START, END
class State(TypedDict):
query: str
results: Annotated[list[dict], operator.add]
def source_a(state: State) -> dict:
return {"results": [
{"content": "Python es un lenguaje de programación", "source": "A"},
{"content": "Python fue creado por Guido van Rossum", "source": "A"},
]}
def source_b(state: State) -> dict:
return {"results": [
{"content": "Python es un lenguaje de programación", "source": "B"},
{"content": "Python es popular en data science", "source": "B"},
]}
def source_c(state: State) -> dict:
return {"results": [
{"content": "Python fue creado por Guido van Rossum", "source": "C"},
{"content": "Python 3.12 es la última versión", "source": "C"},
]}
def merge_dedup(state: State) -> dict:
seen_content = set()
unique = []
for r in state["results"]:
if r["content"] not in seen_content:
seen_content.add(r["content"])
unique.append(r)
return {"results": [{"content": f"Resultados únicos: {len(unique)} de {len(state['results'])} totales", "source": "merge"}]}
graph = StateGraph(State)
graph.add_node("source_a", source_a)
graph.add_node("source_b", source_b)
graph.add_node("source_c", source_c)
graph.add_node("merge", merge_dedup)
graph.add_edge(START, "source_a")
graph.add_edge(START, "source_b")
graph.add_edge(START, "source_c")
graph.add_edge("source_a", "merge")
graph.add_edge("source_b", "merge")
graph.add_edge("source_c", "merge")
graph.add_edge("merge", END)
app = graph.compile()
result = app.invoke({"query": "Python", "results": []})
for r in result["results"]:
print(f" [{r['source']}] {r['content']}")
# Output esperado:
# [A] Python es un lenguaje de programación
# [A] Python fue creado por Guido van Rossum
# [B] Python es popular en data science
# [C] Python 3.12 es la última versión
# [merge] Resultados únicos: 4 de 6 totales
6 resultados totales, pero 2 están duplicados ("Python es un lenguaje de programación" y "Python fue creado por Guido van Rossum"). El merge los reduce a 4 únicos.
Ejercicio 3: Fan-out con error handling por rama (Medio)
Crea 3 nodos de búsqueda paralelos donde uno siempre falla (raise exception). Cada nodo debe manejar su propio error internamente y retornar un resultado con status "error" o "ok". El merge reporta cuántas fuentes fueron exitosas.
Ver solución
from dotenv import load_dotenv
load_dotenv()
from typing import TypedDict, Annotated
import operator
from langgraph.graph import StateGraph, START, END
class State(TypedDict):
query: str
results: Annotated[list[dict], operator.add]
def safe_search(source_name: str, should_fail: bool = False):
def search(state: State) -> dict:
try:
if should_fail:
raise ConnectionError(f"{source_name} is down")
return {"results": [{
"source": source_name,
"content": f"Results from {source_name} for '{state['query']}'",
"status": "ok",
}]}
except Exception as e:
return {"results": [{
"source": source_name,
"content": str(e),
"status": "error",
}]}
return search
def merge(state: State) -> dict:
ok = [r for r in state["results"] if r["status"] == "ok"]
errors = [r for r in state["results"] if r["status"] == "error"]
summary = f"{len(ok)} fuentes OK, {len(errors)} con error"
return {"results": [{"source": "merge", "content": summary, "status": "done"}]}
graph = StateGraph(State)
graph.add_node("search_a", safe_search("source_a"))
graph.add_node("search_b", safe_search("source_b", should_fail=True))
graph.add_node("search_c", safe_search("source_c"))
graph.add_node("merge", merge)
graph.add_edge(START, "search_a")
graph.add_edge(START, "search_b")
graph.add_edge(START, "search_c")
graph.add_edge("search_a", "merge")
graph.add_edge("search_b", "merge")
graph.add_edge("search_c", "merge")
graph.add_edge("merge", END)
app = graph.compile()
result = app.invoke({"query": "test", "results": []})
for r in result["results"]:
print(f" [{r['source']}] ({r['status']}) {r['content']}")
# Output esperado:
# [source_a] (ok) Results from source_a for 'test'
# [source_b] (error) source_b is down
# [source_c] (ok) Results from source_c for 'test'
# [merge] (done) 2 fuentes OK, 1 con error
El nodo search_b falla, pero captura la excepción internamente y retorna un resultado con status "error". El grafo no se cae — el merge recibe los 3 resultados y puede decidir cómo proceder.
Ejercicio 4: Functional API con merge personalizado (Medio)
Implementa el mismo patrón de búsqueda paralela con la Functional API. 3 tasks de búsqueda en paralelo, recolección con error handling, y una task de merge que ordena resultados por prioridad de fuente (academic > web > news).
Ver solución
from dotenv import load_dotenv
load_dotenv()
import time
from langgraph.func import entrypoint, task
@task
def search_web(query: str) -> dict:
time.sleep(0.5)
return {"source": "web", "content": f"Web results for '{query}'", "priority": 2}
@task
def search_academic(query: str) -> dict:
time.sleep(0.5)
return {"source": "academic", "content": f"Papers on '{query}'", "priority": 3}
@task
def search_news(query: str) -> dict:
time.sleep(0.5)
return {"source": "news", "content": f"News about '{query}'", "priority": 1}
@task
def merge_ranked(results: list) -> dict:
sorted_results = sorted(results, key=lambda r: r["priority"], reverse=True)
return {
"ranked_results": sorted_results,
"order": [r["source"] for r in sorted_results],
}
@entrypoint()
def ranked_research(query: str) -> dict:
web_fut = search_web(query)
academic_fut = search_academic(query)
news_fut = search_news(query)
results = []
errors = []
for name, fut in [("web", web_fut), ("academic", academic_fut), ("news", news_fut)]:
try:
results.append(fut.result())
except Exception as e:
errors.append({"source": name, "error": str(e)})
merged = merge_ranked(results).result()
return {
"results": merged["ranked_results"],
"order": merged["order"],
"errors": errors,
}
start = time.time()
result = ranked_research.invoke("machine learning")
elapsed = time.time() - start
print(f"Tiempo: {elapsed:.1f}s")
print(f"Orden de prioridad: {result['order']}")
for r in result["results"]:
print(f" [{r['source']}] (prioridad {r['priority']}) {r['content']}")
# Output esperado:
# Tiempo: ~0.6s
# Orden de prioridad: ['academic', 'web', 'news']
# [academic] (prioridad 3) Papers on 'machine learning'
# [web] (prioridad 2) Web results for 'machine learning'
# [news] (prioridad 1) News about 'machine learning'
La Functional API logra lo mismo con menos código. El try/except al recoger futuros maneja errores individuales sin detener las demás búsquedas.
Ejercicio 5: Research Agent completo con planner + branching + LLM merge (Avanzado)
Construye un Research Agent con Graph API que tenga: (1) un nodo planner que analiza la query, (2) fan-out a 3 fuentes de búsqueda en paralelo, (3) un nodo de merge que usa un LLM para sintetizar los resultados en un resumen ejecutivo. Usa init_chat_model para el merge. Visualiza el grafo.
Ver solución
from dotenv import load_dotenv
load_dotenv()
from typing import TypedDict, Annotated
import operator
import time
from langgraph.graph import StateGraph, START, END
from langchain.chat_models import init_chat_model
from IPython.display import Image, display
class AgentState(TypedDict):
query: str
plan: str
results: Annotated[list[dict], operator.add]
summary: str
def planner(state: AgentState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
response = model.invoke(
f"Analiza esta query de investigación y genera un plan breve (1-2 oraciones) "
f"de cómo buscar información:\n\n{state['query']}"
)
return {"plan": response.content}
def search_web(state: AgentState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
response = model.invoke(
f"Simula una búsqueda web. Genera un párrafo informativo sobre: {state['query']}"
)
return {"results": [{"source": "web", "content": response.content}]}
def search_academic(state: AgentState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
response = model.invoke(
f"Simula una búsqueda académica. Genera un abstract técnico sobre: {state['query']}"
)
return {"results": [{"source": "academic", "content": response.content}]}
def search_news(state: AgentState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
response = model.invoke(
f"Simula una búsqueda de noticias. Genera un resumen periodístico sobre: {state['query']}"
)
return {"results": [{"source": "news", "content": response.content}]}
def merge_synthesize(state: AgentState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
context = "\n\n".join(
f"[{r['source'].upper()}]:\n{r['content']}" for r in state["results"]
)
response = model.invoke(
f"Plan de investigación: {state['plan']}\n\n"
f"Fuentes recolectadas:\n\n{context}\n\n"
f"Genera un resumen ejecutivo de 4-5 oraciones que sintetice las fuentes."
)
return {"summary": response.content}
graph = StateGraph(AgentState)
graph.add_node("planner", planner)
graph.add_node("search_web", search_web)
graph.add_node("search_academic", search_academic)
graph.add_node("search_news", search_news)
graph.add_node("merge", merge_synthesize)
graph.add_edge(START, "planner")
graph.add_edge("planner", "search_web")
graph.add_edge("planner", "search_academic")
graph.add_edge("planner", "search_news")
graph.add_edge("search_web", "merge")
graph.add_edge("search_academic", "merge")
graph.add_edge("search_news", "merge")
graph.add_edge("merge", END)
app = graph.compile()
display(Image(app.get_graph().draw_mermaid_png()))
start = time.time()
result = app.invoke({"query": "¿Cómo impacta RAG en la precisión de los LLMs?", "plan": "", "results": [], "summary": ""})
elapsed = time.time() - start
print(f"Tiempo: {elapsed:.1f}s")
print(f"Plan: {result['plan'][:100]}...")
print(f"Fuentes: {len(result['results'])}")
print(f"Resumen: {result['summary'][:200]}...")
# Output esperado:
# Tiempo: ~3-5s
# Plan: Buscar información sobre RAG y su impacto en la precisión...
# Fuentes: 3
# Resumen: RAG (Retrieval-Augmented Generation) mejora significativamente...
El diagrama muestra: START → planner → [search_web, search_academic, search_news] → merge → END. Es el Research Agent completo con planificación, búsqueda paralela y síntesis por LLM.
Ejercicio 6: Branching dinámico con número variable de fuentes (Avanzado)
Adapta el ejercicio anterior para que el planner decida cuántas y cuáles fuentes consultar. Si la query es técnica, busca en web + academic. Si es sobre eventos recientes, busca en web + news. Si es general, busca en las 3. Usa conditional edges que retornen una lista de nodos.
Ver solución
from dotenv import load_dotenv
load_dotenv()
from typing import TypedDict, Annotated
import operator
from langgraph.graph import StateGraph, START, END
from IPython.display import Image, display
class State(TypedDict):
query: str
query_type: str
results: Annotated[list[str], operator.add]
def classifier(state: State) -> dict:
q = state["query"].lower()
if any(w in q for w in ["paper", "estudio", "técnico", "algoritmo"]):
return {"query_type": "technical"}
elif any(w in q for w in ["noticia", "hoy", "reciente", "lanzamiento"]):
return {"query_type": "current"}
return {"query_type": "general"}
def route_to_sources(state: State) -> list[str]:
qt = state["query_type"]
if qt == "technical":
return ["search_web", "search_academic"]
elif qt == "current":
return ["search_web", "search_news"]
return ["search_web", "search_academic", "search_news"]
def search_web(state: State) -> dict:
return {"results": [f"[Web] Resultados para '{state['query']}'"]}
def search_academic(state: State) -> dict:
return {"results": [f"[Academic] Papers sobre '{state['query']}'"]}
def search_news(state: State) -> dict:
return {"results": [f"[News] Noticias sobre '{state['query']}'"]}
def merge(state: State) -> dict:
sources_used = len(state["results"])
return {"results": [f"[Merge] Sintetizados {sources_used} resultados"]}
graph = StateGraph(State)
graph.add_node("classifier", classifier)
graph.add_node("search_web", search_web)
graph.add_node("search_academic", search_academic)
graph.add_node("search_news", search_news)
graph.add_node("merge", merge)
graph.add_edge(START, "classifier")
graph.add_conditional_edges(
"classifier",
route_to_sources,
["search_web", "search_academic", "search_news"],
)
graph.add_edge("search_web", "merge")
graph.add_edge("search_academic", "merge")
graph.add_edge("search_news", "merge")
graph.add_edge("merge", END)
app = graph.compile()
display(Image(app.get_graph().draw_mermaid_png()))
test_queries = [
"Busco un paper sobre transformers",
"¿Qué noticias hay hoy sobre IA?",
"Explica qué es machine learning",
]
for query in test_queries:
result = app.invoke({"query": query, "query_type": "", "results": []})
sources = [r for r in result["results"] if not r.startswith("[Merge]")]
print(f"'{query}' → tipo: {result['query_type']}, fuentes: {len(sources)}")
for s in sources:
print(f" {s}")
# Output esperado:
# 'Busco un paper sobre transformers' → tipo: technical, fuentes: 2
# [Web] Resultados para 'Busco un paper sobre transformers'
# [Academic] Papers sobre 'Busco un paper sobre transformers'
# '¿Qué noticias hay hoy sobre IA?' → tipo: current, fuentes: 2
# [Web] Resultados para '¿Qué noticias hay hoy sobre IA?'
# [News] Noticias sobre '¿Qué noticias hay hoy sobre IA?'
# 'Explica qué es machine learning' → tipo: general, fuentes: 3
# [Web] Resultados para 'Explica qué es machine learning'
# [Academic] Papers sobre 'Explica qué es machine learning'
# [News] Noticias sobre 'Explica qué es machine learning'
La función route_to_sources retorna una lista de nodos. LangGraph ejecuta en paralelo solo los nodos seleccionados. El merge espera solo a las fuentes que realmente se ejecutaron.
Resumen
En esta cápsula aprendiste:
- Fan-out es cuando un nodo dispara múltiples nodos en paralelo — se declara con múltiples
add_edgedesde el mismo nodo origen - Fan-in es cuando múltiples nodos paralelos convergen en un único nodo de merge — se declara con múltiples
add_edgehacia el mismo nodo destino - El merge espera a todos: el nodo de convergencia solo se ejecuta cuando el último nodo paralelo termina
- Annotated con operator.add es crítico: sin un reducer, los nodos paralelos sobreescriben los resultados de los demás
- Estrategias de merge: concatenación (junta todo), deduplicación (elimina repetidos), ranking (ordena por prioridad), síntesis LLM (genera resumen coherente)
- En la Functional API, fan-out = lanzar
@tasksin.result(), fan-in = recoger.result()después - Branching condicional permite que un planner decida qué ramas ejecutar, retornando una lista de nodos destino
- Retry + branching se combinan: cada rama paralela puede tener su propio retry independiente
- Cuándo usar branching: operaciones I/O-bound independientes (APIs, LLMs, búsquedas). No para operaciones rápidas o CPU-bound
Próxima cápsula: Subgraphs — cómo encapsular flujos completos como nodos reutilizables. Así como una función puede llamar a otra función, un grafo puede contener otro grafo.
Recursos adicionales
- LangGraph Branching — How-to Guide — Fan-out y fan-in con ejemplos completos
- LangGraph Concepts: Nodes and Edges — Documentación de edges, incluyendo branching
- State Reducers — Cómo funcionan los reducers como
operator.add - Functional API Parallelism — Futures y paralelismo en Functional API
- LangGraph Visualization — Visualizar grafos con draw_mermaid_png
- Python operator module — Referencia de
operator.addy otros operadores - Concurrent Futures (Python docs) — Concepto análogo de Futures en Python estándar
Módulo 7 — LangChain & LangGraph: From Chains to Agents