William RodriguezEvitando el GIL en Pipelines de Datos: Ejecución Paralela de DAGs en Wpipe Día 11 de la...
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.
| 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) |
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"))