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:

  1. Identifica qué nodos pueden ejecutarse sin dependencias pendientes
  2. Los lanza en paralelo (usando threads internamente)
  3. 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 reducerCon 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

  1. Fan-out: 3 edges desde START hacia los 3 nodos de búsqueda
  2. Ejecución paralela: LangGraph ejecuta los 3 nodos concurrentemente
  3. Fan-in: 3 edges desde los nodos de búsqueda hacia merge_results
  4. Reducer: Annotated[list[dict], operator.add] acumula resultados de todos los nodos
  5. Merge: merge_results recibe 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}]}
EstrategiaCuándo usarlaCosto
ConcatenaciónResultados ya están limpios, solo necesitas juntarlosMínimo
DeduplicaciónFuentes con overlap (web + news pueden tener lo mismo)Bajo
RankingNecesitas priorizar fuentes más confiablesBajo
Síntesis LLMNecesitas un resumen coherente de fuentes diversasAlto (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

AspectoGraph API (branching)Functional API (Futures)
Cómo declaras paralelismoEdges de un nodo a múltiples nodosLanzar @task sin .result()
Cómo declaras mergeEdges de múltiples nodos a unoRecoger .result() en una lista
Estado compartidoTypedDict con reducer (operator.add)Variables locales del @entrypoint
Visualizacióndraw_mermaid_png() muestra el diamanteNo hay visualización nativa
SchedulerAutomático (LangGraph detecta el paralelismo)Automático (Futures se ejecutan en background)
CheckpointingPor nodo (cada nodo guarda su resultado)Por @task (cada task guarda su resultado)
Manejo de erroresNodo de merge recibe lo que haya completadotry/except al llamar .result()
EscalabilidadAgregar nodos = agregar edgesAgregar 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❌ No150ms → 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❌ NoSolo 1 operación, nada que paralelizar
10+ búsquedas con rate limit⚠️ DependeParalelo 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_edge desde 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_edge hacia 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 @task sin .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

  1. LangGraph Branching — How-to Guide — Fan-out y fan-in con ejemplos completos
  2. LangGraph Concepts: Nodes and Edges — Documentación de edges, incluyendo branching
  3. State Reducers — Cómo funcionan los reducers como operator.add
  4. Functional API Parallelism — Futures y paralelismo en Functional API
  5. LangGraph Visualization — Visualizar grafos con draw_mermaid_png
  6. Python operator module — Referencia de operator.add y otros operadores
  7. Concurrent Futures (Python docs) — Concepto análogo de Futures en Python estándar

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