Módulo 8: Memoria y Persistencia
Persistencia con PostgresSaver y Redis
Descripción de la cápsula
MemorySaver funciona perfecto en tu máquina. Abres el notebook, interactúas con el agente, los checkpoints se guardan, el time-travel funciona. Cierras el notebook. Lo abres al día siguiente. Todo desapareció. Cada checkpoint, cada conversación, cada historial de estados — evaporado.
Eso es aceptable en desarrollo. En producción es inaceptable. Tu usuario habló con el agente ayer, cerró la pestaña, y hoy espera continuar. Tu servidor hizo un restart a las 3am para aplicar un parche de seguridad, y los 200 threads activos se perdieron. Tu aplicación corre en 4 instancias de Kubernetes, y el usuario que estaba en la instancia 2 ahora llega a la instancia 3 — sin contexto.
PostgresSaver resuelve todo esto. Los checkpoints se guardan en una base de datos PostgreSQL. Sobreviven a crashes, restarts, deploys, y migraciones de instancia. Y la migración desde MemorySaver es trivial: cambias una línea de código. Una. Todo lo demás — thread_id, get_state(), get_state_history(), time-travel — funciona exactamente igual.
La migración: una línea de código
Esta es la transición completa de desarrollo a producción:
Antes (desarrollo con MemorySaver)
from langgraph.checkpoint.memory import MemorySaver
checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)
Después (producción con PostgresSaver)
from langgraph.checkpoint.postgres import PostgresSaver
checkpointer = PostgresSaver.from_conn_string("postgresql://user:pass@localhost:5432/mydb")
checkpointer.setup()
graph = graph_builder.compile(checkpointer=checkpointer)
Eso es todo. El resto de tu código — los nodos, los edges, los invokes con thread_id, las llamadas a get_state(), el time-travel — no cambia ni una línea. Esa es la potencia de la abstracción de checkpointer en LangGraph: el grafo no sabe ni le importa dónde se guardan los checkpoints. Solo sabe que se guardan.
PostgresSaver: setup completo
Instalación
pip install langgraph-checkpoint-postgres psycopg[binary]
langgraph-checkpoint-postgres es el paquete que contiene PostgresSaver. psycopg[binary] es el driver de PostgreSQL para Python — la versión binary incluye las librerías C precompiladas para no necesitar compilación.
Connection string
El formato es estándar PostgreSQL:
postgresql://usuario:contraseña@host:puerto/base_de_datos
| Componente | Ejemplo | Descripción |
|---|---|---|
usuario | postgres | Usuario de la base de datos |
contraseña | mysecretpass | Contraseña del usuario |
host | localhost | Dirección del servidor |
puerto | 5432 | Puerto de PostgreSQL (default: 5432) |
base_de_datos | langgraph_app | Nombre de la base de datos |
En producción, nunca pongas la contraseña en el código. Usa variables de entorno:
import os
conn_string = os.environ["DATABASE_URL"]
# Ejemplo: DATABASE_URL=postgresql://user:pass@db.example.com:5432/langgraph_prod
Crear las tablas necesarias
PostgresSaver necesita tablas para almacenar los checkpoints. El método setup() las crea automáticamente:
from langgraph.checkpoint.postgres import PostgresSaver
conn_string = "postgresql://postgres:postgres@localhost:5432/langgraph_dev"
checkpointer = PostgresSaver.from_conn_string(conn_string)
checkpointer.setup()
setup() es idempotente — puedes ejecutarlo múltiples veces sin problema. Si las tablas ya existen, no hace nada. Es seguro llamarlo en el startup de tu aplicación.
Versión asíncrona
Para aplicaciones async (FastAPI, por ejemplo), usa la versión AsyncPostgresSaver:
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver
conn_string = "postgresql://postgres:postgres@localhost:5432/langgraph_dev"
async def create_graph():
checkpointer = AsyncPostgresSaver.from_conn_string(conn_string)
await checkpointer.setup()
graph = graph_builder.compile(checkpointer=checkpointer)
return graph
Ejemplo completo: chatbot con PostgresSaver
import os
from langchain.chat_models import init_chat_model
from langgraph.checkpoint.postgres import PostgresSaver
from langgraph.graph import StateGraph, MessagesState, START, END
def chatbot(state: MessagesState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
response = model.invoke(state["messages"])
return {"messages": [response]}
graph_builder = StateGraph(MessagesState)
graph_builder.add_node("chatbot", chatbot)
graph_builder.add_edge(START, "chatbot")
graph_builder.add_edge("chatbot", END)
conn_string = os.environ.get(
"DATABASE_URL",
"postgresql://postgres:postgres@localhost:5432/langgraph_dev"
)
checkpointer = PostgresSaver.from_conn_string(conn_string)
checkpointer.setup()
graph = graph_builder.compile(checkpointer=checkpointer)
config = {"configurable": {"thread_id": "user_001"}}
result = graph.invoke({"messages": [("user", "¿Qué es RAG?")]}, config)
print(result["messages"][-1].content)
# Output: "RAG (Retrieval-Augmented Generation) es una técnica que..."
result = graph.invoke({"messages": [("user", "Dame ejemplos")]}, config)
print(result["messages"][-1].content)
# Output: "Aquí tienes ejemplos de RAG: 1) Un chatbot que consulta
# documentación interna..."
El código es casi idéntico al de MemorySaver. Las únicas diferencias:
from langgraph.checkpoint.postgres import PostgresSaveren vez defrom langgraph.checkpoint.memory import MemorySaverPostgresSaver.from_conn_string(conn_string)en vez deMemorySaver()checkpointer.setup()para crear las tablas
Todo lo demás — thread_id, invoke, get_state(), get_state_history() — es idéntico. Si mañana necesitas cambiar de Postgres a otro backend, cambias esas 2-3 líneas.
Durabilidad: la diferencia real
La diferencia entre MemorySaver y PostgresSaver no es el API — es lo que pasa cuando algo sale mal.
| Escenario | MemorySaver | PostgresSaver |
|---|---|---|
| Deploy a las 3am | 200 threads activos: PERDIDOS. Usuarios vuelven y el agente no recuerda nada | 200 threads en PostgreSQL: INTACTOS. Usuarios continúan donde quedaron |
| Crash durante investigación | 3 de 5 fuentes procesadas, crash. Empezar de cero: 5 min + doble costo API | Resume desde fuente 4: 1 min, sin costo duplicado |
| Multi-instancia (k8s) | Load balancer envía request a otra instancia → sin contexto | Todas las instancias comparten la misma DB → continuidad perfecta |
Cuándo usar cada checkpointer
| Checkpointer | Caso de uso | Durabilidad | Performance | Escalabilidad |
|---|---|---|---|---|
| MemorySaver | Desarrollo, testing, notebooks | ❌ Proceso restart = perdido | Más rápido (RAM) | Solo un proceso |
| PostgresSaver | Producción, multi-instancia | ✅ Sobrevive restarts | Bueno (red + disco) | Múltiples instancias |
| RedisSaver | Alta frecuencia, caché de sesiones | ⚠️ Configurable (TTL) | Muy rápido (RAM + red) | Múltiples instancias |
Árbol de decisión
¿Estás en desarrollo/testing?
└── Sí → MemorySaver (sin setup, sin dependencias)
└── No → ¿Necesitas durabilidad total?
└── Sí → PostgresSaver (producción estándar)
└── No → ¿Necesitas máxima velocidad?
└── Sí → RedisSaver (sesiones efímeras, alta frecuencia)
└── No → PostgresSaver (default de producción)
La recomendación para el 90% de los casos: MemorySaver para desarrollo, PostgresSaver para producción. RedisSaver es para escenarios específicos donde la latencia de lectura/escritura importa más que la durabilidad absoluta.
RedisSaver: cuando necesitas velocidad
RedisSaver guarda checkpoints en Redis en lugar de PostgreSQL. Redis opera en memoria (como MemorySaver) pero es un servicio separado que persiste datos a disco y sobrevive a restarts del proceso Python.
Instalación
pip install langgraph-checkpoint-redis
Uso básico
from langgraph.checkpoint.redis import RedisSaver
checkpointer = RedisSaver.from_conn_string("redis://localhost:6379")
graph = graph_builder.compile(checkpointer=checkpointer)
PostgresSaver vs RedisSaver
| Dimensión | PostgresSaver | RedisSaver |
|---|---|---|
| Latencia de lectura | ~1-5ms (disco/SSD) | ~0.1-1ms (memoria) |
| Latencia de escritura | ~2-10ms | ~0.1-1ms |
| Durabilidad | Total (WAL + fsync) | Configurable (RDB/AOF) |
| Capacidad | Terabytes (disco) | Limitada por RAM |
| TTL automático | Manual (necesitas cron/trigger) | Nativo (EXPIRE) |
| Costo | Menor (disco es barato) | Mayor (RAM es cara) |
| Queries complejas | SQL completo | Key-value solamente |
RedisSaver brilla cuando:
- Tienes miles de requests por segundo y cada ms importa
- Los checkpoints son efímeros (sesiones de chat que expiran en 24h)
- Ya tienes Redis en tu infraestructura
PostgresSaver gana cuando:
- Necesitas durabilidad garantizada
- Quieres hacer queries sobre los checkpoints (analytics, debugging)
- Los checkpoints deben persistir indefinidamente
- El volumen de datos crece más allá de lo que cabe en RAM
Connection pooling: manejo eficiente de conexiones
En producción, cada request que necesita el checkpointer abre y cierra una conexión a la base de datos. Con 100 requests concurrentes, eso son 100 conexiones abiertas simultáneamente. PostgreSQL tiene un límite (default: 100), y abrirlas/cerrarlas constantemente es costoso.
Connection pooling resuelve esto: mantiene un pool de conexiones reutilizables.
Con psycopg pool
import os
from psycopg_pool import ConnectionPool
from langgraph.checkpoint.postgres import PostgresSaver
conn_string = os.environ.get(
"DATABASE_URL",
"postgresql://postgres:postgres@localhost:5432/langgraph_dev"
)
pool = ConnectionPool(
conninfo=conn_string,
min_size=5,
max_size=20,
)
checkpointer = PostgresSaver(conn=pool)
checkpointer.setup()
graph = graph_builder.compile(checkpointer=checkpointer)
| Parámetro | Valor sugerido | Descripción |
|---|---|---|
min_size | 5 | Conexiones mínimas abiertas permanentemente |
max_size | 20 | Conexiones máximas (limita picos) |
Regla de dedo: max_size ≤ max_connections de PostgreSQL / número de instancias de tu app. Si PostgreSQL tiene max_connections=100 y corres 4 instancias, max_size=20 por instancia (4 × 20 = 80, deja margen).
Para async (FastAPI), usa AsyncConnectionPool de psycopg_pool con AsyncPostgresSaver — mismos parámetros, misma lógica, pero con await pool.open().
Retención y limpieza de checkpoints
Los checkpoints se acumulan. Sin política de retención, la base de datos crece indefinidamente. Tres estrategias: por tiempo (borrar > N días), por cantidad (últimos N por thread), o por actividad (threads inactivos).
En PostgreSQL, limpieza directa con SQL:
DELETE FROM checkpoints WHERE created_at < NOW() - INTERVAL '30 days';
En Redis, TTL nativo — los checkpoints expiran automáticamente:
checkpointer = RedisSaver.from_conn_string("redis://localhost:6379", ttl={"default": 86400})
Reglas de almacenamiento: no guardes datos grandes en el estado (usa referencias, no contenido), aplica message trimming (cápsula 02), y monitorea el tamaño con queries SQL periódicas.
Patrón de producción: selección de checkpointer por entorno
import os
def get_checkpointer():
env = os.environ.get("ENVIRONMENT", "development")
if env == "development":
from langgraph.checkpoint.memory import MemorySaver
return MemorySaver()
elif env in ("production", "staging"):
from langgraph.checkpoint.postgres import PostgresSaver
checkpointer = PostgresSaver.from_conn_string(os.environ["DATABASE_URL"])
checkpointer.setup()
return checkpointer
raise ValueError(f"Unknown environment: {env}")
graph = graph_builder.compile(checkpointer=get_checkpointer())
El grafo no cambia. Los nodos no cambian. Los tests no cambian. Solo el checkpointer varía según el entorno.
Ejemplo completo: FastAPI con PostgresSaver async
Un patrón de producción con FastAPI: el lifespan crea el pool y el checkpointer al startup, los cierra al shutdown. El pool se comparte entre todos los requests. Cada request pasa su thread_id y usa ainvoke (async). El servidor puede reiniciarse y las conversaciones persisten.
import os
from contextlib import asynccontextmanager
from fastapi import FastAPI
from pydantic import BaseModel
from langchain.chat_models import init_chat_model
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver
from langgraph.graph import StateGraph, MessagesState, START, END
from psycopg_pool import AsyncConnectionPool
graph = None
def chatbot(state: MessagesState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
response = model.invoke(state["messages"])
return {"messages": [response]}
graph_builder = StateGraph(MessagesState)
graph_builder.add_node("chatbot", chatbot)
graph_builder.add_edge(START, "chatbot")
graph_builder.add_edge("chatbot", END)
@asynccontextmanager
async def lifespan(app: FastAPI):
global graph
conn_string = os.environ["DATABASE_URL"]
pool = AsyncConnectionPool(conninfo=conn_string, min_size=5, max_size=20)
await pool.open()
checkpointer = AsyncPostgresSaver(conn=pool)
await checkpointer.setup()
graph = graph_builder.compile(checkpointer=checkpointer)
yield
await pool.close()
app = FastAPI(lifespan=lifespan)
class ChatRequest(BaseModel):
message: str
thread_id: str
@app.post("/chat")
async def chat(request: ChatRequest):
config = {"configurable": {"thread_id": request.thread_id}}
result = await graph.ainvoke(
{"messages": [("user", request.message)]}, config
)
return {"response": result["messages"][-1].content, "thread_id": request.thread_id}
# POST /chat {"message": "¿Qué es RAG?", "thread_id": "user_001"}
# → {"response": "RAG es...", "thread_id": "user_001"}
Troubleshooting
Problema 1: "relation 'checkpoints' does not exist"
Síntoma: Error al hacer el primer invoke después de configurar PostgresSaver.
Causa: No llamaste checkpointer.setup() para crear las tablas.
Solución: Agrega checkpointer.setup() (o await checkpointer.setup() en async) después de crear el checkpointer. Es idempotente — seguro de llamar en cada startup.
Problema 2: "connection refused" al conectar a PostgreSQL
Síntoma: psycopg.OperationalError: connection to server at "localhost"... refused.
Causa: PostgreSQL no está corriendo, o el puerto/host es incorrecto.
Solución: Verifica que PostgreSQL está activo (pg_isready -h localhost -p 5432). Si usas Docker: docker run -d -p 5432:5432 -e POSTGRES_PASSWORD=postgres postgres:16.
Problema 3: "too many clients already" en producción
Síntoma: psycopg.OperationalError: too many clients already bajo carga.
Causa: Cada request abre una nueva conexión y el pool no está configurado, o max_size del pool excede max_connections de PostgreSQL.
Solución: Usa ConnectionPool con max_size apropiado (ver sección de connection pooling). Regla: max_size × num_instancias < max_connections.
Problema 4: Performance degradado con muchos checkpoints
Síntoma: get_state_history() es lento después de miles de invocaciones en un thread.
Causa: La tabla de checkpoints creció sin indexación adicional o sin política de retención.
Solución: Implementa una política de retención (borrar checkpoints viejos) y verifica que los índices de la tabla están activos. Ver sección de gestión de checkpoints.
Problema 5: MemorySaver en producción "funciona" pero pierde datos
Síntoma: Todo funciona bien hasta que el servidor se reinicia y las conversaciones desaparecen. Causa: MemorySaver está en uso en producción. Solución: Migra a PostgresSaver. Es literalmente cambiar una línea (ver sección "La migración").
Ejercicios
Ejercicio 1: Migración MemorySaver → PostgresSaver (Fácil)
Tienes este código con MemorySaver. Modifícalo para usar PostgresSaver con la connection string "postgresql://postgres:postgres@localhost:5432/langgraph_dev". No cambies nada más que las líneas del checkpointer.
from langchain.chat_models import init_chat_model
from langgraph.checkpoint.memory import MemorySaver
from langgraph.graph import StateGraph, MessagesState, START, END
def chatbot(state: MessagesState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
response = model.invoke(state["messages"])
return {"messages": [response]}
graph_builder = StateGraph(MessagesState)
graph_builder.add_node("chatbot", chatbot)
graph_builder.add_edge(START, "chatbot")
graph_builder.add_edge("chatbot", END)
checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)
Ver solución
from langchain.chat_models import init_chat_model
from langgraph.checkpoint.postgres import PostgresSaver
from langgraph.graph import StateGraph, MessagesState, START, END
def chatbot(state: MessagesState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
response = model.invoke(state["messages"])
return {"messages": [response]}
graph_builder = StateGraph(MessagesState)
graph_builder.add_node("chatbot", chatbot)
graph_builder.add_edge(START, "chatbot")
graph_builder.add_edge("chatbot", END)
conn_string = "postgresql://postgres:postgres@localhost:5432/langgraph_dev"
checkpointer = PostgresSaver.from_conn_string(conn_string)
checkpointer.setup()
graph = graph_builder.compile(checkpointer=checkpointer)
config = {"configurable": {"thread_id": "migration_test"}}
result = graph.invoke({"messages": [("user", "¿Funciona la migración?")]}, config)
print(result["messages"][-1].content)
# Output esperado: Sí, la migración funciona. El agente responde normalmente.
state = graph.get_state(config)
print(f"Checkpoint guardado en Postgres: {len(state.values['messages'])} mensajes")
print("✅ Migración exitosa — los checkpoints ahora persisten en PostgreSQL")
# Output esperado:
# Checkpoint guardado en Postgres: 2 mensajes
# ✅ Migración exitosa — los checkpoints ahora persisten en PostgreSQL
Explicación: Tres cambios: (1) cambiar el import de MemorySaver a PostgresSaver, (2) usar PostgresSaver.from_conn_string() en vez de MemorySaver(), (3) llamar checkpointer.setup(). Todo lo demás — nodos, edges, invoke, config — idéntico.
Ejercicio 2: Selección de checkpointer por entorno (Fácil)
Escribe una función get_checkpointer(env: str) que retorne MemorySaver para "development", PostgresSaver para "production" (con connection string desde variable de entorno), y lance ValueError para cualquier otro valor. Testea con "development" y verifica que retorna un MemorySaver.
Ver solución
import os
from langgraph.checkpoint.memory import MemorySaver
def get_checkpointer(env: str):
if env == "development":
return MemorySaver()
elif env == "production":
from langgraph.checkpoint.postgres import PostgresSaver
conn_string = os.environ.get("DATABASE_URL")
if not conn_string:
raise ValueError("DATABASE_URL environment variable required for production")
checkpointer = PostgresSaver.from_conn_string(conn_string)
checkpointer.setup()
return checkpointer
raise ValueError(f"Unknown environment: {env}")
dev_checkpointer = get_checkpointer("development")
print(f"Dev checkpointer: {type(dev_checkpointer).__name__}")
assert type(dev_checkpointer).__name__ == "MemorySaver"
try:
get_checkpointer("unknown")
except ValueError as e:
print(f"Error esperado: {e}")
print("✅ Selector de checkpointer funciona correctamente")
# Output esperado:
# Dev checkpointer: MemorySaver
# Error esperado: Unknown environment: unknown
# ✅ Selector de checkpointer funciona correctamente
Explicación: El import de PostgresSaver es lazy (dentro del elif) para que no falle en desarrollo si langgraph-checkpoint-postgres no está instalado. La validación de DATABASE_URL previene errores crípticos en producción.
Ejercicio 3: Verificar persistencia post-restart simulado (Medio)
Crea un grafo con MemorySaver. Haz 3 invocaciones en un thread. Luego simula un "restart" creando un nuevo MemorySaver() y compilando un nuevo grafo. Verifica que el nuevo grafo NO tiene los checkpoints anteriores. Después, repite el experimento pero reutilizando el mismo objeto checkpointer — verifica que SÍ mantiene los checkpoints. Esto demuestra por qué MemorySaver no sirve para producción.
Ver solución
from langchain.chat_models import init_chat_model
from langgraph.checkpoint.memory import MemorySaver
from langgraph.graph import StateGraph, MessagesState, START, END
def build_graph(checkpointer):
def chatbot(state: MessagesState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
response = model.invoke(state["messages"])
return {"messages": [response]}
builder = StateGraph(MessagesState)
builder.add_node("chatbot", chatbot)
builder.add_edge(START, "chatbot")
builder.add_edge("chatbot", END)
return builder.compile(checkpointer=checkpointer)
config = {"configurable": {"thread_id": "persist_test"}}
print("=== Escenario 1: Nuevo MemorySaver (simula restart) ===")
checkpointer_v1 = MemorySaver()
graph_v1 = build_graph(checkpointer_v1)
graph_v1.invoke({"messages": [("user", "Hola, soy el turno 1")]}, config)
graph_v1.invoke({"messages": [("user", "Este es el turno 2")]}, config)
graph_v1.invoke({"messages": [("user", "Y este el turno 3")]}, config)
state_v1 = graph_v1.get_state(config)
print(f"Antes del restart: {len(state_v1.values['messages'])} mensajes")
checkpointer_v2 = MemorySaver()
graph_v2 = build_graph(checkpointer_v2)
state_v2 = graph_v2.get_state(config)
has_state = state_v2.values is not None and len(state_v2.values.get("messages", [])) > 0
print(f"Después del restart: {'tiene datos' if has_state else 'vacío'}")
assert not has_state, "No debería tener datos después de un restart simulado"
print("❌ Checkpoints perdidos — MemorySaver no persiste entre restarts\n")
print("=== Escenario 2: Mismo MemorySaver (sin restart real) ===")
shared_checkpointer = MemorySaver()
graph_a = build_graph(shared_checkpointer)
graph_a.invoke({"messages": [("user", "Mensaje en grafo A")]}, config)
graph_b = build_graph(shared_checkpointer)
state_b = graph_b.get_state(config)
has_state_b = state_b.values is not None and len(state_b.values.get("messages", [])) > 0
print(f"Grafo B con mismo checkpointer: {'tiene datos' if has_state_b else 'vacío'}")
assert has_state_b
print("✅ Checkpoints preservados — mismo objeto en memoria")
print("Conclusión: MemorySaver pierde datos al crear nueva instancia. PostgresSaver no.")
# Output esperado:
# Antes del restart: 6 mensajes → Después: vacío
# ❌ MemorySaver no persiste → ✅ Mismo objeto sí persiste
Explicación: Crear un nuevo MemorySaver() equivale a un restart del proceso — la RAM se limpia. Reutilizar el mismo objeto mantiene los datos, pero en producción no puedes garantizar eso. PostgresSaver desacopla los datos del proceso.
Ejercicio 4: Connection pool con validación (Medio)
Escribe create_pooled_checkpointer(conn_string, min_size, max_size) con validación: min_size >= 1, max_size >= min_size, max_size <= 50. Testea que las validaciones rechazan configuraciones peligrosas.
Ver solución
from psycopg_pool import ConnectionPool
from langgraph.checkpoint.postgres import PostgresSaver
def create_pooled_checkpointer(conn_string: str, min_size: int = 5, max_size: int = 20):
if min_size < 1:
raise ValueError(f"min_size debe ser >= 1, recibido: {min_size}")
if max_size < min_size:
raise ValueError(f"max_size ({max_size}) debe ser >= min_size ({min_size})")
if max_size > 50:
raise ValueError(f"max_size debe ser <= 50, recibido: {max_size}")
pool = ConnectionPool(conninfo=conn_string, min_size=min_size, max_size=max_size)
checkpointer = PostgresSaver(conn=pool)
checkpointer.setup()
return checkpointer
invalid_cases = [(0, 10), (10, 5), (5, 100)]
for min_s, max_s in invalid_cases:
try:
create_pooled_checkpointer("postgresql://localhost/test", min_s, max_s)
print(f" ❌ min={min_s}, max={max_s}: should have failed")
except ValueError as e:
print(f" ✅ min={min_s}, max={max_s}: rejected ({e})")
print("✅ Validaciones del pool funcionan correctamente")
# Output esperado:
# ✅ min=0, max=10: rejected (min_size debe ser >= 1...)
# ✅ min=10, max=5: rejected (max_size (5) debe ser >= min_size...)
# ✅ min=5, max=100: rejected (max_size debe ser <= 50...)
# ✅ Validaciones del pool funcionan correctamente
Explicación: max_size=50 como tope evita agotar las conexiones de PostgreSQL. En producción, ajusta según max_connections de tu instancia.
Ejercicio 5: Función de limpieza de checkpoints (Medio-Avanzado)
Escribe una función cleanup_old_threads(graph, known_thread_ids, max_age_turns, active_thread_ids) que retorne los thread_ids que deberían limpiarse (inactivos con más de max_age_turns turnos). Simula con MemorySaver.
Ver solución
from langchain.chat_models import init_chat_model
from langgraph.checkpoint.memory import MemorySaver
from langgraph.graph import StateGraph, MessagesState, START, END
def chatbot(state: MessagesState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
return {"messages": [model.invoke(state["messages"])]}
graph_builder = StateGraph(MessagesState)
graph_builder.add_node("chatbot", chatbot)
graph_builder.add_edge(START, "chatbot")
graph_builder.add_edge("chatbot", END)
checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)
def populate_thread(graph, thread_id: str, num_turns: int):
config = {"configurable": {"thread_id": thread_id}}
for i in range(num_turns):
graph.invoke({"messages": [("user", f"Turno {i+1}")]}, config)
def cleanup_old_threads(graph, known_ids, max_age_turns, active_ids):
to_clean = []
for tid in known_ids:
if tid in active_ids:
continue
state = graph.get_state({"configurable": {"thread_id": tid}})
if state.values and len(state.values.get("messages", [])) // 2 > max_age_turns:
to_clean.append(tid)
return to_clean
populate_thread(graph, "alice", 3)
populate_thread(graph, "bob", 15)
populate_thread(graph, "carol", 8)
populate_thread(graph, "dave", 2)
populate_thread(graph, "eve", 20)
to_clean = cleanup_old_threads(
graph, ["alice", "bob", "carol", "dave", "eve"],
max_age_turns=10, active_ids={"alice", "carol"}
)
print(f"Threads a limpiar: {to_clean}")
assert "bob" in to_clean and "eve" in to_clean
assert "alice" not in to_clean and "dave" not in to_clean
print("✅ Lógica de limpieza verificada")
# Output esperado:
# Threads a limpiar: ['bob', 'eve']
# ✅ Lógica de limpieza verificada
Explicación: En producción con PostgresSaver, esto sería SQL (DELETE FROM checkpoints WHERE thread_id = ...). Los threads activos siempre se protegen, los inactivos con muchos turnos se marcan para limpieza.
Ejercicio 6: Endpoint FastAPI con persistencia (Avanzado)
Escribe un endpoint POST /chat en FastAPI que reciba message y thread_id, use un grafo con MemorySaver, y retorne la respuesta. El checkpointer debe inicializarse en el lifespan. Incluye comentarios indicando qué cambiar para migrar a AsyncPostgresSaver.
Ver solución
from contextlib import asynccontextmanager
from fastapi import FastAPI
from pydantic import BaseModel
from langchain.chat_models import init_chat_model
from langgraph.checkpoint.memory import MemorySaver
from langgraph.graph import StateGraph, MessagesState, START, END
graph = None
def chatbot(state: MessagesState) -> dict:
model = init_chat_model("openai:gpt-4.1-mini")
response = model.invoke(state["messages"])
return {"messages": [response]}
graph_builder = StateGraph(MessagesState)
graph_builder.add_node("chatbot", chatbot)
graph_builder.add_edge(START, "chatbot")
graph_builder.add_edge("chatbot", END)
@asynccontextmanager
async def lifespan(app: FastAPI):
global graph
# Para migrar a PostgresSaver, cambia estas líneas:
# from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver
# checkpointer = AsyncPostgresSaver.from_conn_string(os.environ["DATABASE_URL"])
# await checkpointer.setup()
checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)
yield
app = FastAPI(lifespan=lifespan)
class ChatRequest(BaseModel):
message: str
thread_id: str
@app.post("/chat")
async def chat(request: ChatRequest):
config = {"configurable": {"thread_id": request.thread_id}}
result = await graph.ainvoke(
{"messages": [("user", request.message)]}, config
)
return {
"response": result["messages"][-1].content,
"thread_id": request.thread_id,
"total_messages": len(result["messages"])
}
# uvicorn module_name:app --reload
# curl -X POST http://localhost:8000/chat -H "Content-Type: application/json" \
# -d '{"message": "¿Qué es RAG?", "thread_id": "test_001"}'
# → {"response": "RAG es...", "thread_id": "test_001", "total_messages": 2}
Explicación: El checkpointer se crea en lifespan y se comparte entre todos los requests. Para migrar a Postgres, cambias 2-3 líneas en lifespan. Los endpoints no cambian. ainvoke (en vez de invoke) es para compatibilidad async con FastAPI.
Resumen
En esta cápsula aprendiste:
- La migración de MemorySaver a PostgresSaver es una línea de código. Cambiar
MemorySaver()porPostgresSaver.from_conn_string(conn_string)— todo lo demás (thread_id, get_state, get_state_history, time-travel) funciona idéntico. Esa es la potencia de la abstracción de checkpointer - PostgresSaver = durabilidad real. Los checkpoints sobreviven a crashes, restarts, deploys, y migraciones de instancia. En multi-instancia (Kubernetes), todos los pods comparten la misma base de datos
- RedisSaver es para alta frecuencia. Latencia sub-milisegundo, TTL nativo para expirar checkpoints automáticamente. Ideal para sesiones efímeras. Pero la durabilidad no es absoluta como Postgres
- Connection pooling es obligatorio en producción. Sin pool, cada request abre una conexión nueva. Con 100 requests concurrentes, agotas las conexiones de PostgreSQL.
ConnectionPool(min_size=5, max_size=20)resuelve esto - Los checkpoints se acumulan y necesitan retención. Política de limpieza por tiempo, por cantidad, o por actividad. En Redis, TTL nativo. En PostgreSQL, queries SQL periódicas
- El patrón de producción: MemorySaver para desarrollo (sin dependencias), PostgresSaver + connection pooling para producción (durabilidad + escala), selección por variable de entorno
Recursos adicionales
- LangGraph — PostgresSaver — API reference de PostgresSaver con todos los métodos y opciones de configuración
- LangGraph — Persistence How-to — Guía oficial para implementar persistencia con diferentes backends
- psycopg3 — Connection Pools — Documentación del connection pool de psycopg3, el driver usado por PostgresSaver
- Redis Persistence — redis.io — Cómo Redis persiste datos a disco: RDB snapshots vs AOF log. Importante para entender la durabilidad de RedisSaver
- PostgreSQL Connection Management — Documentación oficial de PostgreSQL sobre
max_connectionsy gestión de conexiones
Módulo 8 — LangChain & LangGraph: From Chains to Agents