Módulo 7: Flujos Avanzados
Patterns de Producción
Descripción de la cápsula
Los patterns de las cápsulas 02-06 hacen que tu sistema sea robusto: retry con backoff, branching paralelo, subgraphs reutilizables, map-reduce para colecciones, y error handling con fallback. Pero hay una capa más que separa un prototipo funcional de un sistema de producción real: timeouts por nodo, rate limiting interno, logging estructurado, y métricas de rendimiento.
Son los "últimos 10%" que hacen que tu sistema sea operable — que puedas entender qué pasa cuando algo falla a las 3am, que puedas responder "¿por qué esta request tardó 12 segundos?" con datos concretos, y que puedas escalar de 1 usuario en desarrollo a 100 usuarios concurrentes sin que todo se caiga.
El problema
Tu agente funciona perfecto en desarrollo. Lo pruebas con un tema, ejecuta en 3 segundos, genera un reporte limpio. Listo para producción, ¿verdad?
No.
En producción con 100 usuarios concurrentes:
- Un nodo de búsqueda tarda 45 segundos porque la API externa está lenta. Los otros 99 usuarios esperan.
- Tu agente hace 200 llamadas por minuto a una API que permite 60. Te bloquean el acceso.
- Un usuario reporta "mi investigación falló". No tienes idea de dónde ni por qué — solo sabes que falló.
- El equipo pregunta "¿cuántos tokens gastamos ayer?" y no puedes responder.
Estos son problemas de operabilidad, no de funcionalidad. Tu agente hace lo correcto — pero no puedes operar lo, monitorear lo, ni escalar lo. Esta cápsula cierra esa brecha.
Timeouts por nodo
El problema que resuelve
Un nodo de tu grafo llama a una API externa. Normalmente responde en 2 segundos. Pero un día la API está lenta y tarda 60 segundos. Sin timeout, ese nodo bloquea todo el pipeline. Con 100 requests concurrentes, tienes 100 threads bloqueados esperando una API que no va a mejorar.
Implementación con asyncio
from dotenv import load_dotenv
load_dotenv()
import asyncio
import time
from langgraph.func import entrypoint, task
@task
async def fast_search(query: str) -> dict:
"""Búsqueda que responde rápido."""
await asyncio.sleep(0.5)
return {"source": "fast_api", "content": f"Resultados rápidos para '{query}'"}
@task
async def slow_search(query: str) -> dict:
"""Búsqueda que simula una API lenta."""
await asyncio.sleep(10)
return {"source": "slow_api", "content": f"Resultados lentos para '{query}'"}
@task
async def search_with_timeout(query: str, timeout_seconds: float = 3.0) -> dict:
"""Ejecuta búsqueda con timeout. Si excede el tiempo, retorna error."""
try:
result = await asyncio.wait_for(
slow_search.acall(query),
timeout=timeout_seconds,
)
return {"status": "success", "result": result}
except asyncio.TimeoutError:
return {
"status": "timeout",
"error": f"Búsqueda excedió {timeout_seconds}s",
"query": query,
}
@entrypoint()
async def research_with_timeouts(topic: str) -> dict:
start = time.time()
fast_future = fast_search(topic)
slow_future = search_with_timeout(topic, timeout_seconds=3.0)
fast_result = await fast_future
slow_result = await slow_future
elapsed = time.time() - start
return {
"topic": topic,
"fast_result": fast_result,
"slow_result": slow_result,
"total_time_seconds": round(elapsed, 2),
}
result = asyncio.run(
research_with_timeouts.ainvoke("machine learning en medicina")
)
print(f"Tiempo total: {result['total_time_seconds']}s")
print(f"Fast: {result['fast_result']['source']}")
print(f"Slow: {result['slow_result']['status']}")
# Output esperado:
# Tiempo total: ~3.0s (no 10s)
# Fast: fast_api
# Slow: timeout
Sin el timeout, el pipeline tardaría 10 segundos esperando a slow_search. Con timeout de 3 segundos, falla rápido y el pipeline continúa con los resultados disponibles.
Timeout con threading (para código síncrono)
Si tus tasks son síncronas, puedes usar concurrent.futures para timeouts:
from dotenv import load_dotenv
load_dotenv()
import time
from concurrent.futures import ThreadPoolExecutor, TimeoutError as FuturesTimeoutError
from langgraph.func import entrypoint, task
def _slow_api_call(query: str) -> str:
"""Simula una API que tarda mucho."""
time.sleep(15)
return f"Resultado para {query}"
@task
def search_with_thread_timeout(query: str, timeout_seconds: float = 3.0) -> dict:
"""Ejecuta una búsqueda síncrona con timeout usando ThreadPoolExecutor."""
with ThreadPoolExecutor(max_workers=1) as executor:
future = executor.submit(_slow_api_call, query)
try:
result = future.result(timeout=timeout_seconds)
return {"status": "success", "content": result}
except FuturesTimeoutError:
return {
"status": "timeout",
"error": f"API excedió {timeout_seconds}s",
}
@entrypoint()
def pipeline_with_sync_timeout(topic: str) -> dict:
result = search_with_thread_timeout(topic, timeout_seconds=2.0).result()
return result
result = pipeline_with_sync_timeout.invoke("quantum computing")
print(f"Status: {result['status']}")
print(f"Error: {result.get('error', 'ninguno')}")
# Output esperado:
# Status: timeout
# Error: API excedió 2.0s
Cuándo usar cada approach
| Escenario | Approach | Por qué |
|---|---|---|
Tasks async (APIs con aiohttp) | asyncio.wait_for | Nativo, eficiente, no crea threads extra |
| Tasks síncronas (requests, SDK) | ThreadPoolExecutor | Funciona con cualquier código blocking |
| Timeout global del pipeline | Timeout en el .invoke() caller | No modifica el grafo |
Rate limiting interno
El problema que resuelve
Tu agente busca en 3 fuentes para cada sub-query, con 4 sub-queries. Son 12 llamadas a APIs externas en paralelo. Si la API permite 10 requests por minuto, las últimas 2 llamadas fallan con 429 Too Many Requests. Peor: si tienes 10 usuarios concurrentes, son 120 llamadas simultáneas.
Implementación con InMemoryRateLimiter
LangChain incluye InMemoryRateLimiter para controlar la tasa de llamadas:
from dotenv import load_dotenv
load_dotenv()
import time
from langchain_core.rate_limiters import InMemoryRateLimiter
from langchain.chat_models import init_chat_model
from langgraph.func import entrypoint, task
rate_limiter = InMemoryRateLimiter(
requests_per_second=2,
check_every_n_seconds=0.1,
max_bucket_size=5,
)
model = init_chat_model(
"openai:gpt-4.1-mini",
rate_limiter=rate_limiter,
)
@task
def analyze_topic(topic: str, index: int) -> dict:
"""Analiza un sub-tema con rate limiting automático."""
start = time.time()
response = model.invoke(
f"Resume en una oración el tema: {topic}"
)
elapsed = time.time() - start
return {
"index": index,
"topic": topic,
"summary": response.content,
"elapsed_seconds": round(elapsed, 2),
}
@entrypoint()
def rate_limited_pipeline(topics: list) -> dict:
start = time.time()
futures = [analyze_topic(t, i) for i, t in enumerate(topics)]
results = [f.result() for f in futures]
total_elapsed = time.time() - start
return {
"results": results,
"total_time_seconds": round(total_elapsed, 2),
"requests_made": len(results),
}
topics = [
"inteligencia artificial",
"computación cuántica",
"energías renovables",
"biotecnología",
"exploración espacial",
"ciberseguridad",
]
result = rate_limited_pipeline.invoke(topics)
print(f"Total: {result['total_time_seconds']}s para {result['requests_made']} requests")
for r in result["results"]:
print(f" [{r['index']}] {r['topic']}: {r['elapsed_seconds']}s")
# Output esperado:
# Total: ~3-4s para 6 requests (rate limited a 2/s)
# [0] inteligencia artificial: 0.8s
# [1] computación cuántica: 1.2s
# [2] energías renovables: 1.5s
# ...
Rate limiter por servicio
En producción, diferentes APIs tienen diferentes límites. Crea un rate limiter por servicio:
from langchain_core.rate_limiters import InMemoryRateLimiter
rate_limiters = {
"openai": InMemoryRateLimiter(
requests_per_second=5,
check_every_n_seconds=0.1,
max_bucket_size=10,
),
"search_api": InMemoryRateLimiter(
requests_per_second=1,
check_every_n_seconds=0.1,
max_bucket_size=3,
),
"news_api": InMemoryRateLimiter(
requests_per_second=2,
check_every_n_seconds=0.1,
max_bucket_size=5,
),
}
Cada modelo o cliente usa su propio rate limiter. Así, las llamadas a OpenAI no bloquean las llamadas a la API de búsqueda y viceversa.
Logging estructurado
El problema que resuelve
Tu agente falla. El log dice:
ERROR: Something went wrong
Inútil. No sabes qué nodo falló, cuánto tardó antes de fallar, qué input recibió, ni qué request del usuario causó el error. Con logging estructurado, el log dice:
{"node": "search_academic", "duration_ms": 4500, "status": "error", "error": "TimeoutError", "request_id": "abc123", "query": "quantum computing", "timestamp": "2026-03-08T14:30:00"}
Ahora puedes filtrar por request_id, ver que search_academic es el nodo problemático, y saber que tardó 4.5 segundos antes de fallar.
Implementación completa
from dotenv import load_dotenv
load_dotenv()
import json
import time
import uuid
import logging
from langgraph.func import entrypoint, task
logging.basicConfig(
level=logging.INFO,
format="%(message)s",
)
logger = logging.getLogger("research_agent")
class StructuredLogger:
"""Logger que emite JSON estructurado con contexto de request."""
def __init__(self, logger_instance: logging.Logger):
self._logger = logger_instance
self._request_id = None
def set_request_id(self, request_id: str):
self._request_id = request_id
def _log(self, level: str, node: str, **kwargs):
entry = {
"timestamp": time.strftime("%Y-%m-%dT%H:%M:%S"),
"level": level,
"request_id": self._request_id,
"node": node,
**kwargs,
}
self._logger.info(json.dumps(entry, ensure_ascii=False))
def node_start(self, node: str, **kwargs):
self._log("INFO", node, event="node_start", **kwargs)
def node_end(self, node: str, duration_ms: float, **kwargs):
self._log("INFO", node, event="node_end", duration_ms=round(duration_ms, 1), **kwargs)
def node_error(self, node: str, error: str, duration_ms: float, **kwargs):
self._log("ERROR", node, event="node_error", error=error, duration_ms=round(duration_ms, 1), **kwargs)
slog = StructuredLogger(logger)
@task
def search_with_logging(query: str, source: str) -> dict:
"""Búsqueda con logging estructurado automático."""
slog.node_start("search", source=source, query=query)
start = time.time()
try:
time.sleep(0.3)
if source == "academic" and "quantum" in query.lower():
raise ConnectionError("Academic API unavailable")
result = {
"source": source,
"content": f"Resultados de {source} para '{query}'",
}
duration_ms = (time.time() - start) * 1000
slog.node_end("search", duration_ms=duration_ms, source=source, status="success")
return result
except Exception as e:
duration_ms = (time.time() - start) * 1000
slog.node_error("search", error=str(e), duration_ms=duration_ms, source=source)
raise
@task
def synthesize_with_logging(results: list) -> str:
"""Síntesis con logging estructurado."""
slog.node_start("synthesize", num_sources=len(results))
start = time.time()
summary = f"Síntesis de {len(results)} fuentes completada."
duration_ms = (time.time() - start) * 1000
slog.node_end("synthesize", duration_ms=duration_ms, num_sources=len(results))
return summary
@entrypoint()
def logged_pipeline(topic: str) -> dict:
request_id = uuid.uuid4().hex[:8]
slog.set_request_id(request_id)
slog.node_start("pipeline", topic=topic)
pipeline_start = time.time()
sources = ["web", "academic", "news"]
futures = [search_with_logging(topic, s) for s in sources]
results = []
errors = []
for i, future in enumerate(futures):
try:
results.append(future.result())
except Exception as e:
errors.append({"source": sources[i], "error": str(e)})
summary = synthesize_with_logging(results).result()
pipeline_duration = (time.time() - pipeline_start) * 1000
slog.node_end("pipeline", duration_ms=pipeline_duration, sources_ok=len(results), sources_failed=len(errors))
return {
"request_id": request_id,
"summary": summary,
"sources_ok": len(results),
"errors": errors,
"pipeline_duration_ms": round(pipeline_duration, 1),
}
result = logged_pipeline.invoke("quantum computing")
print(f"\nRequest {result['request_id']}: {result['sources_ok']} fuentes OK, {len(result['errors'])} errores")
print(f"Duración total: {result['pipeline_duration_ms']}ms")
# Output esperado (logs en stderr, resultado en stdout):
# {"timestamp": "2026-03-08T14:30:00", "level": "INFO", "request_id": "a1b2c3d4", "node": "pipeline", "event": "node_start", "topic": "quantum computing"}
# {"timestamp": "2026-03-08T14:30:00", "level": "INFO", "request_id": "a1b2c3d4", "node": "search", "event": "node_start", "source": "web", "query": "quantum computing"}
# {"timestamp": "2026-03-08T14:30:00", "level": "INFO", "request_id": "a1b2c3d4", "node": "search", "event": "node_end", "duration_ms": 302.1, "source": "web", "status": "success"}
# {"timestamp": "2026-03-08T14:30:01", "level": "ERROR", "request_id": "a1b2c3d4", "node": "search", "event": "node_error", "error": "Academic API unavailable", "duration_ms": 300.5, "source": "academic"}
# {"timestamp": "2026-03-08T14:30:01", "level": "INFO", "request_id": "a1b2c3d4", "node": "search", "event": "node_end", "duration_ms": 301.3, "source": "news", "status": "success"}
# {"timestamp": "2026-03-08T14:30:01", "level": "INFO", "request_id": "a1b2c3d4", "node": "synthesize", "event": "node_start", "num_sources": 2}
# {"timestamp": "2026-03-08T14:30:01", "level": "INFO", "request_id": "a1b2c3d4", "node": "pipeline", "event": "node_end", "duration_ms": 920.5, "sources_ok": 2, "sources_failed": 1}
#
# Request a1b2c3d4: 2 fuentes OK, 1 errores
# Duración total: 920.5ms
Correlation IDs: trazando una request completa
El request_id es la pieza clave. Cuando un usuario reporta "mi investigación falló", te da el request_id y puedes filtrar todos los logs de esa request:
# Filtrar todos los logs de una request específica
cat logs.jsonl | jq 'select(.request_id == "a1b2c3d4")'
Verás cada nodo que se ejecutó, cuánto tardó, cuáles fallaron, y en qué orden. Sin correlation IDs, los logs de 100 requests concurrentes se mezclan y son imposibles de analizar.
Métricas de rendimiento
El problema que resuelve
"¿Cuánto tarda nuestro pipeline en promedio?" "¿Qué nodo es el cuello de botella?" "¿Cuántos tokens gastamos ayer?" Sin métricas, respondes con "no sé" o con una anécdota de la última vez que probaste. Con métricas, respondes con números.
Implementación: colector de métricas in-memory
from dotenv import load_dotenv
load_dotenv()
import time
import statistics
from collections import defaultdict
from langgraph.func import entrypoint, task
from langchain.chat_models import init_chat_model
class MetricsCollector:
"""Colector de métricas en memoria para pipelines de LangGraph."""
def __init__(self):
self._latencies: dict[str, list[float]] = defaultdict(list)
self._success_count: dict[str, int] = defaultdict(int)
self._error_count: dict[str, int] = defaultdict(int)
self._token_usage: dict[str, int] = defaultdict(int)
def record_latency(self, node: str, duration_ms: float):
self._latencies[node].append(duration_ms)
def record_success(self, node: str):
self._success_count[node] += 1
def record_error(self, node: str):
self._error_count[node] += 1
def record_tokens(self, node: str, tokens: int):
self._token_usage[node] += tokens
def get_summary(self) -> dict:
summary = {}
for node in set(
list(self._latencies.keys())
+ list(self._success_count.keys())
+ list(self._error_count.keys())
):
latencies = self._latencies.get(node, [])
successes = self._success_count.get(node, 0)
errors = self._error_count.get(node, 0)
total = successes + errors
summary[node] = {
"total_calls": total,
"successes": successes,
"errors": errors,
"success_rate": round(successes / total, 2) if total > 0 else 0,
"avg_latency_ms": round(statistics.mean(latencies), 1) if latencies else 0,
"p50_latency_ms": round(statistics.median(latencies), 1) if latencies else 0,
"p95_latency_ms": round(
sorted(latencies)[int(len(latencies) * 0.95)] if latencies else 0, 1
),
"max_latency_ms": round(max(latencies), 1) if latencies else 0,
"total_tokens": self._token_usage.get(node, 0),
}
return summary
def reset(self):
self._latencies.clear()
self._success_count.clear()
self._error_count.clear()
self._token_usage.clear()
metrics = MetricsCollector()
model = init_chat_model("openai:gpt-4.1-mini")
@task
def search_node(query: str, source: str) -> dict:
start = time.time()
try:
time.sleep(0.2)
result = {"source": source, "content": f"Resultado de {source}"}
duration_ms = (time.time() - start) * 1000
metrics.record_latency(f"search_{source}", duration_ms)
metrics.record_success(f"search_{source}")
return result
except Exception:
duration_ms = (time.time() - start) * 1000
metrics.record_latency(f"search_{source}", duration_ms)
metrics.record_error(f"search_{source}")
raise
@task
def synthesize_node(topic: str, results: list) -> str:
start = time.time()
context = "\n".join(r["content"] for r in results)
response = model.invoke(
f"Resume en una oración estos resultados sobre '{topic}': {context}"
)
duration_ms = (time.time() - start) * 1000
metrics.record_latency("synthesize", duration_ms)
metrics.record_success("synthesize")
token_count = response.usage_metadata.get("total_tokens", 0) if response.usage_metadata else 0
metrics.record_tokens("synthesize", token_count)
return response.content
@entrypoint()
def pipeline_with_metrics(topic: str) -> dict:
start = time.time()
sources = ["web", "academic", "news"]
futures = [search_node(topic, s) for s in sources]
results = [f.result() for f in futures]
summary = synthesize_node(topic, results).result()
pipeline_ms = (time.time() - start) * 1000
metrics.record_latency("pipeline_total", pipeline_ms)
metrics.record_success("pipeline_total")
return {"summary": summary}
for topic in ["AI en educación", "quantum computing", "energías renovables"]:
pipeline_with_metrics.invoke(topic)
summary = metrics.get_summary()
print("\n📊 MÉTRICAS DE RENDIMIENTO")
print("=" * 60)
for node, data in sorted(summary.items()):
print(f"\n {node}:")
print(f" Calls: {data['total_calls']} ({data['success_rate']:.0%} éxito)")
print(f" Latencia: avg={data['avg_latency_ms']}ms, p50={data['p50_latency_ms']}ms, p95={data['p95_latency_ms']}ms")
if data['total_tokens'] > 0:
print(f" Tokens: {data['total_tokens']}")
# Output esperado:
# 📊 MÉTRICAS DE RENDIMIENTO
# ============================================================
#
# pipeline_total:
# Calls: 3 (100% éxito)
# Latencia: avg=1200.5ms, p50=1180.3ms, p95=1250.1ms
#
# search_academic:
# Calls: 3 (100% éxito)
# Latencia: avg=201.2ms, p50=200.8ms, p95=202.1ms
#
# search_news:
# Calls: 3 (100% éxito)
# Latencia: avg=200.5ms, p50=200.3ms, p95=201.0ms
#
# search_web:
# Calls: 3 (100% éxito)
# Latencia: avg=200.8ms, p50=200.5ms, p95=201.5ms
#
# synthesize:
# Calls: 3 (100% éxito)
# Latencia: avg=980.3ms, p50=950.1ms, p95=1050.2ms
# Tokens: 450
Las métricas revelan que synthesize es el cuello de botella (980ms vs 200ms de búsqueda). Con esta información, puedes decidir si optimizar el prompt, usar un modelo más rápido para síntesis, o cachear resultados similares.
Combinando todo: grafo production-ready
Ahora combina los 4 patterns en un solo pipeline. Este es el patrón completo que usarías en producción:
from dotenv import load_dotenv
load_dotenv()
import json
import time
import uuid
import logging
import statistics
from collections import defaultdict
from concurrent.futures import ThreadPoolExecutor, TimeoutError as FuturesTimeoutError
from langchain_core.rate_limiters import InMemoryRateLimiter
from langchain.chat_models import init_chat_model
from langgraph.func import entrypoint, task
# === INFRAESTRUCTURA ===
logging.basicConfig(level=logging.INFO, format="%(message)s")
log = logging.getLogger("production_agent")
class ProdLogger:
def __init__(self):
self.request_id = None
def set_request(self, rid: str):
self.request_id = rid
def info(self, node: str, event: str, **kw):
entry = {"ts": time.strftime("%H:%M:%S"), "rid": self.request_id, "node": node, "event": event, **kw}
log.info(json.dumps(entry, ensure_ascii=False))
def error(self, node: str, event: str, **kw):
entry = {"ts": time.strftime("%H:%M:%S"), "rid": self.request_id, "node": node, "event": event, "level": "ERROR", **kw}
log.info(json.dumps(entry, ensure_ascii=False))
class ProdMetrics:
def __init__(self):
self.latencies = defaultdict(list)
self.counts = defaultdict(lambda: {"ok": 0, "err": 0})
self.tokens = defaultdict(int)
def record(self, node: str, duration_ms: float, success: bool, tokens: int = 0):
self.latencies[node].append(duration_ms)
self.counts[node]["ok" if success else "err"] += 1
if tokens:
self.tokens[node] += tokens
def summary(self) -> dict:
out = {}
for node in self.latencies:
lats = self.latencies[node]
c = self.counts[node]
total = c["ok"] + c["err"]
out[node] = {
"calls": total,
"success_rate": f"{c['ok']/total:.0%}" if total else "N/A",
"avg_ms": round(statistics.mean(lats), 1),
"p95_ms": round(sorted(lats)[int(len(lats) * 0.95)], 1) if lats else 0,
"tokens": self.tokens.get(node, 0),
}
return out
plog = ProdLogger()
pmetrics = ProdMetrics()
rate_limiter = InMemoryRateLimiter(
requests_per_second=3,
check_every_n_seconds=0.1,
max_bucket_size=5,
)
model = init_chat_model("openai:gpt-4.1-mini", rate_limiter=rate_limiter)
SEARCH_TIMEOUT_SECONDS = 5.0
# === TASKS ===
def _mock_search(query: str, source: str) -> str:
"""Simula búsqueda externa con latencia variable."""
import random
delay = random.uniform(0.2, 0.8)
time.sleep(delay)
return f"[{source}] Resultados para '{query}': información relevante encontrada."
@task
def production_search(query: str, source: str) -> dict:
"""Búsqueda con timeout + logging + métricas."""
plog.info("search", "start", source=source, query=query[:50])
start = time.time()
with ThreadPoolExecutor(max_workers=1) as executor:
future = executor.submit(_mock_search, query, source)
try:
content = future.result(timeout=SEARCH_TIMEOUT_SECONDS)
duration_ms = (time.time() - start) * 1000
plog.info("search", "end", source=source, duration_ms=round(duration_ms))
pmetrics.record("search", duration_ms, success=True)
return {"source": source, "content": content, "status": "ok"}
except FuturesTimeoutError:
duration_ms = (time.time() - start) * 1000
plog.error("search", "timeout", source=source, duration_ms=round(duration_ms))
pmetrics.record("search", duration_ms, success=False)
return {"source": source, "content": "", "status": "timeout"}
@task
def production_synthesize(topic: str, results: list) -> dict:
"""Síntesis con rate limiting + logging + métricas + tokens."""
plog.info("synthesize", "start", num_sources=len(results))
start = time.time()
context = "\n".join(r["content"] for r in results if r["content"])
response = model.invoke(
f"Resume en 2-3 oraciones estos resultados sobre '{topic}':\n{context}"
)
duration_ms = (time.time() - start) * 1000
tokens = response.usage_metadata.get("total_tokens", 0) if response.usage_metadata else 0
plog.info("synthesize", "end", duration_ms=round(duration_ms), tokens=tokens)
pmetrics.record("synthesize", duration_ms, success=True, tokens=tokens)
return {"summary": response.content, "tokens_used": tokens}
# === PIPELINE ===
@entrypoint()
def production_pipeline(topic: str) -> dict:
request_id = uuid.uuid4().hex[:8]
plog.set_request(request_id)
plog.info("pipeline", "start", topic=topic)
pipeline_start = time.time()
sources = ["web", "academic", "news"]
search_futures = [production_search(topic, s) for s in sources]
search_results = [f.result() for f in search_futures]
successful = [r for r in search_results if r["status"] == "ok"]
failed = [r for r in search_results if r["status"] != "ok"]
if not successful:
plog.error("pipeline", "all_sources_failed")
pmetrics.record("pipeline", (time.time() - pipeline_start) * 1000, success=False)
return {"error": "Todas las fuentes fallaron", "request_id": request_id}
synthesis = production_synthesize(topic, successful).result()
pipeline_ms = (time.time() - pipeline_start) * 1000
plog.info("pipeline", "end", duration_ms=round(pipeline_ms), sources_ok=len(successful), sources_failed=len(failed))
pmetrics.record("pipeline", pipeline_ms, success=True)
return {
"request_id": request_id,
"summary": synthesis["summary"],
"sources_ok": len(successful),
"sources_failed": len(failed),
"tokens_used": synthesis["tokens_used"],
"duration_ms": round(pipeline_ms),
}
result = production_pipeline.invoke("impacto de LLMs en producción de software")
print(f"\n{'=' * 50}")
print(f"Request: {result['request_id']}")
print(f"Fuentes: {result['sources_ok']} OK, {result['sources_failed']} fallidas")
print(f"Tokens: {result['tokens_used']}")
print(f"Duración: {result['duration_ms']}ms")
print(f"Resumen: {result['summary'][:120]}...")
print(f"\n{'=' * 50}")
print("MÉTRICAS ACUMULADAS:")
for node, data in pmetrics.summary().items():
print(f" {node}: {data['calls']} calls, {data['success_rate']} éxito, avg {data['avg_ms']}ms")
# Output esperado:
# {"ts": "14:30:00", "rid": "a1b2c3d4", "node": "pipeline", "event": "start", "topic": "impacto de LLMs en producción de software"}
# {"ts": "14:30:00", "rid": "a1b2c3d4", "node": "search", "event": "start", "source": "web", ...}
# {"ts": "14:30:00", "rid": "a1b2c3d4", "node": "search", "event": "end", "source": "web", "duration_ms": 450}
# ... (más logs)
#
# ==================================================
# Request: a1b2c3d4
# Fuentes: 3 OK, 0 fallidas
# Tokens: 150
# Duración: 1850ms
# Resumen: Los LLMs están transformando la producción de software...
#
# ==================================================
# MÉTRICAS ACUMULADAS:
# search: 3 calls, 100% éxito, avg 420.5ms
# synthesize: 1 calls, 100% éxito, avg 980.3ms
# pipeline: 1 calls, 100% éxito, avg 1850.2ms
Este pipeline tiene todo lo que necesitas para producción:
- ✅ Timeouts por búsqueda (5s máximo, no bloquea el pipeline)
- ✅ Rate limiting en las llamadas al modelo (3 req/s, respeta límites de OpenAI)
- ✅ Logging estructurado con correlation ID (filtrable por request)
- ✅ Métricas de latencia, success rate, y tokens (responde preguntas de negocio)
- ✅ Graceful degradation (si todas las fuentes fallan, no crashea)
Qué necesitarías sin LangGraph
Sin LangGraph, implementar este pipeline requiere:
| Componente | Con LangGraph | Sin LangGraph |
|---|---|---|
| Paralelismo | Futures de @task | concurrent.futures.ThreadPoolExecutor manual, gestión de threads |
| Checkpointing | Automático con MemorySaver | Base de datos + lógica de checkpoint + serialización custom |
| Retry | Python try/except dentro de @task con checkpoint | Retry library + state management manual para saber qué reintentar |
| Composición | @entrypoint anida @task naturalmente | Funciones anidadas con contexto compartido via globals o inyección |
| Timeout | asyncio.wait_for o ThreadPool en el @task | Igual, pero sin el beneficio de checkpoint automático si timeout |
| State recovery | Resume desde el último checkpoint | Reimplementar todo el state management desde cero |
El boilerplate sin LangGraph es 3-5x más código para la misma funcionalidad, y no incluye checkpointing (que es la parte más difícil de implementar correctamente).
Troubleshooting
Problema 1: "El rate limiter bloquea todo y el pipeline es muy lento"
Síntoma: Con requests_per_second=1 y 10 tasks paralelas, el pipeline tarda 10 segundos.
Causa: El rate limiter serializa las requests. Si tienes 10 tasks y solo 1 req/s permitida, cada task espera su turno.
Solución: Ajusta requests_per_second al límite real de tu API. Si la API permite 60 req/min, usa requests_per_second=1 (correcto). Si permite 600 req/min, usa requests_per_second=10. El max_bucket_size permite ráfagas cortas:
rate_limiter = InMemoryRateLimiter(
requests_per_second=10,
check_every_n_seconds=0.05,
max_bucket_size=20,
)
Problema 2: "Los logs de diferentes requests se mezclan y no puedo rastrear una"
Síntoma: Con 10 requests concurrentes, los logs JSON están intercalados y no sabes cuáles pertenecen a cuál request.
Causa: El request_id se pierde entre nodos, o estás usando un logger global sin correlation.
Solución: El StructuredLogger de esta cápsula usa set_request_id() al inicio del pipeline. Para concurrencia real, usa contextvars de Python para que cada thread tenga su propio request_id:
import contextvars
_request_id_var = contextvars.ContextVar("request_id", default="unknown")
class ThreadSafeLogger:
def set_request_id(self, rid: str):
_request_id_var.set(rid)
def _log(self, node: str, **kwargs):
entry = {"request_id": _request_id_var.get(), "node": node, **kwargs}
print(json.dumps(entry))
Problema 3: "Las métricas de p95 no son confiables"
Síntoma: El p95 reportado no coincide con la realidad — a veces es más bajo que el promedio.
Causa: Con pocas muestras (menos de 20), los percentiles estadísticos no son representativos. Con 3 muestras, el p95 es básicamente el valor más alto.
Solución: Solo reporta percentiles cuando tienes suficientes muestras:
def get_p95(self, node: str) -> str:
lats = self.latencies.get(node, [])
if len(lats) < 20:
return f"~{max(lats):.1f}ms (solo {len(lats)} muestras)"
return f"{sorted(lats)[int(len(lats) * 0.95)]:.1f}ms"
Problema 4: "El timeout no cancela la tarea, solo ignora el resultado"
Síntoma: Pones timeout de 3 segundos, pero la función sigue ejecutándose en background consumiendo recursos.
Causa: ThreadPoolExecutor.submit().result(timeout=3) lanza TimeoutError en el caller, pero el thread sigue corriendo. Python no puede matar threads de forma limpia.
Solución: Para timeouts reales que liberen recursos, usa asyncio con tasks cancelables:
async def search_with_real_cancel(query: str, timeout: float):
task = asyncio.create_task(async_search(query))
try:
return await asyncio.wait_for(task, timeout=timeout)
except asyncio.TimeoutError:
task.cancel()
return {"status": "timeout"}
Con asyncio, task.cancel() realmente cancela la coroutine. Con threads, no hay equivalente limpio.
Ejercicios
Ejercicio 1: Timeout configurable por fuente (Fácil)
Modifica el pipeline de timeouts para que cada fuente tenga su propio timeout: web = 2s, academic = 5s, news = 3s. Simula que academic tarda 4s (dentro del timeout) y web tarda 3s (fuera de su timeout).
Ver solución
from dotenv import load_dotenv
load_dotenv()
import time
from concurrent.futures import ThreadPoolExecutor, TimeoutError as FuturesTimeoutError
from langgraph.func import entrypoint, task
SOURCE_TIMEOUTS = {
"web": 2.0,
"academic": 5.0,
"news": 3.0,
}
SOURCE_DELAYS = {
"web": 3.0,
"academic": 4.0,
"news": 1.0,
}
def _mock_search(query: str, source: str) -> str:
time.sleep(SOURCE_DELAYS[source])
return f"[{source}] Resultados para '{query}'"
@task
def search_with_custom_timeout(query: str, source: str) -> dict:
timeout = SOURCE_TIMEOUTS[source]
with ThreadPoolExecutor(max_workers=1) as executor:
future = executor.submit(_mock_search, query, source)
try:
content = future.result(timeout=timeout)
return {"source": source, "status": "ok", "content": content}
except FuturesTimeoutError:
return {"source": source, "status": "timeout", "timeout_seconds": timeout}
@entrypoint()
def pipeline(topic: str) -> dict:
sources = ["web", "academic", "news"]
futures = [search_with_custom_timeout(topic, s) for s in sources]
results = [f.result() for f in futures]
ok = [r for r in results if r["status"] == "ok"]
failed = [r for r in results if r["status"] != "ok"]
return {
"ok_sources": [r["source"] for r in ok],
"failed_sources": [f"{r['source']} (timeout {r.get('timeout_seconds', '?')}s)" for r in failed],
}
result = pipeline.invoke("deep learning")
print(f"OK: {result['ok_sources']}")
print(f"Failed: {result['failed_sources']}")
# Output esperado:
# OK: ['academic', 'news'] (academic tarda 4s pero timeout es 5s; news tarda 1s)
# Failed: ['web (timeout 2.0s)'] (web tarda 3s pero timeout es 2s)
Ejercicio 2: Rate limiter con tracking visual (Fácil)
Crea un pipeline que haga 8 llamadas a un modelo con rate limit de 2/segundo. Muestra el timestamp de cada llamada para verificar visualmente que se respeta el rate limit.
Ver solución
from dotenv import load_dotenv
load_dotenv()
import time
from langchain_core.rate_limiters import InMemoryRateLimiter
from langchain.chat_models import init_chat_model
from langgraph.func import entrypoint, task
rate_limiter = InMemoryRateLimiter(
requests_per_second=2,
check_every_n_seconds=0.1,
max_bucket_size=2,
)
model = init_chat_model("openai:gpt-4.1-mini", rate_limiter=rate_limiter)
start_time = time.time()
@task
def timed_call(index: int) -> dict:
elapsed = time.time() - start_time
response = model.invoke(f"Di 'hola {index}' y nada más.")
after = time.time() - start_time
return {
"index": index,
"request_at": round(elapsed, 2),
"response_at": round(after, 2),
"content": response.content.strip(),
}
@entrypoint()
def rate_limit_demo(count: int) -> list:
futures = [timed_call(i) for i in range(count)]
return [f.result() for f in futures]
results = rate_limit_demo.invoke(8)
for r in sorted(results, key=lambda x: x["request_at"]):
print(f" Call {r['index']}: request @{r['request_at']}s → response @{r['response_at']}s | {r['content']}")
# Output esperado (tiempos aproximados — rate limited a 2/s):
# Call 0: request @0.01s → response @0.85s | hola 0
# Call 1: request @0.01s → response @0.90s | hola 1
# Call 2: request @0.51s → response @1.35s | hola 2
# Call 3: request @0.51s → response @1.40s | hola 3
# Call 4: request @1.01s → response @1.85s | hola 4
# ... (cada par de llamadas separado por ~0.5s)
Las llamadas salen en pares de 2 (el rate limit), con ~0.5 segundos entre cada par.
Ejercicio 3: Logger con niveles y filtrado (Medio)
Extiende el StructuredLogger para soportar niveles (DEBUG, INFO, WARN, ERROR) y un nivel mínimo configurable. Si el nivel mínimo es WARN, los logs de INFO y DEBUG no se emiten. Prueba con un pipeline que genere logs de todos los niveles.
Ver solución
from dotenv import load_dotenv
load_dotenv()
import json
import time
import logging
from langgraph.func import entrypoint, task
logging.basicConfig(level=logging.INFO, format="%(message)s")
class LeveledLogger:
LEVELS = {"DEBUG": 0, "INFO": 1, "WARN": 2, "ERROR": 3}
def __init__(self, min_level: str = "INFO"):
self._min = self.LEVELS.get(min_level, 1)
self._logger = logging.getLogger("leveled")
self._rid = None
def set_request_id(self, rid: str):
self._rid = rid
def _emit(self, level: str, node: str, message: str, **kw):
if self.LEVELS.get(level, 0) < self._min:
return
entry = {
"ts": time.strftime("%H:%M:%S"),
"level": level,
"rid": self._rid,
"node": node,
"msg": message,
**kw,
}
self._logger.info(json.dumps(entry, ensure_ascii=False))
def debug(self, node, msg, **kw): self._emit("DEBUG", node, msg, **kw)
def info(self, node, msg, **kw): self._emit("INFO", node, msg, **kw)
def warn(self, node, msg, **kw): self._emit("WARN", node, msg, **kw)
def error(self, node, msg, **kw): self._emit("ERROR", node, msg, **kw)
log_info = LeveledLogger(min_level="INFO")
log_warn = LeveledLogger(min_level="WARN")
@task
def demo_task(value: int, logger: LeveledLogger) -> dict:
logger.debug("demo", f"Procesando valor {value}", value=value)
logger.info("demo", f"Valor recibido: {value}")
if value > 5:
logger.warn("demo", f"Valor alto: {value}", threshold=5)
if value > 8:
logger.error("demo", f"Valor crítico: {value}", threshold=8)
return {"value": value, "processed": True}
@entrypoint()
def logging_demo(config: dict) -> dict:
logger = log_info if config.get("verbose") else log_warn
logger.set_request_id("demo-001")
values = config["values"]
futures = [demo_task(v, logger) for v in values]
results = [f.result() for f in futures]
return {"processed": len(results)}
print("=== Con min_level=INFO (verbose) ===")
logging_demo.invoke({"values": [2, 6, 9], "verbose": True})
print("\n=== Con min_level=WARN (solo warnings+) ===")
logging_demo.invoke({"values": [2, 6, 9], "verbose": False})
# Output esperado:
# === Con min_level=INFO (verbose) ===
# {"ts": "14:30:00", "level": "INFO", "rid": "demo-001", "node": "demo", "msg": "Valor recibido: 2"}
# {"ts": "14:30:00", "level": "INFO", "rid": "demo-001", "node": "demo", "msg": "Valor recibido: 6"}
# {"ts": "14:30:00", "level": "WARN", "rid": "demo-001", "node": "demo", "msg": "Valor alto: 6", "threshold": 5}
# {"ts": "14:30:00", "level": "INFO", "rid": "demo-001", "node": "demo", "msg": "Valor recibido: 9"}
# {"ts": "14:30:00", "level": "WARN", "rid": "demo-001", "node": "demo", "msg": "Valor alto: 9", "threshold": 5}
# {"ts": "14:30:00", "level": "ERROR", "rid": "demo-001", "node": "demo", "msg": "Valor crítico: 9", "threshold": 8}
#
# === Con min_level=WARN (solo warnings+) ===
# {"ts": "14:30:00", "level": "WARN", "rid": "demo-001", "node": "demo", "msg": "Valor alto: 6", "threshold": 5}
# {"ts": "14:30:00", "level": "WARN", "rid": "demo-001", "node": "demo", "msg": "Valor alto: 9", "threshold": 5}
# {"ts": "14:30:00", "level": "ERROR", "rid": "demo-001", "node": "demo", "msg": "Valor crítico: 9", "threshold": 8}
Con min_level="WARN", los logs de DEBUG e INFO se filtran. Solo ves warnings y errores — exactamente lo que quieres en producción cuando no estás debugging.
Ejercicio 4: Dashboard de métricas con resumen (Medio)
Extiende el MetricsCollector para generar un dashboard de texto que incluya: top 3 nodos más lentos, success rate total, total de tokens consumidos, y alertas para nodos con success rate < 90%.
Ver solución
from dotenv import load_dotenv
load_dotenv()
import time
import statistics
from collections import defaultdict
from langgraph.func import entrypoint, task
class DashboardMetrics:
def __init__(self):
self.latencies = defaultdict(list)
self.counts = defaultdict(lambda: {"ok": 0, "err": 0})
self.tokens = defaultdict(int)
def record(self, node: str, ms: float, ok: bool, tokens: int = 0):
self.latencies[node].append(ms)
self.counts[node]["ok" if ok else "err"] += 1
if tokens:
self.tokens[node] += tokens
def dashboard(self) -> str:
lines = []
lines.append("╔══════════════════════════════════════════════════╗")
lines.append("║ 📊 PRODUCTION DASHBOARD ║")
lines.append("╚══════════════════════════════════════════════════╝")
total_ok = sum(c["ok"] for c in self.counts.values())
total_err = sum(c["err"] for c in self.counts.values())
total_all = total_ok + total_err
total_tokens = sum(self.tokens.values())
lines.append(f"\n 🎯 Success rate global: {total_ok}/{total_all} ({total_ok/total_all:.0%})" if total_all else "")
lines.append(f" 🪙 Tokens totales: {total_tokens}")
lines.append(f" 📦 Nodos tracked: {len(self.latencies)}")
avg_by_node = {
n: statistics.mean(lats) for n, lats in self.latencies.items()
}
slowest = sorted(avg_by_node.items(), key=lambda x: x[1], reverse=True)[:3]
lines.append(f"\n 🐌 TOP 3 NODOS MÁS LENTOS:")
for i, (node, avg) in enumerate(slowest, 1):
lines.append(f" {i}. {node}: {avg:.1f}ms avg")
alerts = []
for node, c in self.counts.items():
total = c["ok"] + c["err"]
if total > 0 and c["ok"] / total < 0.9:
rate = c["ok"] / total
alerts.append(f" ⚠️ {node}: {rate:.0%} success rate ({c['err']} errores)")
if alerts:
lines.append(f"\n 🚨 ALERTAS:")
lines.extend(alerts)
else:
lines.append(f"\n ✅ Sin alertas — todos los nodos > 90% success rate")
return "\n".join(lines)
dm = DashboardMetrics()
@task
def reliable_node(name: str, delay: float) -> str:
time.sleep(delay)
dm.record(name, delay * 1000, ok=True, tokens=50)
return f"{name} OK"
@task
def flaky_node(name: str, delay: float, fail_rate: float) -> str:
import random
time.sleep(delay)
if random.random() < fail_rate:
dm.record(name, delay * 1000, ok=False)
raise RuntimeError(f"{name} falló")
dm.record(name, delay * 1000, ok=True, tokens=30)
return f"{name} OK"
@entrypoint()
def demo_dashboard(runs: int) -> str:
for _ in range(runs):
reliable_node("search_web", 0.2).result()
reliable_node("synthesize", 0.8).result()
try:
flaky_node("search_academic", 0.3, 0.3).result()
except Exception:
pass
return dm.dashboard()
print(demo_dashboard.invoke(10))
# Output esperado:
# ╔══════════════════════════════════════════════════╗
# ║ 📊 PRODUCTION DASHBOARD ║
# ╚══════════════════════════════════════════════════╝
#
# 🎯 Success rate global: 27/30 (90%)
# 🪙 Tokens totales: 1210
# 📦 Nodos tracked: 3
#
# 🐌 TOP 3 NODOS MÁS LENTOS:
# 1. synthesize: 800.0ms avg
# 2. search_academic: 300.0ms avg
# 3. search_web: 200.0ms avg
#
# 🚨 ALERTAS:
# ⚠️ search_academic: 70% success rate (3 errores)
Ejercicio 5: Pipeline con circuit breaker (Avanzado)
Implementa un circuit breaker: si un nodo falla más de 3 veces consecutivas, las siguientes llamadas a ese nodo retornan un error inmediatamente (sin intentar) durante 30 segundos. Después de 30 segundos, permite una llamada de prueba ("half-open"). Si esa llamada tiene éxito, el circuit se cierra (normal). Si falla, el circuit se abre de nuevo.
Ver solución
from dotenv import load_dotenv
load_dotenv()
import time
from langgraph.func import entrypoint, task
class CircuitBreaker:
"""Circuit breaker: CLOSED → OPEN → HALF_OPEN → CLOSED/OPEN."""
def __init__(self, failure_threshold: int = 3, recovery_timeout: float = 30.0):
self.failure_threshold = failure_threshold
self.recovery_timeout = recovery_timeout
self._states: dict[str, dict] = {}
def _get_state(self, name: str) -> dict:
if name not in self._states:
self._states[name] = {
"status": "CLOSED",
"consecutive_failures": 0,
"opened_at": 0,
}
return self._states[name]
def can_execute(self, name: str) -> tuple[bool, str]:
state = self._get_state(name)
if state["status"] == "CLOSED":
return True, "CLOSED"
if state["status"] == "OPEN":
elapsed = time.time() - state["opened_at"]
if elapsed >= self.recovery_timeout:
state["status"] = "HALF_OPEN"
return True, "HALF_OPEN"
return False, f"OPEN ({self.recovery_timeout - elapsed:.0f}s remaining)"
if state["status"] == "HALF_OPEN":
return True, "HALF_OPEN"
return False, state["status"]
def record_success(self, name: str):
state = self._get_state(name)
state["consecutive_failures"] = 0
state["status"] = "CLOSED"
def record_failure(self, name: str):
state = self._get_state(name)
state["consecutive_failures"] += 1
if state["consecutive_failures"] >= self.failure_threshold:
state["status"] = "OPEN"
state["opened_at"] = time.time()
cb = CircuitBreaker(failure_threshold=3, recovery_timeout=5.0)
call_count = 0
@task
def protected_search(query: str, source: str) -> dict:
global call_count
can_exec, status = cb.can_execute(source)
if not can_exec:
return {"source": source, "status": "circuit_open", "circuit": status}
try:
call_count += 1
if source == "flaky_api" and call_count <= 4:
raise ConnectionError(f"flaky_api fallo #{call_count}")
cb.record_success(source)
return {"source": source, "status": "ok", "content": f"Resultado de {source}"}
except Exception as e:
cb.record_failure(source)
return {"source": source, "status": "error", "error": str(e)}
@entrypoint()
def circuit_demo(config: dict) -> list:
results = []
for i in range(config["num_calls"]):
r = protected_search(config["query"], config["source"]).result()
results.append({"call": i + 1, **r})
if config.get("sleep_between"):
time.sleep(config["sleep_between"])
return results
results = circuit_demo.invoke({
"query": "test",
"source": "flaky_api",
"num_calls": 8,
"sleep_between": 1.0,
})
for r in results:
print(f" Call {r['call']}: {r['status']} | {r.get('circuit', r.get('error', r.get('content', '')))}")
# Output esperado:
# Call 1: error | flaky_api fallo #1
# Call 2: error | flaky_api fallo #2
# Call 3: error | flaky_api fallo #3 ← circuit OPENS here
# Call 4: circuit_open | OPEN (2s remaining) ← calls blocked
# Call 5: circuit_open | OPEN (1s remaining) ← calls blocked
# Call 6: error | flaky_api fallo #4 ← HALF_OPEN test call fails → re-opens
# Call 7: circuit_open | OPEN (4s remaining) ← blocked again
# Call 8: circuit_open | OPEN (3s remaining) ← blocked again
El circuit breaker protege al sistema de seguir llamando a un servicio que está caído. En vez de fallar 100 veces en 100 requests, falla 3 veces, bloquea las siguientes, y eventualmente prueba si el servicio se recuperó.
Ejercicio 6: Pipeline completo production-ready (Avanzado)
Combina timeouts, rate limiting, logging, métricas y circuit breaker en un pipeline que busque en 3 fuentes. Una fuente debe ser "flaky" (falla el 50% del tiempo). El pipeline debe: respetar timeouts, loguear cada operación, registrar métricas, y usar circuit breaker para la fuente flaky.
Ver solución
from dotenv import load_dotenv
load_dotenv()
import json
import time
import random
import logging
import statistics
from collections import defaultdict
from concurrent.futures import ThreadPoolExecutor, TimeoutError as FuturesTimeoutError
from langgraph.func import entrypoint, task
logging.basicConfig(level=logging.INFO, format="%(message)s")
class MiniCircuitBreaker:
def __init__(self, threshold=3, timeout=10.0):
self.threshold = threshold
self.timeout = timeout
self._fails = defaultdict(int)
self._opened = {}
def allow(self, name):
if name in self._opened:
if time.time() - self._opened[name] > self.timeout:
del self._opened[name]
return True
return False
return True
def success(self, name):
self._fails[name] = 0
self._opened.pop(name, None)
def failure(self, name):
self._fails[name] += 1
if self._fails[name] >= self.threshold:
self._opened[name] = time.time()
class MiniMetrics:
def __init__(self):
self.data = defaultdict(lambda: {"lats": [], "ok": 0, "err": 0})
def record(self, node, ms, ok):
self.data[node]["lats"].append(ms)
self.data[node]["ok" if ok else "err"] += 1
def report(self):
for node, d in sorted(self.data.items()):
total = d["ok"] + d["err"]
rate = d["ok"] / total if total else 0
avg = statistics.mean(d["lats"]) if d["lats"] else 0
print(f" {node}: {total} calls, {rate:.0%} ok, avg {avg:.0f}ms")
cb = MiniCircuitBreaker(threshold=2, timeout=5.0)
met = MiniMetrics()
def _search_impl(query, source):
delays = {"reliable_a": 0.3, "flaky_b": 0.4, "reliable_c": 0.2}
time.sleep(delays.get(source, 0.3))
if source == "flaky_b" and random.random() < 0.5:
raise ConnectionError(f"{source} failed")
return f"[{source}] Results for '{query}'"
@task
def guarded_search(query: str, source: str) -> dict:
if not cb.allow(source):
met.record(source, 0, False)
return {"source": source, "status": "circuit_open"}
start = time.time()
with ThreadPoolExecutor(max_workers=1) as ex:
fut = ex.submit(_search_impl, query, source)
try:
content = fut.result(timeout=2.0)
ms = (time.time() - start) * 1000
cb.success(source)
met.record(source, ms, True)
return {"source": source, "status": "ok", "content": content}
except FuturesTimeoutError:
ms = (time.time() - start) * 1000
cb.failure(source)
met.record(source, ms, False)
return {"source": source, "status": "timeout"}
except Exception as e:
ms = (time.time() - start) * 1000
cb.failure(source)
met.record(source, ms, False)
return {"source": source, "status": "error", "error": str(e)}
@entrypoint()
def full_production_pipeline(config: dict) -> dict:
topic = config["topic"]
sources = config["sources"]
all_run_results = []
for run in range(config.get("runs", 1)):
futures = [guarded_search(topic, s) for s in sources]
results = [f.result() for f in futures]
ok = [r for r in results if r["status"] == "ok"]
all_run_results.append({"run": run + 1, "ok": len(ok), "total": len(results)})
return {"runs": all_run_results}
result = full_production_pipeline.invoke({
"topic": "AI production patterns",
"sources": ["reliable_a", "flaky_b", "reliable_c"],
"runs": 6,
})
print("\n📊 Resultados por run:")
for r in result["runs"]:
print(f" Run {r['run']}: {r['ok']}/{r['total']} fuentes OK")
print("\n📈 Métricas acumuladas:")
met.report()
# Output esperado (varía por el random de flaky_b):
# 📊 Resultados por run:
# Run 1: 3/3 fuentes OK (o 2/3 si flaky_b falló)
# Run 2: 2/3 fuentes OK
# Run 3: 2/3 fuentes OK (circuit opens for flaky_b)
# Run 4: 2/3 fuentes OK (circuit_open — no intenta)
# Run 5: 2/3 fuentes OK (circuit_open)
# Run 6: 2/3 fuentes OK (circuit_open or half_open)
#
# 📈 Métricas acumuladas:
# flaky_b: 6 calls, 33% ok, avg 180ms
# reliable_a: 6 calls, 100% ok, avg 302ms
# reliable_c: 6 calls, 100% ok, avg 201ms
Este es un pipeline production-ready completo: timeout por búsqueda, circuit breaker para fuentes inestables, y métricas para monitorear la salud del sistema.
Resumen
En esta cápsula aprendiste los 4 patterns que separan un prototipo funcional de un sistema de producción:
- Timeouts por nodo — Previenen que una API lenta bloquee todo el pipeline. Usa
asyncio.wait_forpara código async oThreadPoolExecutorpara código síncrono. Sin timeout, un nodo lento congela a todos los usuarios - Rate limiting interno —
InMemoryRateLimiterde LangChain controla cuántas requests por segundo envías a cada servicio. Sin rate limiting, excedes los límites de la API y te bloquean - Logging estructurado — JSON con
request_id, nombre del nodo, duración, y status. Correlation IDs para rastrear una request a través de todos los nodos. Sin logging structured, "algo falló" es tu mejor diagnóstico - Métricas de rendimiento — Latencia por nodo (avg, p50, p95), success rate, tokens consumidos. Identifican cuellos de botella y responden preguntas de negocio. Sin métricas, optimizas a ciegas
Combinados, estos patterns hacen que tu agente sea operable: puedes diagnosticar problemas, responder preguntas de negocio, y escalar con confianza. Son el equivalente a instrumentar un servidor web con logs, métricas, y health checks — estándar en la industria, pero frecuentemente olvidado en AI engineering.
Próxima cápsula: Proyecto Evolutivo — vas a aplicar retry logic, branching paralelo, y estos patterns de producción al Research Agent v1 para convertirlo en v2.
Recursos adicionales
- LangChain Rate Limiters — Documentación de
InMemoryRateLimitery su integración con modelos - Python asyncio.wait_for — Timeouts con asyncio para código asíncrono
- Python logging JSON — Guía oficial de logging con formatters custom
- Circuit Breaker Pattern — Martin Fowler sobre circuit breakers (el artículo original)
- LangGraph Functional API — Referencia de
@entrypointy@task - Python concurrent.futures — ThreadPoolExecutor para timeouts síncronos
Módulo 7 — LangChain & LangGraph: From Chains to Agents