Camino de yuwen-c

Cómo acabé entendiendo pipeline, queue y worker

#backend #architecture #queue notes

Así empezó todo

Me hice cargo de un sistema que venía de un compañero. Él decía que era un «sistema de colas», pero a medida que lo explicaba, acabó convirtiéndose en pipeline.

La arquitectura de verdad era bastante simple: un servicio Flask que recibe la API, más 3 workers; cada uno se encarga, en orden, de una etapa del procesamiento de datos.

Entonces, ¿una queue (cola) es lo mismo que un pipeline? ¿Y qué relación hay entre queue y worker? Así que me puse a investigar…

Dejar claros estos términos primero

¿Un pipeline es una queue? No. En realidad son términos de distinto nivel.

Una queue se parece más a una lista de pendientes: llega una tarea, se añade a la lista y se procesa por orden.

Por comparación, hay otra forma de tratar las tareas: las que llegan después se procesan primero.

Un pipeline, en cambio, es el flujo de procesamiento en sí: los datos tienen que pasar, en orden, por una serie de pasos para considerarse terminados.

flowchart TB
    root["¿Un pipeline es<br> una queue?"]
    root --> Q["queue<br>FIFO — primero en entrar, primero en salir<br>Hacer fila; lo que llega antes se procesa primero"]
    root --> S["stack<br>LIFO — último en entrar, primero en salir<br>También llamada stack-based queue"]
    root --> P["pipeline<br>Un flujo de procesamiento ordenado<br>No choca con la idea de queue"]

    style root fill:#dce9f5,stroke:#5d8aa8

Entre recibir y ejecutar: el broker

En un sistema de colas, la queue se encarga de guardar las tareas, el worker de ejecutarlas, y en medio hay un broker que se ocupa del paso de mensajes.

Si tomamos Celery, el worker framework habitual en Python, la combinación típica es usar Redis como «broker». Antes, cuando trabajaba en comercio internacional, veía mucho esta palabra: significaba «intermediario» o «agencia de aduanas». En software, en cambio, suele llamarse «intermediario de mensajes».

El broker incluye la queue de la que hablaba antes: esa estructura donde dejar la lista de pendientes. Cuando nace una tarea, hay un sitio donde guardarla, y el worker la recoge de ahí para ejecutarla.

El conjunto queda más o menos así:

flowchart TB
    A["Programa principal<br>(p. ej. API Flask)"] -->|"llamar task.delay()"| C1["Celery client/app<br>Convierte la llamada a función en<br>un task message y lo envía al broker"]

    subgraph B["broker (intermediario de mensajes)"]
        Q["queue (contenedor FIFO)<br>Guarda las tareas pendientes"]
    end

    C1 -->|enviar task message| Q
    Q -->|sacar tarea| C2["Celery worker<br>Saca tareas del broker<br>y las ejecuta"]
    C2 --> T["Ejecutar la tarea"]

    classDef celery fill:#f9e79f,stroke:#b7950b
    class C1,C2 celery

Redis no es realmente un broker formal

Mientras investigaba me encontré con esta frase: «Redis, en sentido estricto, no cumple la definición de broker». Entonces… ¿por qué lo usa todo el mundo?

Más adelante, investigando LLM resumable streaming, vi que el autor de este artículo también usa Redis — y encima el patrón Redis Pub/Sub. Al tirar del hilo, justamente respondía a la duda de antes:

Cuando Redis salió, no era más que una estructura de datos list en memoria: muy cómoda de usar, y encajaba justo con lo que necesita una queue (un sitio temporal para las tareas), así que se popularizó. Pero en un entorno de producción más exigente, donde las tareas no se pueden perder, Redis ya no encaja tan bien.

Un broker que cumpla la definición estricta tiene que poder «gestionar las tareas». Más adelante sí apareció Redis Stream para usarse como message broker.

flowchart TB
    DEF["Definición estricta de un broker"]
    DEF --> A1["ACK"]
    DEF --> A2["Reintentos ante error"]
    DEF --> A3["delivery guarantee<br>at-least-once /<br>exactly-once"]

    A1 -.-> R
    A2 -.->|"Usar lo anterior para juzgar Redis"| R
    A3 -.-> R
    R["Redis:<br>No es un broker «formal».<br>Evolución histórica:"]

    R --> L["Redis List<br>Diseño inicial; cómodo, a menudo usado como queue.<br>Sin ACK."]
    R --> ST["Redis Stream (posterior a List)<br>Soporta ACK<br>Pensado para message broker."]
    R --> PS["Redis Pub/Sub<br>Sin ACK<br>Difusión en tiempo real fire-and-forget."]

    style DEF fill:#dce9f5,stroke:#5d8aa8
    style ST fill:#d5e8d4,stroke:#82b366
    style L fill:#f5f5f5,stroke:#999999
    style PS fill:#f5f5f5,stroke:#999999

Parece que worker y queue siempre van juntos… ¿o no?

Volviendo al principio: el sistema de colas que disparó toda esta investigación sí parece la combinación clásica, queue más worker. Entonces, ¿cada vez que hay que procesar tareas, encajo ya esta arquitectura y listo?

No necesariamente. Se pueden pensar por separado: según el escenario usas una tecnología u otra, y no hace falta atarlas siempre.

He dibujado cuándo hace falta una queue, cuándo conviene añadir un worker y qué lo decide, para ayudarme a construir el modelo mental —

En pocas palabras: el volumen de las tareas decide si hace falta una queue; lo pesadas que son, si hace falta un worker.

Solo queue

Al abrir plazas: API responde «recibido»; sign-ups van a una cola para escribirse después.

queue + worker

El usuario sube una imagen → redimensionar / convertir

Procesar en el proceso principal

Exportación manual de pocos datos en el panel (decenas de filas)

Worker en segundo plano

Informes por lotes nocturnos (hora fija, volumen previsible)

Diagrama de decisión queue / worker: el eje vertical es el volumen e incertidumbre de las tareas (arriba: entrada a ráfagas, hace falta una queue; abajo: origen fijo); el eje horizontal es la carga de ejecución (izquierda: termina rápido; derecha: bloquearía el proceso principal, hace falta un worker). Cuadrantes: ① proceso principal, ② worker en segundo plano, ③ solo queue, ④ queue más worker.

El diagrama es ancho: desliza horizontalmente para verlo; pellizca para ampliar.