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:

  1. Clasificar urgencia y tipo
  2. Extraer información clave (cliente, problema, contexto)
  3. 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

Recursos adicionales

  1. LangChain Pipeline documentation
  2. Circuit Breaker pattern
  3. Python shelve for simple caching
  4. Exponential backoff - AWS guide
  5. OpenAI Error handling best practices
  6. Pydantic validators