Módulo 8: Memoria y Persistencia

Durable Execution

Descripción de la cápsula

Durable execution significa que tu agente sobrevive a crashes, reinicios e interrupciones. Si el proceso muere a mitad de una ejecución, el agente resume exactamente donde se detuvo — sin trabajo repetido, sin datos perdidos, sin costos desperdiciados.

En las cápsulas anteriores configuraste checkpointing con MemorySaver y PostgresSaver. Aprendiste que cada paso del grafo se guarda como un checkpoint. Pero hay una pregunta que no respondimos: ¿qué pasa cuando el proceso se cae a mitad de la ejecución?

La respuesta es durable execution. El checkpointer ya tiene guardado el progreso hasta el último nodo completado. Cuando el proceso se reinicia y usas el mismo thread_id, LangGraph detecta el checkpoint existente y resume desde donde quedó. No repite trabajo. No pierde datos. No gasta dinero de más.

Esto no es una feature de demo — es la diferencia entre un prototipo que "funciona en mi laptop" y un sistema que puedes poner en producción con confianza.


El escenario que lo hace tangible

Piensa en tu Research Assistant haciendo una investigación larga:

Investigación: "Estado del arte en AI Agents"
  → Paso 1: Buscar en Wikipedia       ($0.10 en API calls)  ✅ completado
  → Paso 2: Buscar en arXiv           ($0.10 en API calls)  ✅ completado
  → Paso 3: Buscar en noticias        ($0.10 en API calls)  ✅ completado
  → Paso 4: Buscar en blogs técnicos  ($0.10 en API calls)  💥 CRASH del proceso
  → Paso 5: Sintetizar reporte        ($0.10 en API calls)  ⏳ pendiente

Sin durable execution:

  • El proceso se reinicia
  • No hay forma de saber que ya hiciste los pasos 1-3
  • Repites TODO desde el inicio: 5 pasos × $0.10 = $0.50 desperdiciados
  • El usuario espera otros 5 minutos

Con durable execution:

  • El proceso se reinicia
  • Usas el mismo thread_id
  • LangGraph detecta: "Ya completé hasta el paso 3"
  • Resume desde el paso 4: $0.00 desperdiciados
  • El usuario espera solo 2 minutos (pasos 4 y 5)

Multiplica esto por 100 usuarios diarios, investigaciones de 10+ pasos, y APIs que cuestan $0.50+ por llamada. Durable execution no es un "nice to have" — es dinero real y tiempo real que ahorras.


Cómo funciona durable execution

El mecanismo es sorprendentemente simple porque ya tienes todos los componentes:

Ejecución normal:
  Nodo 1 → checkpoint guardado → Nodo 2 → checkpoint guardado → Nodo 3 → checkpoint guardado

Crash después del Nodo 3:
  [proceso muere]

Reinicio con mismo thread_id:
  LangGraph pregunta: "¿Hay un checkpoint para este thread_id?"
  → Sí: checkpoint del Nodo 3
  → "¿La ejecución estaba completa?"
  → No: faltaban Nodos 4 y 5
  → Resume desde Nodo 4

Cada checkpoint contiene:

  • ✅ El estado completo del grafo en ese punto
  • ✅ Qué nodo acaba de ejecutarse
  • ✅ Qué nodo sigue en la ejecución
  • ✅ El historial completo de mensajes y datos acumulados

Implementación: checkpointing + thread management

Durable execution no requiere código nuevo. Es checkpointing (que ya conoces) usado correctamente:

from dotenv import load_dotenv
load_dotenv()

import time
from typing import TypedDict, Annotated
import operator
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver

class ResearchState(TypedDict):
    topic: str
    sources: Annotated[list[str], operator.add]
    steps_completed: Annotated[list[str], operator.add]
    final_report: str

def search_wikipedia(state: ResearchState) -> dict:
    time.sleep(0.5)
    result = f"Wikipedia: información enciclopédica sobre '{state['topic']}'"
    return {
        "sources": [result],
        "steps_completed": ["wikipedia"],
    }

def search_arxiv(state: ResearchState) -> dict:
    time.sleep(0.5)
    result = f"arXiv: papers académicos sobre '{state['topic']}'"
    return {
        "sources": [result],
        "steps_completed": ["arxiv"],
    }

def search_news(state: ResearchState) -> dict:
    time.sleep(0.5)
    result = f"News: noticias recientes sobre '{state['topic']}'"
    return {
        "sources": [result],
        "steps_completed": ["news"],
    }

def synthesize(state: ResearchState) -> dict:
    report = f"Reporte sobre '{state['topic']}' basado en {len(state['sources'])} fuentes:\n"
    for i, source in enumerate(state["sources"], 1):
        report += f"  {i}. {source}\n"
    report += f"Pasos completados: {', '.join(state['steps_completed'])}"
    return {
        "final_report": report,
        "steps_completed": ["synthesis"],
    }

graph_builder = StateGraph(ResearchState)
graph_builder.add_node("wikipedia", search_wikipedia)
graph_builder.add_node("arxiv", search_arxiv)
graph_builder.add_node("news", search_news)
graph_builder.add_node("synthesize", synthesize)

graph_builder.add_edge(START, "wikipedia")
graph_builder.add_edge("wikipedia", "arxiv")
graph_builder.add_edge("arxiv", "news")
graph_builder.add_edge("news", "synthesize")
graph_builder.add_edge("synthesize", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

config = {"configurable": {"thread_id": "research-001"}}

result = graph.invoke(
    {"topic": "AI agents", "sources": [], "steps_completed": [], "final_report": ""},
    config,
)
print("Ejecución completa:")
print(f"Pasos: {result['steps_completed']}")
print(f"Fuentes: {len(result['sources'])}")
print(f"Reporte:\n{result['final_report']}")
# Output esperado:
# Ejecución completa:
# Pasos: ['wikipedia', 'arxiv', 'news', 'synthesis']
# Fuentes: 3
# Reporte:
# Reporte sobre 'AI agents' basado en 3 fuentes:
#   1. Wikipedia: información enciclopédica sobre 'AI agents'
#   2. arXiv: papers académicos sobre 'AI agents'
#   3. News: noticias recientes sobre 'AI agents'
# Pasos completados: wikipedia, arxiv, news, synthesis

Lo clave: MemorySaver() + thread_id = durable execution. Cada nodo que completa guarda un checkpoint automáticamente. Si el proceso muere, la próxima invocación con el mismo thread_id encuentra el checkpoint y continúa.


Resume from checkpoint: la mecánica del reinicio

Cuando invocas un grafo con un thread_id que ya tiene checkpoints, LangGraph hace esta evaluación:

graph.invoke(input, {"configurable": {"thread_id": "research-001"}})

  1. ¿Existe un checkpoint para "research-001"?
     → No: ejecuta desde el inicio (ejecución nueva)
     → Sí: continúa al paso 2

  2. ¿La ejecución anterior terminó (llegó a END)?
     → Sí: esta es una nueva invocación sobre un thread existente
     → No: la ejecución anterior fue interrumpida → RESUME

Veamos cómo inspeccionar el estado de un thread para entender dónde quedó:

from dotenv import load_dotenv
load_dotenv()

from typing import TypedDict, Annotated
import operator
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver

class State(TypedDict):
    task: str
    progress: Annotated[list[str], operator.add]
    result: str

def step_1(state: State) -> dict:
    return {"progress": ["step_1_done"]}

def step_2(state: State) -> dict:
    return {"progress": ["step_2_done"]}

def step_3(state: State) -> dict:
    return {"progress": ["step_3_done"], "result": "Tarea completada."}

graph_builder = StateGraph(State)
graph_builder.add_node("step_1", step_1)
graph_builder.add_node("step_2", step_2)
graph_builder.add_node("step_3", step_3)

graph_builder.add_edge(START, "step_1")
graph_builder.add_edge("step_1", "step_2")
graph_builder.add_edge("step_2", "step_3")
graph_builder.add_edge("step_3", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

config = {"configurable": {"thread_id": "task-resume-demo"}}

result = graph.invoke(
    {"task": "Procesar datos", "progress": [], "result": ""},
    config,
)

current_state = graph.get_state(config)
print(f"Estado actual: {current_state.values['progress']}")
print(f"Siguiente nodo: {current_state.next}")
print(f"Resultado: {current_state.values['result']}")
# Output esperado:
# Estado actual: ['step_1_done', 'step_2_done', 'step_3_done']
# Siguiente nodo: ()
# Resultado: Tarea completada.

Cuando current_state.next es una tupla vacía (), la ejecución terminó normalmente. Si contiene un nombre de nodo, la ejecución fue interrumpida ahí — y puedes resumirla.


Simulando un crash y verificando el resume

Para entender el valor real, simulemos una interrupción. Usaremos interrupt de LangGraph para pausar la ejecución (simula un crash controlado), y luego verificaremos que el resume funciona:

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 interrupt
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import Command

class State(TypedDict):
    task: str
    steps_done: Annotated[list[str], operator.add]
    cost_usd: float

def expensive_step_1(state: State) -> dict:
    """Paso 1: cuesta $0.10 en API calls."""
    print("  Ejecutando paso 1 (costo: $0.10)...")
    return {
        "steps_done": ["paso_1"],
        "cost_usd": state.get("cost_usd", 0) + 0.10,
    }

def expensive_step_2(state: State) -> dict:
    """Paso 2: cuesta $0.10 en API calls."""
    print("  Ejecutando paso 2 (costo: $0.10)...")
    return {
        "steps_done": ["paso_2"],
        "cost_usd": state.get("cost_usd", 0) + 0.10,
    }

def crash_point(state: State) -> dict:
    """Simula un crash: el proceso se interrumpe aquí."""
    print("  💥 Interrupción en paso 3...")
    interrupt("Simulando crash del proceso")
    return {
        "steps_done": ["paso_3"],
        "cost_usd": state.get("cost_usd", 0) + 0.10,
    }

def final_step(state: State) -> dict:
    """Paso 4: genera el resultado final."""
    print("  Ejecutando paso 4 (costo: $0.10)...")
    return {
        "steps_done": ["paso_4"],
        "cost_usd": state.get("cost_usd", 0) + 0.10,
    }

graph_builder = StateGraph(State)
graph_builder.add_node("step_1", expensive_step_1)
graph_builder.add_node("step_2", expensive_step_2)
graph_builder.add_node("step_3", crash_point)
graph_builder.add_node("step_4", final_step)

graph_builder.add_edge(START, "step_1")
graph_builder.add_edge("step_1", "step_2")
graph_builder.add_edge("step_2", "step_3")
graph_builder.add_edge("step_3", "step_4")
graph_builder.add_edge("step_4", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

config = {"configurable": {"thread_id": "crash-demo-001"}}

print("=== PRIMERA EJECUCIÓN (se interrumpe en paso 3) ===")
result = graph.invoke(
    {"task": "Investigación larga", "steps_done": [], "cost_usd": 0.0},
    config,
)

state_after_crash = graph.get_state(config)
print(f"\nEstado después del crash:")
print(f"  Pasos completados: {state_after_crash.values['steps_done']}")
print(f"  Costo acumulado: ${state_after_crash.values['cost_usd']:.2f}")
print(f"  Siguiente nodo: {state_after_crash.next}")

print("\n=== RESUME (continúa desde donde quedó) ===")
result = graph.invoke(Command(resume="continuar"), config)

state_after_resume = graph.get_state(config)
print(f"\nEstado después del resume:")
print(f"  Pasos completados: {state_after_resume.values['steps_done']}")
print(f"  Costo acumulado: ${state_after_resume.values['cost_usd']:.2f}")
print(f"  Siguiente nodo: {state_after_resume.next}")
# Output esperado:
# === PRIMERA EJECUCIÓN (se interrumpe en paso 3) ===
#   Ejecutando paso 1 (costo: $0.10)...
#   Ejecutando paso 2 (costo: $0.10)...
#   💥 Interrupción en paso 3...
#
# Estado después del crash:
#   Pasos completados: ['paso_1', 'paso_2']
#   Costo acumulado: $0.20
#   Siguiente nodo: ('step_3',)
#
# === RESUME (continúa desde donde quedó) ===
#   Ejecutando paso 3 (costo: $0.10)...
#   Ejecutando paso 4 (costo: $0.10)...
#
# Estado después del resume:
#   Pasos completados: ['paso_1', 'paso_2', 'paso_3', 'paso_4']
#   Costo acumulado: $0.40
#   Siguiente nodo: ()

Observa lo que pasó:

  1. La primera ejecución completó pasos 1 y 2, y se interrumpió en el paso 3
  2. El estado registra que llevamos $0.20 gastados y 2 pasos completados
  3. next dice ('step_3',) — el grafo sabe exactamente dónde retomar
  4. El resume ejecuta pasos 3 y 4 sin repetir los pasos 1 y 2
  5. Costo total: $0.40 (sin repetición) en vez de $0.80 (si repitiéramos todo)

Agentes de larga duración

Algunos agentes ejecutan por minutos u horas. Un Research Assistant que analiza 50 papers, un agente que procesa 1000 registros de una base de datos, o un pipeline de generación que crea 20 secciones de un documento. Estos agentes son los que más se benefician de durable execution:

from dotenv import load_dotenv
load_dotenv()

import time
import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver

class BatchState(TypedDict):
    items_to_process: list[str]
    processed: Annotated[list[dict], operator.add]
    current_batch: int
    total_batches: int

def process_batch(state: BatchState) -> dict:
    """Procesa un batch de items. Cada batch es un checkpoint."""
    batch_num = state["current_batch"]
    batch_size = 3
    start_idx = batch_num * batch_size
    end_idx = min(start_idx + batch_size, len(state["items_to_process"]))

    items = state["items_to_process"][start_idx:end_idx]

    results = []
    for item in items:
        time.sleep(0.1)
        results.append({
            "item": item,
            "result": f"Procesado: {item}",
            "batch": batch_num,
        })

    return {
        "processed": results,
        "current_batch": batch_num + 1,
    }

def should_continue(state: BatchState) -> str:
    if state["current_batch"] >= state["total_batches"]:
        return "done"
    return "process"

def finalize(state: BatchState) -> dict:
    return {}

graph_builder = StateGraph(BatchState)
graph_builder.add_node("process", process_batch)
graph_builder.add_node("done", finalize)

graph_builder.add_edge(START, "process")
graph_builder.add_conditional_edges(
    "process", should_continue,
    {"process": "process", "done": "done"},
)
graph_builder.add_edge("done", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

items = [f"item_{i}" for i in range(12)]
config = {"configurable": {"thread_id": "batch-001"}}

result = graph.invoke(
    {
        "items_to_process": items,
        "processed": [],
        "current_batch": 0,
        "total_batches": 4,
    },
    config,
)

print(f"Items procesados: {len(result['processed'])}")
print(f"Batches completados: {result['current_batch']}")
for item in result["processed"]:
    print(f"  Batch {item['batch']}: {item['result']}")
# Output esperado:
# Items procesados: 12
# Batches completados: 4
#   Batch 0: Procesado: item_0
#   Batch 0: Procesado: item_1
#   Batch 0: Procesado: item_2
#   Batch 1: Procesado: item_3
#   Batch 1: Procesado: item_4
#   Batch 1: Procesado: item_5
#   Batch 2: Procesado: item_6
#   Batch 2: Procesado: item_7
#   Batch 2: Procesado: item_8
#   Batch 3: Procesado: item_9
#   Batch 3: Procesado: item_10
#   Batch 3: Procesado: item_11

El patrón clave: el grafo usa un loop con conditional edge que procesa un batch por iteración. Cada iteración genera un checkpoint. Si el proceso se cae después del batch 2, al reiniciar con el mismo thread_id, el grafo resume desde el batch 3 — los 6 items ya procesados no se repiten.


Nodos idempotentes: diseño seguro para re-ejecución

Hay un caso borde importante: ¿qué pasa si el proceso se cae durante la ejecución de un nodo, antes de que el checkpoint se guarde? En ese caso, el nodo se re-ejecutará al reiniciar. Por eso necesitas que tus nodos sean idempotentes — seguros de ejecutar más de una vez sin efectos duplicados.

from dotenv import load_dotenv
load_dotenv()

from typing import TypedDict, Annotated
import operator
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver

class State(TypedDict):
    task: str
    processed_ids: Annotated[list[str], operator.add]
    results: Annotated[list[str], operator.add]

def idempotent_process(state: State) -> dict:
    """Nodo idempotente: verifica si ya procesó antes de actuar."""
    task_id = "api-call-001"

    if task_id in state.get("processed_ids", []):
        print(f"  ⏭️  {task_id} ya procesado, saltando...")
        return {"processed_ids": [], "results": []}

    print(f"  🔄 Procesando {task_id}...")
    result = f"Resultado de {task_id}"

    return {
        "processed_ids": [task_id],
        "results": [result],
    }

def non_idempotent_process(state: State) -> dict:
    """Nodo NO idempotente: cada ejecución agrega un duplicado."""
    task_id = "api-call-002"
    print(f"  🔄 Procesando {task_id} (sin verificar duplicados)...")
    return {
        "processed_ids": [task_id],
        "results": [f"Resultado de {task_id}"],
    }

graph_builder = StateGraph(State)
graph_builder.add_node("safe", idempotent_process)
graph_builder.add_node("unsafe", non_idempotent_process)

graph_builder.add_edge(START, "safe")
graph_builder.add_edge("safe", "unsafe")
graph_builder.add_edge("unsafe", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

config = {"configurable": {"thread_id": "idempotent-demo"}}

print("=== Primera ejecución ===")
result = graph.invoke(
    {"task": "demo", "processed_ids": [], "results": []},
    config,
)
print(f"Resultados: {result['results']}")
print(f"IDs procesados: {result['processed_ids']}")
# Output esperado:
# === Primera ejecución ===
#   🔄 Procesando api-call-001...
#   🔄 Procesando api-call-002 (sin verificar duplicados)...
# Resultados: ['Resultado de api-call-001', 'Resultado de api-call-002']
# IDs procesados: ['api-call-001', 'api-call-002']

Patrones de idempotencia

OperaciónNo idempotenteIdempotente
API callLlama siempreVerifica si ya tienes el resultado
Insertar en DBINSERT (duplica)INSERT ... ON CONFLICT DO NOTHING
Enviar emailEnvía siempreVerifica flag email_sent en estado
Generar archivoSobrescribe siempreVerifica si el archivo ya existe
Cobrar pagoCobra siempreUsa idempotency_key del payment API

La regla: si un nodo tiene efectos secundarios (envía datos, cobra dinero, modifica bases de datos), hazlo idempotente. Los nodos que solo transforman datos internos del estado no necesitan protección adicional porque el checkpoint los maneja.


Retry automático a nivel de grafo

LangGraph te permite configurar retry policies directamente en los nodos. Esto es complementario a durable execution — un nodo que falla por un error transitorio se reintenta automáticamente antes de que el grafo lo considere un fallo:

from dotenv import load_dotenv
load_dotenv()

from typing import TypedDict
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver

CALL_COUNT = 0

class State(TypedDict):
    query: str
    result: str
    attempts: int

def flaky_api_call(state: State) -> dict:
    """Simula una API que falla intermitentemente."""
    global CALL_COUNT
    CALL_COUNT += 1

    if CALL_COUNT <= 2:
        raise ConnectionError(f"API timeout (intento #{CALL_COUNT})")

    return {
        "result": f"Datos obtenidos exitosamente en intento #{CALL_COUNT}",
        "attempts": CALL_COUNT,
    }

graph_builder = StateGraph(State)
graph_builder.add_node(
    "api_call",
    flaky_api_call,
    retry={"max_attempts": 4, "delay": 0.5, "multiplier": 2.0},
)

graph_builder.add_edge(START, "api_call")
graph_builder.add_edge("api_call", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

CALL_COUNT = 0
config = {"configurable": {"thread_id": "retry-demo"}}
result = graph.invoke({"query": "test", "result": "", "attempts": 0}, config)
print(f"Resultado: {result['result']}")
print(f"Intentos necesarios: {result['attempts']}")
# Output esperado:
# Resultado: Datos obtenidos exitosamente en intento #3
# Intentos necesarios: 3

El parámetro retry en add_node configura:

  • max_attempts: cuántas veces intentar (4 = 1 original + 3 retries)
  • delay: tiempo base entre intentos (en segundos)
  • multiplier: factor de backoff exponencial (2.0 = delay se duplica cada vez)

Esto es automático — no necesitas escribir lógica de retry en cada nodo. LangGraph lo maneja por ti. Si después de todos los reintentos el nodo sigue fallando, ahí sí se propaga el error para que tu error handling lo capture.


Verificación de checkpoint antes de operaciones costosas

Un patrón avanzado: antes de ejecutar una operación costosa (API call de $1, procesamiento de 10 minutos), verifica si el resultado ya existe en el estado:

from dotenv import load_dotenv
load_dotenv()

import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver

class State(TypedDict):
    topic: str
    cache: dict
    results: Annotated[list[str], operator.add]
    cost_saved: float

def expensive_analysis(state: State) -> dict:
    """Operación costosa: $1.00 por ejecución. Verifica caché primero."""
    cache_key = f"analysis_{state['topic']}"
    cache = state.get("cache", {})

    if cache_key in cache:
        print(f"  💰 Caché hit para '{cache_key}' — ahorramos $1.00")
        return {
            "results": [cache[cache_key]],
            "cost_saved": state.get("cost_saved", 0) + 1.00,
        }

    print(f"  🔄 Ejecutando análisis costoso para '{state['topic']}'...")
    result = f"Análisis profundo de '{state['topic']}': 15 papers revisados, 3 tendencias identificadas."

    new_cache = {**cache, cache_key: result}
    return {
        "results": [result],
        "cache": new_cache,
        "cost_saved": state.get("cost_saved", 0),
    }

def secondary_analysis(state: State) -> dict:
    print(f"  🔄 Análisis secundario...")
    return {"results": [f"Análisis complementario sobre '{state['topic']}'"]}

graph_builder = StateGraph(State)
graph_builder.add_node("expensive", expensive_analysis)
graph_builder.add_node("secondary", secondary_analysis)

graph_builder.add_edge(START, "expensive")
graph_builder.add_edge("expensive", "secondary")
graph_builder.add_edge("secondary", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

config = {"configurable": {"thread_id": "cost-aware-001"}}

print("=== Primera ejecución ===")
result1 = graph.invoke(
    {"topic": "transformer architectures", "cache": {}, "results": [], "cost_saved": 0},
    config,
)
print(f"Resultados: {len(result1['results'])}")
print(f"Costo ahorrado: ${result1['cost_saved']:.2f}")

config2 = {"configurable": {"thread_id": "cost-aware-002"}}

print("\n=== Segunda ejecución (con caché de la primera) ===")
result2 = graph.invoke(
    {
        "topic": "transformer architectures",
        "cache": result1["cache"],
        "results": [],
        "cost_saved": 0,
    },
    config2,
)
print(f"Resultados: {len(result2['results'])}")
print(f"Costo ahorrado: ${result2['cost_saved']:.2f}")
# Output esperado:
# === Primera ejecución ===
#   🔄 Ejecutando análisis costoso para 'transformer architectures'...
#   🔄 Análisis secundario...
# Resultados: 2
# Costo ahorrado: $0.00
#
# === Segunda ejecución (con caché de la primera) ===
#   💰 Caché hit para 'analysis_transformer architectures' — ahorramos $1.00
#   🔄 Análisis secundario...
# Resultados: 2
# Costo ahorrado: $1.00

El caché vive en el estado del grafo. Si estás usando PostgresSaver, el caché se persiste automáticamente con el checkpoint. La próxima ejecución con los mismos datos puede saltarse la operación costosa.


De MemorySaver a PostgresSaver: durabilidad real

MemorySaver es perfecto para desarrollo y testing, pero tiene una limitación fundamental: vive en memoria. Si el proceso se reinicia, los checkpoints se pierden. Para durable execution real en producción, necesitas PostgresSaver:

# Desarrollo: MemorySaver (rápido, sin dependencias)
from langgraph.checkpoint.memory import MemorySaver
checkpointer = MemorySaver()

# Producción: PostgresSaver (durable, sobrevive reinicios)
# pip install langgraph-checkpoint-postgres
from langgraph.checkpoint.postgres import PostgresSaver
checkpointer = PostgresSaver.from_conn_string(
    "postgresql://user:password@localhost:5432/mydb"
)

El cambio es una línea. Todo tu código del grafo permanece idéntico. Los nodos no saben ni les importa qué checkpointer estás usando.

MemorySaver                          PostgresSaver
  ├── Rápido (in-memory)               ├── Durable (disco/red)
  ├── Sin dependencias                  ├── Requiere PostgreSQL
  ├── Se pierde al reiniciar            ├── Sobrevive reinicios
  ├── Perfecto para dev/test            ├── Perfecto para producción
  └── Un solo proceso                   └── Multi-proceso/multi-servidor

Con PostgresSaver:

  • ✅ El proceso puede morir y reiniciar → los checkpoints están en PostgreSQL
  • ✅ Múltiples servidores pueden acceder a los mismos threads
  • ✅ Puedes inspeccionar checkpoints directamente en la base de datos
  • ✅ Backups de PostgreSQL = backup de todo el estado de tus agentes

Design patterns para agentes durables

Pattern 1: Progress tracking explícito

Siempre incluye un campo de progreso en tu estado para saber exactamente dónde está el agente:

from typing import TypedDict, Annotated
import operator

class DurableState(TypedDict):
    task: str
    steps_completed: Annotated[list[str], operator.add]
    total_steps: int
    current_step: int
    result: str

Pattern 2: Operaciones atómicas por nodo

Cada nodo debe hacer una sola cosa bien definida. Si un nodo hace 3 operaciones y falla en la segunda, las 3 se re-ejecutan. Sepáralas en nodos independientes:

# ❌ Un nodo hace demasiado
def do_everything(state):
    a = call_api_a()    # $0.50
    b = call_api_b()    # $0.50 ← falla aquí
    c = process(a, b)   # nunca se ejecuta
    return {"result": c}
# Si falla en api_b, la próxima vez repite api_a ($0.50 perdido)

# ✅ Cada nodo es una operación atómica
def call_a(state):
    return {"data_a": call_api_a()}  # checkpoint después de esto

def call_b(state):
    return {"data_b": call_api_b()}  # si falla, solo repite esto

def process(state):
    return {"result": process(state["data_a"], state["data_b"])}

Pattern 3: Checkpoint verification

Antes de una operación costosa, verifica si ya se completó:

def smart_node(state):
    if state.get("expensive_result"):
        return {}
    result = expensive_operation()
    return {"expensive_result": result}

Pattern 4: Graceful shutdown

Diseña tu agente para que pueda detenerse limpiamente en cualquier punto:

import signal

shutdown_requested = False

def handle_shutdown(signum, frame):
    global shutdown_requested
    shutdown_requested = True
    print("Shutdown solicitado, terminando después del nodo actual...")

signal.signal(signal.SIGTERM, handle_shutdown)
signal.signal(signal.SIGINT, handle_shutdown)

El checkpoint se guarda después de cada nodo, así que un shutdown entre nodos siempre deja el estado consistente.


Troubleshooting

Problema 1: "El grafo repite todos los pasos al reiniciar"

Síntoma: Usas el mismo thread_id pero el grafo ejecuta desde el inicio.

Causa: Estás usando MemorySaver y el proceso se reinició (los datos en memoria se perdieron).

Solución: Usa PostgresSaver para durabilidad real:

# MemorySaver pierde datos al reiniciar
checkpointer = MemorySaver()  # ← solo para desarrollo

# PostgresSaver persiste entre reinicios
from langgraph.checkpoint.postgres import PostgresSaver
checkpointer = PostgresSaver.from_conn_string("postgresql://...")

Problema 2: "El agente procesa items duplicados después de un crash"

Síntoma: Al reiniciar, el agente reprocesa items que ya había procesado.

Causa: El nodo no es idempotente — no verifica si el trabajo ya se hizo.

Solución: Agrega verificación de idempotencia:

def process_item(state):
    item_id = state["current_item_id"]
    if item_id in state.get("processed_ids", []):
        return {}  # ya procesado
    result = do_work(item_id)
    return {"processed_ids": [item_id], "results": [result]}

Problema 3: "El checkpoint ocupa demasiado espacio en PostgreSQL"

Síntoma: La tabla de checkpoints crece sin control.

Causa: Cada invocación de cada thread guarda múltiples checkpoints.

Solución: Implementa limpieza periódica de threads antiguos:

# Mantener solo los últimos N checkpoints por thread
# o eliminar threads completados hace más de X días
# Esto se maneja a nivel de base de datos con queries SQL

Problema 4: "El retry automático no se activa"

Síntoma: El nodo falla y el grafo termina con error inmediatamente.

Causa: La configuración de retry no está en add_node.

Solución: Verifica que retry está configurado correctamente:

# ❌ Sin retry
graph_builder.add_node("my_node", my_func)

# ✅ Con retry
graph_builder.add_node(
    "my_node", my_func,
    retry={"max_attempts": 3, "delay": 1.0, "multiplier": 2.0},
)

Ejercicios

Ejercicio 1: Durable execution básico (Fácil)

Crea un grafo con 4 nodos secuenciales que representan pasos de una investigación. Usa MemorySaver. Ejecuta el grafo y luego inspecciona el estado final con get_state() para verificar que todos los pasos se completaron. El estado debe incluir steps_completed y total_cost.

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.checkpoint.memory import MemorySaver

class State(TypedDict):
    topic: str
    steps_completed: Annotated[list[str], operator.add]
    total_cost: float

def search(state: State) -> dict:
    return {
        "steps_completed": ["search"],
        "total_cost": state.get("total_cost", 0) + 0.05,
    }

def analyze(state: State) -> dict:
    return {
        "steps_completed": ["analyze"],
        "total_cost": state.get("total_cost", 0) + 0.10,
    }

def synthesize(state: State) -> dict:
    return {
        "steps_completed": ["synthesize"],
        "total_cost": state.get("total_cost", 0) + 0.15,
    }

def format_report(state: State) -> dict:
    return {
        "steps_completed": ["format"],
        "total_cost": state.get("total_cost", 0) + 0.02,
    }

graph_builder = StateGraph(State)
graph_builder.add_node("search", search)
graph_builder.add_node("analyze", analyze)
graph_builder.add_node("synthesize", synthesize)
graph_builder.add_node("format", format_report)

graph_builder.add_edge(START, "search")
graph_builder.add_edge("search", "analyze")
graph_builder.add_edge("analyze", "synthesize")
graph_builder.add_edge("synthesize", "format")
graph_builder.add_edge("format", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

config = {"configurable": {"thread_id": "exercise-1"}}
result = graph.invoke(
    {"topic": "durable execution", "steps_completed": [], "total_cost": 0.0},
    config,
)

state = graph.get_state(config)
print(f"Pasos completados: {state.values['steps_completed']}")
print(f"Costo total: ${state.values['total_cost']:.2f}")
print(f"Ejecución completa: {state.next == ()}")
# Output esperado:
# Pasos completados: ['search', 'analyze', 'synthesize', 'format']
# Costo total: $0.32
# Ejecución completa: True

Ejercicio 2: Simulación de crash con interrupt (Fácil)

Crea un grafo de 3 nodos donde el segundo nodo usa interrupt() para simular un crash. Inspecciona el estado después de la interrupción (qué pasos se completaron, cuál es el siguiente nodo). Luego resume la ejecución con Command(resume=...).

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.checkpoint.memory import MemorySaver
from langgraph.types import interrupt, Command

class State(TypedDict):
    data: str
    log: Annotated[list[str], operator.add]

def step_a(state: State) -> dict:
    return {"log": ["step_a completado"]}

def step_b(state: State) -> dict:
    interrupt("Simulando crash en step_b")
    return {"log": ["step_b completado"]}

def step_c(state: State) -> dict:
    return {"log": ["step_c completado"], "data": "resultado final"}

graph_builder = StateGraph(State)
graph_builder.add_node("a", step_a)
graph_builder.add_node("b", step_b)
graph_builder.add_node("c", step_c)

graph_builder.add_edge(START, "a")
graph_builder.add_edge("a", "b")
graph_builder.add_edge("b", "c")
graph_builder.add_edge("c", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)
config = {"configurable": {"thread_id": "crash-exercise"}}

print("=== Ejecución inicial (se interrumpe) ===")
result = graph.invoke({"data": "", "log": []}, config)

state = graph.get_state(config)
print(f"Log: {state.values['log']}")
print(f"Siguiente nodo: {state.next}")

print("\n=== Resume ===")
result = graph.invoke(Command(resume="ok"), config)

state = graph.get_state(config)
print(f"Log: {state.values['log']}")
print(f"Siguiente nodo: {state.next}")
print(f"Resultado: {state.values['data']}")
# Output esperado:
# === Ejecución inicial (se interrumpe) ===
# Log: ['step_a completado']
# Siguiente nodo: ('b',)
#
# === Resume ===
# Log: ['step_a completado', 'step_b completado', 'step_c completado']
# Siguiente nodo: ()
# Resultado: resultado final

Ejercicio 3: Procesamiento por batches con loop (Medio)

Crea un grafo que procese una lista de 10 items en batches de 3. El grafo debe usar un loop (conditional edge) que repite el nodo de procesamiento hasta que todos los batches estén completos. Usa checkpointing y muestra el progreso en cada batch.

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.checkpoint.memory import MemorySaver

class State(TypedDict):
    items: list[str]
    processed: Annotated[list[str], operator.add]
    batch_num: int
    batch_size: int

def process_batch(state: State) -> dict:
    batch_num = state["batch_num"]
    batch_size = state["batch_size"]
    start = batch_num * batch_size
    end = min(start + batch_size, len(state["items"]))

    batch_items = state["items"][start:end]
    results = [f"✅ {item}" for item in batch_items]

    print(f"  Batch {batch_num}: procesados {len(results)} items ({start}-{end-1})")

    return {
        "processed": results,
        "batch_num": batch_num + 1,
    }

def should_continue(state: State) -> str:
    total_batches = -(-len(state["items"]) // state["batch_size"])
    if state["batch_num"] >= total_batches:
        return "done"
    return "next_batch"

def finalize(state: State) -> dict:
    return {}

graph_builder = StateGraph(State)
graph_builder.add_node("process", process_batch)
graph_builder.add_node("done", finalize)

graph_builder.add_edge(START, "process")
graph_builder.add_conditional_edges(
    "process", should_continue,
    {"next_batch": "process", "done": "done"},
)
graph_builder.add_edge("done", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

items = [f"doc_{i}" for i in range(10)]
config = {"configurable": {"thread_id": "batch-exercise"}}

result = graph.invoke(
    {"items": items, "processed": [], "batch_num": 0, "batch_size": 3},
    config,
)

print(f"\nTotal procesados: {len(result['processed'])}")
print(f"Batches ejecutados: {result['batch_num']}")
for item in result["processed"]:
    print(f"  {item}")
# Output esperado:
#   Batch 0: procesados 3 items (0-2)
#   Batch 1: procesados 3 items (3-5)
#   Batch 2: procesados 3 items (6-8)
#   Batch 3: procesados 1 items (9-9)
#
# Total procesados: 10
# Batches ejecutados: 4
#   ✅ doc_0
#   ✅ doc_1
#   ...
#   ✅ doc_9

Ejercicio 4: Nodo idempotente con verificación (Medio)

Crea un grafo con un nodo que simula enviar un email. El nodo debe ser idempotente: si email_sent ya está en True en el estado, no envía de nuevo. Ejecuta el grafo dos veces con el mismo thread y verifica que el email solo se "envía" una vez.

Ver solución
from dotenv import load_dotenv
load_dotenv()

from typing import TypedDict
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver

class State(TypedDict):
    recipient: str
    email_sent: bool
    send_count: int
    log: str

def prepare_email(state: State) -> dict:
    return {"log": f"Email preparado para {state['recipient']}"}

def send_email(state: State) -> dict:
    if state.get("email_sent", False):
        print(f"  ⏭️  Email a {state['recipient']} ya enviado. Saltando.")
        return {}

    print(f"  📧 Enviando email a {state['recipient']}...")
    return {
        "email_sent": True,
        "send_count": state.get("send_count", 0) + 1,
        "log": f"Email enviado a {state['recipient']}",
    }

def confirm(state: State) -> dict:
    return {}

graph_builder = StateGraph(State)
graph_builder.add_node("prepare", prepare_email)
graph_builder.add_node("send", send_email)
graph_builder.add_node("confirm", confirm)

graph_builder.add_edge(START, "prepare")
graph_builder.add_edge("prepare", "send")
graph_builder.add_edge("send", "confirm")
graph_builder.add_edge("confirm", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

config = {"configurable": {"thread_id": "email-idempotent"}}

print("=== Primera ejecución ===")
result1 = graph.invoke(
    {"recipient": "usuario@email.com", "email_sent": False, "send_count": 0, "log": ""},
    config,
)
print(f"Email enviado: {result1['email_sent']}")
print(f"Veces enviado: {result1['send_count']}")

config2 = {"configurable": {"thread_id": "email-idempotent-2"}}

print("\n=== Segunda ejecución (simula re-ejecución) ===")
result2 = graph.invoke(
    {
        "recipient": "usuario@email.com",
        "email_sent": result1["email_sent"],
        "send_count": result1["send_count"],
        "log": "",
    },
    config2,
)
print(f"Email enviado: {result2['email_sent']}")
print(f"Veces enviado: {result2['send_count']}")
# Output esperado:
# === Primera ejecución ===
#   📧 Enviando email a usuario@email.com...
# Email enviado: True
# Veces enviado: 1
#
# === Segunda ejecución (simula re-ejecución) ===
#   ⏭️  Email a usuario@email.com ya enviado. Saltando.
# Email enviado: True
# Veces enviado: 1

Ejercicio 5: Retry automático con configuración por nodo (Medio)

Crea un grafo con 2 nodos: uno estable (siempre funciona) y uno inestable (falla las primeras 2 veces). Configura retry automático en el nodo inestable con max_attempts=4. Verifica que el grafo completa exitosamente sin error handling manual.

Ver solución
from dotenv import load_dotenv
load_dotenv()

from typing import TypedDict
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver

UNSTABLE_CALLS = 0

class State(TypedDict):
    input: str
    stable_result: str
    unstable_result: str

def stable_node(state: State) -> dict:
    return {"stable_result": f"Análisis estable de '{state['input']}'"}

def unstable_node(state: State) -> dict:
    global UNSTABLE_CALLS
    UNSTABLE_CALLS += 1

    if UNSTABLE_CALLS <= 2:
        raise ConnectionError(f"Servicio inestable: fallo #{UNSTABLE_CALLS}")

    return {"unstable_result": f"Datos externos obtenidos (intento #{UNSTABLE_CALLS})"}

graph_builder = StateGraph(State)
graph_builder.add_node("stable", stable_node)
graph_builder.add_node(
    "unstable",
    unstable_node,
    retry={"max_attempts": 4, "delay": 0.2, "multiplier": 2.0},
)

graph_builder.add_edge(START, "stable")
graph_builder.add_edge("stable", "unstable")
graph_builder.add_edge("unstable", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

UNSTABLE_CALLS = 0
config = {"configurable": {"thread_id": "retry-exercise"}}
result = graph.invoke(
    {"input": "test query", "stable_result": "", "unstable_result": ""},
    config,
)
print(f"Estable: {result['stable_result']}")
print(f"Inestable: {result['unstable_result']}")
print(f"Intentos del nodo inestable: {UNSTABLE_CALLS}")
# Output esperado:
# Estable: Análisis estable de 'test query'
# Inestable: Datos externos obtenidos (intento #3)
# Intentos del nodo inestable: 3

Ejercicio 6: Pipeline durable completo con tracking de costos (Avanzado)

Construye un grafo de investigación con 5 nodos secuenciales. Cada nodo tiene un costo simulado diferente. Usa interrupt() para simular un crash después del nodo 3. Muestra el costo acumulado antes y después del crash. Resume la ejecución y verifica que el costo total es correcto (sin costos duplicados). Incluye un campo execution_log que registra cada paso con timestamp.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import time
import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import interrupt, Command

class State(TypedDict):
    topic: str
    cost_usd: float
    results: Annotated[list[str], operator.add]
    execution_log: Annotated[list[dict], operator.add]

def make_step(name: str, cost: float, crash: bool = False):
    def node(state: State) -> dict:
        if crash:
            interrupt(f"Crash simulado en {name}")

        log_entry = {
            "step": name,
            "cost": cost,
            "timestamp": time.time(),
            "cumulative_cost": state.get("cost_usd", 0) + cost,
        }

        return {
            "results": [f"{name}: completado (${cost:.2f})"],
            "cost_usd": state.get("cost_usd", 0) + cost,
            "execution_log": [log_entry],
        }
    return node

graph_builder = StateGraph(State)
graph_builder.add_node("search", make_step("search", 0.05))
graph_builder.add_node("fetch", make_step("fetch_papers", 0.15))
graph_builder.add_node("analyze", make_step("analyze", 0.25))
graph_builder.add_node("crash_point", make_step("deep_analysis", 0.30, crash=True))
graph_builder.add_node("report", make_step("generate_report", 0.20))

graph_builder.add_edge(START, "search")
graph_builder.add_edge("search", "fetch")
graph_builder.add_edge("fetch", "analyze")
graph_builder.add_edge("analyze", "crash_point")
graph_builder.add_edge("crash_point", "report")
graph_builder.add_edge("report", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

config = {"configurable": {"thread_id": "durable-pipeline"}}

print("=== Ejecución inicial (crash después de analyze) ===")
result = graph.invoke(
    {"topic": "AI safety", "cost_usd": 0.0, "results": [], "execution_log": []},
    config,
)

state = graph.get_state(config)
print(f"Pasos completados: {len(state.values['results'])}")
print(f"Costo acumulado: ${state.values['cost_usd']:.2f}")
print(f"Siguiente nodo: {state.next}")
for entry in state.values["execution_log"]:
    print(f"  ${entry['cost']:.2f}{entry['step']}")

print("\n=== Resume (continúa desde deep_analysis) ===")
result = graph.invoke(Command(resume="continuar"), config)

state = graph.get_state(config)
print(f"Pasos completados: {len(state.values['results'])}")
print(f"Costo total: ${state.values['cost_usd']:.2f}")
print(f"Siguiente nodo: {state.next}")
print(f"\nLog completo:")
for entry in state.values["execution_log"]:
    print(f"  ${entry['cost']:.2f}{entry['step']} (acumulado: ${entry['cumulative_cost']:.2f})")
# Output esperado:
# === Ejecución inicial (crash después de analyze) ===
# Pasos completados: 3
# Costo acumulado: $0.45
# Siguiente nodo: ('crash_point',)
#   $0.05 — search
#   $0.15 — fetch_papers
#   $0.25 — analyze
#
# === Resume (continúa desde deep_analysis) ===
# Pasos completados: 5
# Costo total: $0.95
# Siguiente nodo: ()
#
# Log completo:
#   $0.05 — search (acumulado: $0.05)
#   $0.15 — fetch_papers (acumulado: $0.20)
#   $0.25 — analyze (acumulado: $0.45)
#   $0.30 — deep_analysis (acumulado: $0.75)
#   $0.20 — generate_report (acumulado: $0.95)

Resumen

En esta cápsula aprendiste:

  • Durable execution = checkpointing + thread management — no es una feature separada, es el uso correcto de lo que ya sabías. Cada nodo completado guarda un checkpoint, y al reiniciar con el mismo thread_id, el grafo resume desde donde quedó
  • El valor es tangible y medible: una investigación de 5 pasos a $0.10 cada uno que crashea en el paso 3 cuesta $0.00 extra con durable execution vs $0.50 sin él. Multiplica por 100 usuarios diarios
  • Nodos idempotentes son obligatorios para operaciones con efectos secundarios — un nodo que cobra un pago o envía un email debe verificar si ya se ejecutó antes de actuar
  • Retry automático a nivel de nodo maneja errores transitorios sin código manual — retry={"max_attempts": 3, "delay": 1.0} en add_node es todo lo que necesitas
  • De MemorySaver a PostgresSaver es una línea — el cambio es trivial pero el valor es enorme: durabilidad real que sobrevive reinicios del proceso
  • El procesamiento por batches con loops es el patrón natural para agentes de larga duración — cada iteración del loop es un checkpoint, cada batch procesado es progreso guardado
  • Design patterns clave: progress tracking explícito, operaciones atómicas por nodo, checkpoint verification antes de operaciones costosas, y graceful shutdown para mantener consistencia

Próxima cápsula: Time-Travel Debugging — navegar el historial completo de estados de tu agente, retroceder a cualquier paso, y entender exactamente por qué tu agente tomó cada decisión.


Recursos adicionales

  1. LangGraph Persistence — Conceptos de persistencia y checkpointing en LangGraph
  2. LangGraph Checkpointers — MemorySaver, PostgresSaver y otros backends
  3. How to use LangGraph's built-in retry policy — Retry automático a nivel de nodo
  4. LangGraph Interrupt — Breakpoints e interrupciones para control de ejecución
  5. PostgresSaver Setup — Configuración de PostgreSQL como backend de checkpoints
  6. Idempotency Patterns — AWS: Making retries safe with idempotent APIs

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