Módulo 6: Prompt Composition y Chaining
4. Pipelines Multi-Etapa
Descripción
Un pipeline multi-etapa es un sistema de procesamiento donde múltiples prompts se ejecutan en secuencia, pasando estado entre ellos, con manejo de errores por etapa, retry automático, y monitoreo. Es el patrón arquitectónico central para construir sistemas de LLM de producción.
En esta cápsula aprenderás a diseñar pipelines robustos con gestión de estado, retry inteligente, logging, y estrategias de fallback que garantizan que el sistema funcione incluso cuando partes individuales fallan.
La Diferencia entre un Chain Simple y un Pipeline
Chain simple:
A → B → C
Si B falla, todo falla.
No hay state management.
No hay retry.
No hay logging.
Pipeline de producción:
[Step A] → [Validate] → [Step B] → [Validate] → [Step C]
↑ ↑ ↑
Retry(2) Retry(2) Retry(2)
↑ ↑ ↑
Logging Logging Logging
↑
Fallback
State: {input_original, output_A, output_B, output_C, errores, timestamps}
Arquitectura Base del Pipeline
from openai import OpenAI
from typing import Callable, Any, Optional
from dataclasses import dataclass, field
from enum import Enum
import json
import time
import logging
client = OpenAI()
logger = logging.getLogger(__name__)
class EstadoEtapa(Enum):
PENDIENTE = "pendiente"
EN_PROGRESO = "en_progreso"
COMPLETADO = "completado"
FALLIDO = "fallido"
OMITIDO = "omitido" # Cuando se usa fallback
@dataclass
class MetricasEtapa:
"""Métricas de ejecución de una etapa."""
tiempo_inicio: float = 0.0
tiempo_fin: float = 0.0
intentos: int = 0
tokens_entrada: int = 0
tokens_salida: int = 0
estado: EstadoEtapa = EstadoEtapa.PENDIENTE
error: Optional[str] = None
@property
def duracion(self) -> float:
return self.tiempo_fin - self.tiempo_inicio
@dataclass
class PipelineState:
"""
Estado completo del pipeline durante su ejecución.
Contiene el input original, todos los outputs intermedios,
y métricas de ejecución.
"""
input_original: str
outputs: dict[str, Any] = field(default_factory=dict)
metricas: dict[str, MetricasEtapa] = field(default_factory=dict)
metadatos: dict[str, Any] = field(default_factory=dict)
def get_output(self, nombre_etapa: str, default: str = "") -> Any:
"""Obtiene el output de una etapa, o default si no existe."""
return self.outputs.get(nombre_etapa, default)
def get_ultimo_output(self) -> str:
"""Retorna el output de la última etapa exitosa."""
for nombre in reversed(list(self.outputs.keys())):
if nombre != "_error":
return self.outputs[nombre]
return self.input_original
def resumen_metricas(self) -> dict:
"""Genera un resumen de métricas del pipeline."""
etapas_completadas = sum(1 for m in self.metricas.values() if m.estado == EstadoEtapa.COMPLETADO)
etapas_fallidas = sum(1 for m in self.metricas.values() if m.estado == EstadoEtapa.FALLIDO)
tiempo_total = sum(m.duracion for m in self.metricas.values())
return {
"etapas_totales": len(self.metricas),
"etapas_completadas": etapas_completadas,
"etapas_fallidas": etapas_fallidas,
"tiempo_total_segundos": tiempo_total,
"tasa_exito": etapas_completadas / len(self.metricas) if self.metricas else 0
}
@dataclass
class ConfigEtapa:
"""Configuración de una etapa del pipeline."""
nombre: str
fn: Callable[[PipelineState], str] # Función que ejecuta la etapa
max_retries: int = 2 # Reintentos en caso de error
fallback_fn: Optional[Callable] = None # Función de fallback si falla
validar_fn: Optional[Callable[[str], bool]] = None # Validador del output
timeout_segundos: float = 30.0 # Timeout de la etapa
critica: bool = True # Si falla y no tiene fallback, ¿detener el pipeline?
Motor del Pipeline
class Pipeline:
"""
Motor de ejecución de pipelines de LLM.
Features:
- Retry automático por etapa
- Fallback cuando se agotan retries
- Validación de outputs
- Logging detallado
- Monitoreo de métricas
- State management
"""
def __init__(self, nombre: str = "pipeline"):
self.nombre = nombre
self.etapas: list[ConfigEtapa] = []
def agregar_etapa(self, config: ConfigEtapa) -> 'Pipeline':
"""Agrega una etapa al pipeline (fluent interface)."""
self.etapas.append(config)
return self
def ejecutar(
self,
input_inicial: str,
verbose: bool = True
) -> PipelineState:
"""
Ejecuta el pipeline completo.
Args:
input_inicial: El input de la primera etapa
verbose: Si True, imprime el progreso
Returns:
PipelineState con todos los outputs y métricas
"""
state = PipelineState(input_original=input_inicial)
if verbose:
print(f"\n{'='*60}")
print(f"Pipeline: {self.nombre}")
print(f"Etapas: {[e.nombre for e in self.etapas]}")
print(f"{'='*60}")
for config in self.etapas:
metricas = MetricasEtapa(
tiempo_inicio=time.time(),
estado=EstadoEtapa.EN_PROGRESO
)
state.metricas[config.nombre] = metricas
if verbose:
print(f"\n[{config.nombre}] Iniciando...")
exito = False
ultimo_error = None
# Ciclo de retry
for intento in range(config.max_retries + 1):
metricas.intentos = intento + 1
try:
output = config.fn(state)
# Validación opcional
if config.validar_fn and not config.validar_fn(output):
raise ValueError(f"Validación falló para etapa '{config.nombre}'")
# Éxito
state.outputs[config.nombre] = output
metricas.estado = EstadoEtapa.COMPLETADO
metricas.tiempo_fin = time.time()
exito = True
if verbose:
print(f" ✓ Completado en {metricas.duracion:.2f}s (intento {intento+1})")
print(f" Output: {str(output)[:100]}...")
break
except Exception as e:
ultimo_error = str(e)
logger.warning(f"Intento {intento+1} fallido en '{config.nombre}': {e}")
if verbose:
print(f" ⚠ Error intento {intento+1}: {str(e)[:80]}")
if not exito:
metricas.error = ultimo_error
# Intentar fallback
if config.fallback_fn:
try:
output_fallback = config.fallback_fn(state)
state.outputs[config.nombre] = output_fallback
metricas.estado = EstadoEtapa.OMITIDO
metricas.tiempo_fin = time.time()
if verbose:
print(f" ⚡ Usando fallback para '{config.nombre}'")
except Exception as fb_error:
metricas.estado = EstadoEtapa.FALLIDO
if verbose:
print(f" ✗ Fallback también falló: {fb_error}")
if config.critica:
raise RuntimeError(
f"Etapa crítica '{config.nombre}' falló después de {config.max_retries+1} intentos: {ultimo_error}"
)
else:
metricas.estado = EstadoEtapa.FALLIDO
if config.critica:
raise RuntimeError(
f"Etapa crítica '{config.nombre}' falló: {ultimo_error}"
)
if verbose:
resumen = state.resumen_metricas()
print(f"\n{'='*60}")
print(f"Pipeline completado")
print(f" Tiempo total: {resumen['tiempo_total_segundos']:.2f}s")
print(f" Éxito: {resumen['etapas_completadas']}/{resumen['etapas_totales']} etapas")
print(f"{'='*60}")
return state
Ejemplo Completo: Pipeline de Análisis de Documentos
def construir_pipeline_analisis() -> Pipeline:
"""
Construye un pipeline para análisis de documentos con:
- Extracción de información clave
- Análisis de sentimiento y contexto
- Generación de reporte estructurado
- Formateo de output final
"""
# Funciones de cada etapa
def step_preprocesar(state: PipelineState) -> str:
"""Normaliza y prepara el documento."""
doc = state.input_original
# Si el documento es muy largo, resumir primero
if len(doc) > 8000:
return client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": f"Resume este documento en 2000 palabras preservando toda la información importante:\n\n{doc[:15000]}"}],
temperature=0,
max_tokens=600
).choices[0].message.content
return doc
def step_extraer(state: PipelineState) -> str:
"""Extrae información estructurada del documento."""
doc = state.get_output("preprocesar", state.input_original)
return client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": f"""Extrae información clave del siguiente documento:
1. Tipo de documento
2. Entidades principales (personas, organizaciones, lugares)
3. Fechas importantes
4. Números y métricas clave
5. Temas principales (3-5)
Documento: {doc[:6000]}
Responde en JSON."""}],
temperature=0,
max_tokens=600,
response_format={"type": "json_object"}
).choices[0].message.content
def step_analizar(state: PipelineState) -> str:
"""Análisis profundo basado en la extracción."""
extraccion = state.get_output("extraer")
doc = state.get_output("preprocesar", state.input_original)
return client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": f"""Analiza este documento basándote en la extracción:
Extracción: {extraccion[:1000]}
Documento: {doc[:4000]}
Genera:
1. Análisis de contexto y relevancia
2. Insights más importantes (3-5)
3. Riesgos o puntos de atención
4. Oportunidades si aplica"""}],
temperature=0,
max_tokens=600
).choices[0].message.content
def step_sintetizar(state: PipelineState) -> str:
"""Sintetiza extracción + análisis en conclusiones."""
analisis = state.get_output("analizar")
extraccion = state.get_output("extraer")
return client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": f"""Basándote en este análisis completo:
Extracción: {extraccion[:500]}
Análisis: {analisis[:800]}
Genera:
1. Las 5 conclusiones más importantes (ordenadas por relevancia)
2. Las 3 acciones recomendadas
3. Un párrafo de conclusión ejecutiva
Sé conciso y accionable."""}],
temperature=0.1,
max_tokens=500
).choices[0].message.content
def step_formatear(state: PipelineState) -> str:
"""Formatea el output final en markdown estructurado."""
sintesis = state.get_output("sintetizar")
extraccion_raw = state.get_output("extraer", "{}")
try:
extraccion = json.loads(extraccion_raw)
entidades = extraccion.get("entidades", {})
except Exception:
entidades = {}
return client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": f"""Formatea este análisis en un informe markdown profesional.
Incluye:
- Título descriptivo (H1)
- Sección "Resumen Ejecutivo" (1 párrafo)
- Sección "Entidades Clave" (si hay datos: {str(entidades)[:200]})
- Sección "Conclusiones"
- Sección "Recomendaciones"
Análisis a formatear:
{sintesis}"""}],
temperature=0.2,
max_tokens=700
).choices[0].message.content
# Fallbacks
def fallback_preprocesar(state: PipelineState) -> str:
"""Si el preprocesamiento falla, usar el documento original directamente."""
return state.input_original[:6000] # Truncar si es necesario
def fallback_sintetizar(state: PipelineState) -> str:
"""Si la síntesis falla, concatenar extracción y análisis."""
return f"Extracción:\n{state.get_output('extraer')}\n\nAnálisis:\n{state.get_output('analizar')}"
# Validadores
def validar_json(output: str) -> bool:
try:
json.loads(output)
return True
except Exception:
return False
def validar_no_vacio(output: str) -> bool:
return len(output.strip()) > 50
# Construir pipeline
pipeline = Pipeline(nombre="analisis_documentos")
pipeline.agregar_etapa(ConfigEtapa(
nombre="preprocesar",
fn=step_preprocesar,
max_retries=1,
fallback_fn=fallback_preprocesar,
critica=False # No es crítico; podemos usar el documento original
))
pipeline.agregar_etapa(ConfigEtapa(
nombre="extraer",
fn=step_extraer,
max_retries=2,
validar_fn=validar_json,
critica=True # Crítico: sin extracción no puede continuar correctamente
))
pipeline.agregar_etapa(ConfigEtapa(
nombre="analizar",
fn=step_analizar,
max_retries=2,
validar_fn=validar_no_vacio,
critica=True
))
pipeline.agregar_etapa(ConfigEtapa(
nombre="sintetizar",
fn=step_sintetizar,
max_retries=1,
fallback_fn=fallback_sintetizar,
critica=False
))
pipeline.agregar_etapa(ConfigEtapa(
nombre="formatear",
fn=step_formatear,
max_retries=1,
critica=False # Si falla el formateo, el análisis ya está completo
))
return pipeline
# Ejemplo de uso:
if __name__ == "__main__":
documento = """
Informe Trimestral Q4 2025 - TechStartup S.L.
Resumen Ejecutivo:
El cuarto trimestre de 2025 fue un período de crecimiento acelerado para TechStartup.
Los ingresos alcanzaron €2.4M, un incremento del 45% respecto al Q4 2024.
El número de usuarios activos llegó a 125,000, con una tasa de retención del 87%.
Inversión en I+D:
Se destinó el 23% de los ingresos (€552K) a investigación y desarrollo,
centrado principalmente en el desarrollo de modelos de IA personalizados...
"""
pipeline = construir_pipeline_analisis()
state = pipeline.ejecutar(documento, verbose=True)
print("\n=== OUTPUT FINAL ===")
print(state.get_output("formatear", state.get_ultimo_output()))
print("\n=== MÉTRICAS ===")
for nombre, metricas in state.metricas.items():
print(f"{nombre}: {metricas.estado.value} ({metricas.duracion:.2f}s, {metricas.intentos} intentos)")
Retry con Backoff Exponencial
import random
def ejecutar_con_backoff(
fn: Callable,
state: PipelineState,
max_retries: int = 3,
backoff_base: float = 1.0,
jitter: bool = True
) -> str:
"""
Ejecuta una función con retry y backoff exponencial.
El backoff exponencial previene thundering herd problems
cuando múltiples pipelines fallan simultáneamente.
Args:
fn: Función a ejecutar
state: Estado del pipeline
max_retries: Máximo número de intentos
backoff_base: Base del tiempo de espera (segundos)
jitter: Si True, añade variación aleatoria al wait time
Returns:
Output de la función
"""
for intento in range(max_retries):
try:
return fn(state)
except Exception as e:
if intento == max_retries - 1:
raise # Re-raise en el último intento
# Calcular tiempo de espera: base * 2^intento
wait_time = backoff_base * (2 ** intento)
if jitter:
# Añadir jitter: ±50% del wait time
wait_time *= (1 + random.uniform(-0.5, 0.5))
logger.warning(f"Intento {intento+1} fallido: {e}. Reintentando en {wait_time:.1f}s")
time.sleep(wait_time)
raise RuntimeError("Max retries alcanzado sin éxito")
Monitoring y Observabilidad
class PipelineMonitor:
"""
Sistema de monitoreo para pipelines de LLM.
Registra métricas, errores, y performance.
"""
def __init__(self):
self.ejecuciones: list[dict] = []
def registrar_ejecucion(self, state: PipelineState, pipeline_nombre: str):
"""Registra las métricas de una ejecución."""
resumen = state.resumen_metricas()
entrada = {
"pipeline": pipeline_nombre,
"timestamp": time.time(),
"input_length": len(state.input_original),
**resumen,
"etapas_detalle": {
nombre: {
"estado": m.estado.value,
"duracion": m.duracion,
"intentos": m.intentos,
"error": m.error
}
for nombre, m in state.metricas.items()
}
}
self.ejecuciones.append(entrada)
return entrada
def estadisticas(self) -> dict:
"""Calcula estadísticas agregadas de todas las ejecuciones."""
if not self.ejecuciones:
return {}
tiempos = [e["tiempo_total_segundos"] for e in self.ejecuciones]
tasas_exito = [e["tasa_exito"] for e in self.ejecuciones]
return {
"n_ejecuciones": len(self.ejecuciones),
"tiempo_promedio": sum(tiempos) / len(tiempos),
"tiempo_p95": sorted(tiempos)[int(len(tiempos) * 0.95)],
"tasa_exito_promedio": sum(tasas_exito) / len(tasas_exito),
"n_ejecuciones_con_fallback": sum(
1 for e in self.ejecuciones
if any(d["estado"] == "omitido" for d in e["etapas_detalle"].values())
)
}
def etapas_mas_lentas(self, n: int = 3) -> list[tuple[str, float]]:
"""Identifica las etapas con mayor latencia promedio."""
tiempos_por_etapa: dict[str, list[float]] = {}
for ejecucion in self.ejecuciones:
for etapa, detalle in ejecucion["etapas_detalle"].items():
if etapa not in tiempos_por_etapa:
tiempos_por_etapa[etapa] = []
tiempos_por_etapa[etapa].append(detalle["duracion"])
promedios = [(etapa, sum(ts)/len(ts)) for etapa, ts in tiempos_por_etapa.items()]
return sorted(promedios, key=lambda x: x[1], reverse=True)[:n]
# Uso del monitor:
monitor = PipelineMonitor()
pipeline = construir_pipeline_analisis()
state = pipeline.ejecutar("Documento de prueba...", verbose=False)
monitor.registrar_ejecucion(state, "analisis_documentos")
print(monitor.estadisticas())
Pipeline con Cache Intermedio
import hashlib
import shelve
from pathlib import Path
class PipelineConCache(Pipeline):
"""
Pipeline con cache para evitar re-ejecutar etapas costosas
cuando el input no ha cambiado.
"""
def __init__(self, nombre: str, cache_dir: str = ".pipeline_cache"):
super().__init__(nombre)
self.cache_dir = Path(cache_dir)
self.cache_dir.mkdir(exist_ok=True)
def _cache_key(self, etapa_nombre: str, input_data: str) -> str:
"""Genera una clave única para el cache."""
return hashlib.md5(f"{etapa_nombre}:{input_data}".encode()).hexdigest()
def ejecutar_con_cache(
self,
input_inicial: str,
etapas_cacheables: list[str] = None,
ttl_horas: int = 24,
verbose: bool = True
) -> PipelineState:
"""
Ejecuta el pipeline con cache para etapas especificadas.
Las etapas cacheadas no se re-ejecutan si el output existe en cache.
"""
state = PipelineState(input_original=input_inicial)
cache_file = str(self.cache_dir / f"{self.nombre}.shelve")
with shelve.open(cache_file) as cache:
for config in self.etapas:
input_para_esta_etapa = state.get_ultimo_output()
cache_key = self._cache_key(config.nombre, input_para_esta_etapa)
# Verificar cache
usar_cache = (
etapas_cacheables and
config.nombre in etapas_cacheables and
cache_key in cache
)
if usar_cache:
entrada_cache = cache[cache_key]
edad_horas = (time.time() - entrada_cache["timestamp"]) / 3600
if edad_horas < ttl_horas:
state.outputs[config.nombre] = entrada_cache["output"]
if verbose:
print(f"[{config.nombre}] ✓ Cache hit ({edad_horas:.1f}h de antigüedad)")
continue
# Ejecutar etapa normalmente
try:
output = config.fn(state)
state.outputs[config.nombre] = output
# Guardar en cache si es cacheable
if etapas_cacheables and config.nombre in etapas_cacheables:
cache[cache_key] = {"output": output, "timestamp": time.time()}
if verbose:
print(f"[{config.nombre}] ✓ Completado (guardado en cache)")
except Exception as e:
if verbose:
print(f"[{config.nombre}] ✗ Error: {e}")
raise
return state
Troubleshooting
Problema 1: State explosion (el state crece demasiado)
Síntoma: El state acumula outputs de todas las etapas y se vuelve enorme, causando problemas de memoria y tokens excesivos en prompts posteriores.
Solución:
class PipelineStateOptimizado(PipelineState):
"""State con límite de tamaño para outputs."""
MAX_OUTPUT_CHARS = 2000
def set_output(self, nombre: str, output: str, full_output: str = None):
"""Guarda el output con truncamiento automático."""
if len(output) > self.MAX_OUTPUT_CHARS:
# Guardar versión truncada para uso en prompts
self.outputs[nombre] = output[:self.MAX_OUTPUT_CHARS] + "\n[...truncado...]"
# Opcionalmente guardar versión completa en metadatos
if full_output:
self.metadatos[f"{nombre}_completo"] = full_output
else:
self.outputs[nombre] = output
Problema 2: Latencia alta en pipelines secuenciales
Síntoma: El pipeline tarda 45 segundos para 5 etapas de 10s cada una.
Solución: Paralelizar etapas independientes:
import asyncio
async def pipeline_con_paralelizacion_automatica(
etapas_grupos: list[list[ConfigEtapa]],
input_inicial: str
) -> PipelineState:
"""
Ejecuta grupos de etapas en paralelo cuando es posible.
etapas_grupos: [[etapa1, etapa2_paralela], [etapa3_depende], [etapa4]]
"""
from openai import AsyncOpenAI
client_async = AsyncOpenAI()
state = PipelineState(input_original=input_inicial)
for grupo in etapas_grupos:
if len(grupo) == 1:
# Una sola etapa: ejecutar normalmente
output = grupo[0].fn(state)
state.outputs[grupo[0].nombre] = output
else:
# Múltiples etapas: ejecutar en paralelo
async def run_etapa(config):
loop = asyncio.get_event_loop()
output = await loop.run_in_executor(None, config.fn, state)
return config.nombre, output
tareas = [run_etapa(c) for c in grupo]
resultados = await asyncio.gather(*tareas)
for nombre, output in resultados:
state.outputs[nombre] = output
return state
Problema 3: Retry loop infinito
Síntoma: El pipeline reintenta indefinidamente una etapa que nunca va a tener éxito.
Solución: Implementar circuit breaker:
class CircuitBreaker:
"""
Implementa el patrón Circuit Breaker para pipelines.
Después de N fallos consecutivos, abre el circuito y
no permite más intentos por T segundos.
"""
def __init__(self, umbral_fallos: int = 3, tiempo_reset: float = 60.0):
self.umbral_fallos = umbral_fallos
self.tiempo_reset = tiempo_reset
self.fallos = 0
self.ultimo_fallo = 0
self.abierto = False
def puede_ejecutar(self) -> bool:
if not self.abierto:
return True
# Verificar si ya pasó el tiempo de reset
if time.time() - self.ultimo_fallo > self.tiempo_reset:
self.abierto = False
self.fallos = 0
return True
return False
def registrar_exito(self):
self.fallos = 0
self.abierto = False
def registrar_fallo(self):
self.fallos += 1
self.ultimo_fallo = time.time()
if self.fallos >= self.umbral_fallos:
self.abierto = True
Ejercicios
Ejercicio 1: Pipeline de procesamiento de emails
Construye un pipeline de 3 etapas para procesar emails de soporte:
- Clasificar urgencia y tipo
- Extraer información clave (cliente, problema, contexto)
- Generar respuesta automática personalizada
Incluye retry y un fallback de respuesta genérica.
Ver solución
def pipeline_soporte_email(email: str) -> str:
def step_clasificar(state: PipelineState) -> str:
return client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": f"Clasifica este email:\nurgencia: alta/media/baja\ntipo: bug/consulta/queja/otro\nJSON: {{\"urgencia\": str, \"tipo\": str}}\n\nEmail: {state.input_original}"}],
temperature=0, response_format={"type": "json_object"}
).choices[0].message.content
def step_extraer(state: PipelineState) -> str:
return client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": f"Extrae del email: cliente_nombre, problema_principal, contexto. JSON.\n\nEmail: {state.input_original}"}],
temperature=0, response_format={"type": "json_object"}
).choices[0].message.content
def step_responder(state: PipelineState) -> str:
clasif = json.loads(state.get_output("clasificar", "{}"))
datos = json.loads(state.get_output("extraer", "{}"))
urgencia = clasif.get("urgencia", "media")
return client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": f"Genera respuesta empática para email de soporte (urgencia: {urgencia}). Cliente: {datos.get('cliente_nombre', 'estimado cliente')}. Problema: {datos.get('problema_principal', 'su consulta')}. 2-3 oraciones."}],
temperature=0.2, max_tokens=150
).choices[0].message.content
def fallback_responder(state: PipelineState) -> str:
return "Estimado cliente, hemos recibido su mensaje y nuestro equipo lo atenderá a la brevedad posible. Gracias por contactarnos."
p = Pipeline("soporte_email")
p.agregar_etapa(ConfigEtapa("clasificar", step_clasificar, max_retries=2, critica=False))
p.agregar_etapa(ConfigEtapa("extraer", step_extraer, max_retries=2, critica=False))
p.agregar_etapa(ConfigEtapa("responder", step_responder, max_retries=1, fallback_fn=fallback_responder, critica=True))
state = p.ejecutar(email, verbose=False)
return state.get_output("responder")
email_test = "Llevo 3 días sin poder acceder a mi cuenta. He perdido datos importantes y necesito solución urgente. Soy cliente premium."
print(pipeline_soporte_email(email_test))
Ejercicio 2: Agregar monitoring
Instrumenta el pipeline de análisis de documentos construido en esta cápsula con el PipelineMonitor. Ejecuta 3 documentos diferentes y muestra las estadísticas.
Ver solución
monitor = PipelineMonitor()
pipeline = construir_pipeline_analisis()
documentos = [
"Informe de ventas Q1 2026...",
"Contrato de servicios entre empresa A y empresa B...",
"Artículo sobre tendencias tecnológicas en 2026..."
]
for doc in documentos:
state = pipeline.ejecutar(doc, verbose=False)
monitor.registrar_ejecucion(state, "analisis_documentos")
stats = monitor.estadisticas()
print(f"Ejecuciones: {stats['n_ejecuciones']}")
print(f"Tiempo promedio: {stats['tiempo_promedio']:.2f}s")
print(f"Tasa de éxito: {stats['tasa_exito_promedio']:.0%}")
print(f"Con fallback: {stats['n_ejecuciones_con_fallback']}")
print("\nEtapas más lentas:")
for etapa, tiempo in monitor.etapas_mas_lentas():
print(f" {etapa}: {tiempo:.2f}s promedio")
Resumen
- PipelineState: Objeto centralizado que mantiene el input original, todos los outputs, y métricas de cada etapa
- ConfigEtapa: Configuración declarativa: función, max_retries, fallback, validador, timeout
- Retry inteligente: Backoff exponencial con jitter para evitar thundering herd
- Circuit Breaker: Previene loops infinitos en etapas que sistemáticamente fallan
- Monitoring: Registrar tiempo, intentos, errores por etapa para optimización
- Cache: Para etapas costosas que producen el mismo output para el mismo input