DEV Community

William Rodriguez
William Rodriguez

Posted on

Pipelines tolerantes a fallos en Python: Reanudación con Checkpoints SQLite

¿Qué ocurre cuando tu pipeline de datos en Python falla en el paso 19 de 20 por un timeout de red o una desconexión inesperada? En scripts lineales sin persistencia, debes reiniciar desde cero, duplicando costes de cómputo y generando efectos secundarios no deseados.

WPipe v2.2.0 introduce Checkpoints Inteligentes sobre SQLite (Modo WAL) para garantizar cero pérdida de datos y reanudación automática paso a paso.

Este es el Día 03 de la serie técnica WPipe Open Source (MIT, Python 3.9–3.14).


Los Problemas Reales en Pipelines Convencionales

  1. Ejecución Todo o Nada: Un error al final del pipeline descarta todo el cómputo intermedio.
  2. Fallos por Serialización Pickle: Las referencias circulares o conexiones abiertas rompen el guardado de estado clásico.
  3. Infraestructuras Pesadas: Desplegar orquestadores monolíticos complejos solo para tener puntos de control añade fricción innecesaria.

Implementación en Producción: Recuperación con CheckpointManager

from wpipe import Pipeline, step, CheckpointManager

# 1. Inicializar el pipeline con persistencia de estado
pipeline = Pipeline(pipeline_name="etl_resiliente")

# 2. Definir checkpoint basado en expresión lógica
pipeline.add_checkpoint(
    checkpoint_name="datos_cargados",
    expression="len(registros) > 0"
)

@step(name="extraer_origen")
def extraer_origen(data):
    # Simula ingesta de registros
    return {"registros": [101, 102, 103], "estado": "cargado"}

@step(name="procesamiento_pesado")
def procesamiento_pesado(data):
    # Cómputo intensivo de CPU
    procesados = [r * 2 for r in data["registros"]]
    return {"procesados": procesados}

pipeline.set_steps([extraer_origen, procesamiento_pesado])

# 3. Comprobación y reanudación automática
chk = CheckpointManager("estado_pipeline.db")
if chk.can_resume("etl_resiliente"):
    print("¡Interrupción detectada! Reanudando desde el último checkpoint verificado...")
    result = pipeline.resume()
else:
    result = pipeline.run({})
Enter fullscreen mode Exit fullscreen mode

Por qué WPipe marca la diferencia

  • Motor SQLite en Modo WAL: Escribe checkpoints en microsegundos sin bloquear lecturas concurrentes multihilo.
  • Serialización Resiliente: Tolera grafos de memoria complejos y filtra objetos no persistibles sin colapsar.
  • Puro Python sin Servidores: Sin daemons externos ni dependencias pesadas; embebido directamente en cualquier worker.

Instalación y Recursos

pip install wpipe
Enter fullscreen mode Exit fullscreen mode

Autor: William Steve Rodríguez Villamizar (Wisrovi)

Top comments (0)