¿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
- Ejecución Todo o Nada: Un error al final del pipeline descarta todo el cómputo intermedio.
- Fallos por Serialización Pickle: Las referencias circulares o conexiones abiertas rompen el guardado de estado clásico.
- 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({})
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
- GitHub: https://github.com/wisrovi/wpipe
- PyPI: https://pypi.org/project/wpipe/
- Documentación: https://wpipe.readthedocs.io
Autor: William Steve Rodríguez Villamizar (Wisrovi)
Top comments (0)