Pipeline robusto: 500 llamadas exitosas mediante excepciones, reintentos y circuit breakers
Al finalizar este tema
Podrán construir las herramientas necesarias para que un pipeline que realiza llamadas masivas a una API externa no colapse ante algunos fallos, sino que complete su ejecución. Se trata de patrones prácticos que elevan la tasa de éxito del 95% al 99.9%.
Este texto es un ejemplo educativo general. En entornos de producción se utilizan bibliotecas como tenacity, backoff o resilience4j.
"A la mañana siguiente, abrí los registros": las trampas de un pipeline ingenuo
Imaginen que ejecutaron durante toda la noche un script que llama secuencialmente a NCBI BLAST para 500 genes.
def naive_pipeline(gene_ids: list[str]) -> list[dict]: results = [] for gene_id in gene_ids: response = call_blast_api(gene_id) results.append(response) return resultsAl abrir el registro de al día siguiente, se ve así.
[00:03] gene 1: OK
[00:05] gene 2: OK
...
[01:23] gene 27: TIMEOUT
Traceback (most recent call last):
...
requests.exceptions.TimeoutEl script completo se detuvo en la vigésima séptima llamada, y los restantes 473 intentos ni siquiera se ejecutaron. Un único fallo ha anulado toda la ejecución.
Los problemas de este enfoque:
Problema 1: Falta de aislamiento de fallos. Un solo error detiene todo el proceso.
Problema 2: Sin reintentos. Las fallas de red suelen ser temporales; esperar un momento y reintentar aumenta la probabilidad de éxito.
Problema 3: Sin guardado del progreso. Los resultados exitosos de las 26 llamadas se perdieron en la memoria al producirse el bloqueo, obligando a reiniciar desde cero.
El enfoque correcto es un patrón de canalización robusto. Aislar los fallos de cada llamada, implementar reintentos y abrir el circuito (cortocircuito) ante fallos persistentes. Guardar inmediatamente los resultados exitosos.
De la caja negra a los componentes
Componente 1: Aislamiento de fallos mediante try-except
def robust_pipeline_basic(gene_ids: list[str]) -> tuple[list[dict], list[dict]]: results = [] errors = []
for gene_id in gene_ids: try: response = call_blast_api(gene_id) results.append({"gene_id": gene_id, "data": response}) except Exception as e: errors.append({"gene_id": gene_id, "error": type(e).__name__, "message": str(e)})
return results, errorsClave: Las llamadas fallidas se aíslan mediante errors, y la siguiente iteración continúa. Incluso si algunas de las 500 fallan, el resto completa el recorrido.
Precaución: except Exception captura todas las excepciones. Úsalo únicamente en la parte superior del pipeline. Si deseas capturar solo errores específicos, hazlo explícitamente.
try: response = call_blast_api(gene_id)except requests.exceptions.Timeout: # Errores que permiten reintento passexcept requests.exceptions.HTTPError as e: if e.response.status_code == 429: # límite de tasa → se puede reintentar pass else: # otros 4xx → el reintento no tiene sentido raiseComponente 2: Reintentos con retroceso exponencial
La mayoría de los errores temporales se resuelven mediante reintentos. Sin embargo, los reintentos inmediatos sobrecargan al servidor. Aumenta exponencialmente el intervalo entre cada reintento.
import timeimport random
def with_retry(func, max_attempts: int = 5, base_delay: float = 1.0): """ Llama a func, pero reintenta con backoff exponencial si falla. """ for attempt in range(1, max_attempts + 1): try: return func() except (requests.exceptions.Timeout, requests.exceptions.ConnectionError) as e: if attempt == max_attempts: raise
delay = base_delay * (2 ** (attempt - 1)) jitter = random.uniform(0, delay * 0.1) time.sleep(delay + jitter)Intervalos: 1 segundo, 2 segundos, 4 segundos, 8 segundos, 16 segundos. Hasta 31 segundos de espera con un máximo de 5 reintentos.
Razón para añadir Jitter (error aleatorio): Si varios clientes fallan simultáneamente y reintentan exactamente al mismo tiempo, el servidor vuelve a caerse. Se añade un error aleatorio para que cada uno reintente en momentos ligeramente distintos.
Componente 3: Circuit Breaker (Interruptor de circuito)
Si el servidor permanece inactivo durante mucho tiempo, los reintentos carecen de sentido. Cuando las fallas consecutivas superan un umbral, se detiene temporalmente por completo la invocación.
from dataclasses import dataclassfrom enum import Enum
class CircuitState(Enum): CLOSED = "closed" # Normal - permite llamadas OPEN = "open" # Bloqueado - falla inmediatamente HALF_OPEN = "half_open" # Verificar recuperación - llamada de prueba
@dataclassclass CircuitBreaker: failure_threshold: int = 5 recovery_timeout: float = 60.0 state: CircuitState = CircuitState.CLOSED failure_count: int = 0 last_failure_time: float = 0.0
def call(self, func): now = time.time()
if self.state == CircuitState.OPEN: if now - self.last_failure_time > self.recovery_timeout: self.state = CircuitState.HALF_OPEN else: raise Exception("Circuit breaker OPEN")
try: result = func() if self.state == CircuitState.HALF_OPEN: self.state = CircuitState.CLOSED self.failure_count = 0 return result except Exception: self.failure_count += 1 self.last_failure_time = now if self.failure_count >= self.failure_threshold: self.state = CircuitState.OPEN raiseTres estados:
- CLOSED: Normal. Llamadas permitidas.
- OPEN: Bloqueado. Devuelve fallo inmediato (sin realizar llamadas reales).
- HALF_OPEN: Intento de recuperación. Prueba el estado del servidor con una sola llamada.
Ensamblaje del pipeline
Combina los tres componentes.
def robust_pipeline_full( gene_ids: list[str], output_path: str, checkpoint_every: int = 50) -> dict: breaker = CircuitBreaker(failure_threshold=5, recovery_timeout=60.0) results = [] errors = []
# Cargar punto de control existente already_done = load_checkpoint(output_path)
for i, gene_id in enumerate(gene_ids): if gene_id in already_done: continue
try: def call(): return breaker.call(lambda: call_blast_api(gene_id))
response = with_retry(call, max_attempts=3, base_delay=2.0) results.append({"gene_id": gene_id, "data": response}) except Exception as e: errors.append({ "gene_id": gene_id, "error": type(e).__name__, "message": str(e) })
if (i + 1) % checkpoint_every == 0: save_checkpoint(output_path, results) print(f"Progress: {i+1}/{len(gene_ids)} — checkpoint saved")
save_checkpoint(output_path, results) return { "total": len(gene_ids), "success": len(results), "failed": len(errors), "errors": errors }
def save_checkpoint(path: str, results: list) -> None: import json with open(path, "w") as f: json.dump(results, f, indent=2)
def load_checkpoint(path: str) -> set[str]: import json from pathlib import Path if not Path(path).exists(): return set() with open(path) as f: data = json.load(f) return {r["gene_id"] for r in data}Ahora, incluso si hay un fallo, al volver a ejecutar se omiten los genes ya procesados y solo se intentan los restantes.
Desvanecimiento — Los tres espacios en blanco que deben completar
Espacio en blanco 1: Política de reintento por tipo de excepción
Solo se deben reintentar ciertas excepciones. Por ejemplo, un error 400 Bad Request no tiene sentido volver a intentarlo.
def with_retry_smart(func, max_attempts: int = 5): retryable = ( requests.exceptions.Timeout, requests.exceptions.ConnectionError, )
for attempt in range(1, max_attempts + 1): try: return func() except retryable as e: # TODO: Lógica de reintento (backoff + jitter) pass except requests.exceptions.HTTPError as e: # TODO: Verificar status_code y reintentar si es 429 o 5xx # De lo contrario, lanzar inmediatamente passSugerencia: if e.response.status_code == 429 or 500 <= e.response.status_code < 600: sleep(delay); continue; else: raise.
Cuadro de informe de fallo
Analiza los gene_ids fallidos y clasifícalos según la causa.
def analyze_failures(errors: list[dict]) -> dict: """ Agrupar la lista de errores por tipo. Devolución: {"Timeout": 5, "HTTPError_500": 3, ...} """ from collections import Counter
# TODO: contar por el campo "error" de cada error # HTTPError se subdivide por status_code passSugerencia: counter = Counter(); for e in errors: key = e["error"]; if "status" in e["message"]: key = f"{key}_{parse_status}"; counter[key] += 1.
Espacio en blanco 3: Reintentar solo la parte fallida al reejecutar
Guardar tanto el éxito como el fallo en el punto de control y, al reejecutar, intentar solo lo que falló.
def retry_failed_only( gene_ids: list[str], output_path: str, error_log_path: str) -> None: # TODO 1: cargar error_log_path # TODO 2: extraer solo los gene_ids fallidos # TODO 3: reejecutar robust_pipeline_full (con la lista filtrada) passReflexión: diferencias con herramientas de robustez en la práctica
Biblioteca tenacity: herramienta estándar de Python para reintentos. Sintaxis mucho más limpia gracias al uso de decoradores.
from tenacity import retry, stop_after_attempt, wait_exponential
@retry(stop=stop_after_attempt(5), wait=wait_exponential(multiplier=1, max=60))def call_api(gene_id): return requests.get(f"...").json()Sondas de preparación de Kubernetes: En el despliegue real, la propia aplicación informa de su estado. Al detectar problemas, se evita el tráfico.
Netflix Hystrix / resilience4j: Interruptores de circuito sofisticados del ecosistema Java. Combinación de patrones como Bulkhead, TimeLimiter y RateLimiter.
Cola de mensajes muertos (Dead Letter Queue): Guarda las solicitudes fallidas en una cola separada para su reprocesamiento manual o automático posterior. Función estándar de AWS SQS y RabbitMQ.
Alternativas a transacciones distribuidas: En la práctica, se diseñan con métodos recuperables ante fallos (patrón saga, event sourcing). El fallo se trata como un caso normal, no como una excepción.
Proyecto de expansión
1. Reescritura con tenacity: Sustituir vuestra lógica de backoff por el decorador tenacity.
2. Procesamiento en paralelo: concurrent.futures.ThreadPoolExecutor para realizar varias llamadas simultáneamente. El interruptor de circuito se comparte entre varios hilos.
3. Panel de control: Mostrar en tiempo real el progreso, la tasa de éxito y el estado del interruptor de circuito. Con Streamlit, en 30 minutos.
4. Notificaciones: Si la tasa de fallos supera el umbral, enviar notificaciones por Slack/Telegram.
Mapa de componentes de este capítulo
- [F] Excepciones: El sistema jerárquico de excepciones de Python. Herencia de
Exception, capturar solo tipos específicos. - [F] try-except: Aislamiento y recuperación del fallo. Acceso a la información del error con
except ... as e. - [F] Manejo de errores: Patrón integrado de reintento, backoff e interruptor de circuito.
- [W] E/S de archivos: Guardado y carga de puntos de control (script completo proporcionado).
[F] = Implementación directa por parte del lector / [W] = Código completo proporcionado.