Evitando el GIL en Pipelines de Datos: Ejecución Paralela de DAGs en Wpipe

# architecture
Evitando el GIL en Pipelines de Datos: Ejecución Paralela de DAGs en WpipeWilliam Rodriguez

Evitando el GIL en Pipelines de Datos: Ejecución Paralela de DAGs en Wpipe Día 11 de la...

Evitando el GIL en Pipelines de Datos: Ejecución Paralela de DAGs en Wpipe

Día 11 de la Serie de Arquitectura Open Source de Wisrovi.

Al ejecutar grafos acíclicos dirigidos (DAGs) complejos en Python, los orquestadores tradicionales suelen estrellarse contra el Global Interpreter Lock (GIL) al correr transformaciones numéricas intensivas, o sufren un costo de memoria astronómico al levantar contenedores pesados para procesar ramas paralelas sencillas.

wpipe resuelve este problema mediante un motor de ejecución híbrido: ejecuta tareas atadas a I/O de forma asíncrona mediante hilos asyncio, mientras delega cálculos pesados a procesos de trabajo con memoria contextual atómica y serialización de alta velocidad.


⚡ Arquitecturas de Ejecución Paralela: SaaS en la Nube vs. Wpipe

Dimensión Arquitectónica Orquestadores Pesados (Airflow/Prefect) Motor Híbrido wpipe
Asignación de Workers Contenedor Docker / Worker Celery independiente Hilo o subproceso ligero en el mismo host
Planificación de Ramas Polling entre servicios vía message broker Resolución topológica directa del DAG
Paso de Datos entre Pasos Serialización de red lenta (S3/XCom/JSON) Contexto atómico en memoria / SQLite WAL
Latencia de Arranque Segundos a minutos < 5 milisegundos
Consumo de Memoria Gigabytes por nodo de trabajo Megabytes (Entorno soberano local)

💻 Implementación Práctica: Ejecución de Ramas Paralelas

El siguiente ejemplo muestra cómo construir un pipeline con ramas concurrentes sin lidiar con bloqueos manuales ni colas multiprocessing complejas:

from wpipe import Pipeline, Step, Context

class ExtraccionFuenteA(Step):
    def run(self, ctx: Context) -> None:
        # Simulación de extracción externa
        ctx.set("datos_a", [x * 2 for x in range(10000)])
        print("Fuente A completada.")

class ExtraccionFuenteB(Step):
    def run(self, ctx: Context) -> None:
        # Simulación de extracción independiente
        ctx.set("datos_b", [x * 3 for x in range(10000)])
        print("Fuente B completada.")

class FusionYCalculo(Step):
    def run(self, ctx: Context) -> None:
        # Depende de que ambas fuentes concluyan exitosamente
        a = ctx.get("datos_a")
        b = ctx.get("datos_b")
        total = sum(a) + sum(b)
        ctx.set("total_global", total)
        print(f"Resultado Agregado: {total}")

# Inicialización del pipeline con capacidad de concurrencia
pipeline = Pipeline("MotorIngestaParalela", max_workers=4)

# Definir pasos: Fuente A y B se ejecutan en paralelo
paso_a = ExtraccionFuenteA()
paso_b = ExtraccionFuenteB()
fusion = FusionYCalculo()

pipeline.add_step(paso_a)
pipeline.add_step(paso_b)
# El paso de fusión se activa automáticamente al completar A y B
pipeline.add_step(fusion, depends_on=[paso_a, paso_b])

resultado = pipeline.execute()
print("Estado de Ejecución:", resultado.status)
print("Salida en Contexto:", resultado.context.get("total_global"))
Enter fullscreen mode Exit fullscreen mode

🛡️ Ventajas de Ingeniería para Producción

  1. Sincronización Determinista de Ramas: Los nodos dependientes nunca se ejecutan antes de tiempo; las condiciones de carrera se eliminan en la compilación del DAG.
  2. Aislamiento de Contexto: Las ramas paralelas mutan su espacio de nombres local sin interferir con pasos hermanos hasta la unión.
  3. Cero Infraestructura Extra: Corre pipelines multirama en servidores locales, hardware embebido o runners efímeros de CI/CD sin requerir un cluster de Kubernetes.

spanish #python #dataengineering #architecture #backend