StackPractices
intermediate Por Mathias Paulenko

Patrón Message Queue Load Leveling

Suavizar picos de trafico colocando una cola entre el productor y el consumidor. El productor escribe mensajes a cualquier ritmo; el consumidor los procesa a un ritmo constante.

Descripción General

Cuando un servicio recibe trafico irregular, puede saturar sistemas descendentes que no estan disenados para picos. Una base de datos puede manejar 50 consultas por segundo de forma constante pero fallar a 500 consultas por segundo en un pico. El patron Message Queue Load Leveling coloca una cola entre el productor y el consumidor para que el productor escriba mensajes a cualquier ritmo mientras el consumidor los procesa a un ritmo controlado y constante.

Cuándo Usar

  • For alternatives, see Queue-Based Load Leveling Pattern.

  • El trafico hacia un sistema descendente es irregular y el sistema no maneja picos

  • Necesitas desacoplar la tasa de peticiones de la tasa de procesamiento

  • Las tareas no son sensibles al tiempo (los usuarios no necesitan respuestas inmediatas)

  • Quieres escalar consumidores independientemente de los productores

  • Trabajos en segundo plano como generacion de reportes, envio de emails o procesamiento de archivos

  • Necesitas entrega confiable — los mensajes persisten en la cola incluso si el consumidor esta temporalmente offline

Cuándo Evitar

  • Requests de usuario en tiempo real. Load leveling anade latencia de cola. Si el usuario espera respuesta, procesa sincrono.
  • Orden estricto entre todos los mensajes. Multiples consumidores rompen el orden. Usa un consumidor unico o el patron Sequential Convoy.
  • Trafico bajo sin picos. Si el trafico es consistentemente bajo, la cola anade complejidad sin beneficio.
  • Los mensajes deben procesarse en una ventana de tiempo especifica. Los delays de cola pueden causar que los mensajes pierdan su deadline.
  • No puedes tolerar procesamiento duplicado. Las colas pueden reentregar. Si la idempotencia es imposible, usa una arquitectura diferente.

Solución

Python (Celery + Redis)

from celery import Celery
import time

app = Celery("tasks", broker="redis://localhost:6379", backend="redis://localhost:6379")

# El consumidor procesa una tarea a la vez a su propio ritmo
@app.task(bind=True, max_retries=3)
def process_order(self, order_id):
    try:
        # Simular procesamiento lento (escrituras DB, llamadas API)
        time.sleep(2)
        print(f"Processed order {order_id}")
        return {"status": "done", "order_id": order_id}
    except Exception as exc:
        raise self.retry(exc=exc, countdown=5)

# El productor encola a cualquier ritmo
def submit_orders(order_ids):
    for order_id in order_ids:
        process_order.delay(order_id)
    print(f"Enqueued {len(order_ids)} orders")

# Pico: 1000 ordenes enviadas al instante
# El consumidor las procesa 1 a la vez cada 2 segundos
submit_orders(range(1000))

JavaScript (BullMQ + Redis)

import { Queue, Worker } from "bullmq";

const orderQueue = new Queue("orders", {
  connection: { host: "localhost", port: 6379 },
});

// Productor: encolar a cualquier ritmo
async function submitOrders(orderIds) {
  const jobs = orderIds.map((id) => ({
    name: "process-order",
    data: { orderId: id },
  }));
  await orderQueue.addBulk(jobs);
  console.log(`Enqueued ${orderIds.length} orders`);
}

// Consumidor: procesar a ritmo controlado
const worker = new Worker(
  "orders",
  async (job) => {
    // Simular procesamiento lento
    await new Promise((resolve) => setTimeout(resolve, 2000));
    console.log(`Processed order ${job.data.orderId}`);
    return { status: "done", orderId: job.data.orderId };
  },
  {
    connection: { host: "localhost", port: 6379 },
    concurrency: 1, // Procesar una a la vez
    limiter: { max: 1, duration: 2000 }, // Max 1 job cada 2 segundos
  }
);

// Pico: 1000 ordenes enviadas al instante
await submitOrders(Array.from({ length: 1000 }, (_, i) => i));

Java (RabbitMQ + Spring AMQP)

import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

@Component
public class OrderProcessor {

    private final RabbitTemplate rabbitTemplate;

    public OrderProcessor(RabbitTemplate rabbitTemplate) {
        this.rabbitTemplate = rabbitTemplate;
    }

    // Productor: enviar a cualquier ritmo
    public void submitOrders(List<Integer> orderIds) {
        for (Integer orderId : orderIds) {
            rabbitTemplate.convertAndSend("orders", "order." + orderId, orderId);
        }
        System.out.println("Enqueued " + orderIds.size() + " orders");
    }

    // Consumidor: procesar una a la vez
    @RabbitListener(queues = "orders", concurrency = "1")
    public void processOrder(Integer orderId) {
        try {
            Thread.sleep(2000); // Simular procesamiento lento
            System.out.println("Processed order " + orderId);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

Explicación

La cola actua como buffer entre productor y consumidor. El productor empuja mensajes a la cola tan rapido como puede. El consumidor extrae mensajes a un ritmo que puede manejar. Si el productor envia 1000 mensajes en un segundo pero el consumidor procesa 1 cada 2 segundos, la cola crece a 1000 mensajes y se vacia lentamente en 2000 segundos.

Esto protege al sistema descendente de saturarse. La contrapartida es latencia: los mensajes esperan en la cola hasta que el consumidor puede procesarlos. Para cargas sensibles al tiempo, aumenta la concurrencia del consumidor o usa el patron Priority Queue.

Variantes

VarianteTipo de ColaCaso de UsoCompromiso
Consumidor UnicoCola FIFOOrden estricto, simpleThroughput lento
Multiples ConsumidoresCola FIFOMayor throughputSin garantia de orden
Cola de PrioridadCola prioritariaAlgunos mensajes son urgentesComplejidad en logica de prioridad
Retardo ProgramadoCola con delayProcesar en momentos especificosLos mensajes esperan hasta el horario
Procesamiento por LotesConsumidor batchAgrupar mensajes para eficienciaMayor latencia por mensaje

Qué Funciona

  • Dimensiona la cola segun el volumen de pico esperado y la tasa de procesamiento del consumidor
  • Monitorea la profundidad de cola y alerta cuando crece mas alla de un umbral
  • Escala consumidores horizontalmente cuando la profundidad es consistentemente alta
  • Usa dead-letter queues para mensajes que fallan tras maximos reintentos
  • Establece visibility timeouts para evitar doble procesamiento si un consumidor cae
  • Usa consumidores idempotentes para manejar entregas duplicadas de forma segura

Errores Comunes

  • Sin monitoreo de profundidad de cola: Una cola creciente significa que los consumidores no dan abasto. Sin monitoreo, te enteras cuando la cola se queda sin almacenamiento.
  • Consumidor muy lento para trafico sostenido: Si la tasa promedio de produccion excede la de consumo, la cola crece infinitamente.
  • No manejar mensajes venenosos: Un mensaje que siempre falla bloquea al consumidor.
  • Productor sincrono esperando al consumidor: Derrota el proposito. El productor debe fire-and-forget.
  • Ignorar el orden de mensajes: Si el orden importa, se necesita un consumidor unico o estrategia de particion. Multiples consumidores rompen el orden.

Como Funciona

  1. El productor escribe a la cola: El productor envia mensajes a la cola sin esperar al consumidor. La cola acusa recibo inmediatamente.
  2. La cola bufferiza mensajes: Los mensajes persisten en la cola hasta que un consumidor esta disponible. La cola garantiza entrega incluso si el consumidor esta offline.
  3. El consumidor tira a su ritmo: El consumidor lee mensajes uno a uno (o en lotes) a una tasa que puede manejar. El tiempo de procesamiento por mensaje determina el throughput efectivo.
  4. Acknowledgment cierra el ciclo: Tras procesar, el consumidor acusa recibo del mensaje. Si el consumidor cae antes de acusar, el broker reentrega el mensaje a otro consumidor.

El insight clave es desacoplamiento de tasas: la tasa del productor y la del consumidor son independientes. La cola absorbe la diferencia durante picos.

Mejores Practicas

  • Establece alerta de profundidad maxima. Cuando la profundidad excede 80% de capacidad, dispara una alerta. Te da tiempo para escalar consumidores antes de que la cola se llene.
  • Usa backoff exponencial para reintentos. Si un mensaje falla, reintenta con delays crecientes (1s, 2s, 4s, 8s). Previene que tormentas de reintentos saturen al consumidor.
  • Separa colas por prioridad. Usa la variante Priority Queue para mensajes urgentes. Una cola unica trata todos los mensajes igual.
  • Consumidores idempotentes. Los brokers pueden reentregar mensajes. Disena consumidores para que procesar el mismo mensaje dos veces produzca el mismo resultado.
  • Dimensiona consumidores para carga sostenida, no pico. Si el pico es 10x el promedio, dimensionar para pico desperdicia recursos. Dimensiona para promedio + 20% de margen y deja que la cola absorba picos.

Ejemplos del Mundo Real

Amazon SQS + Lambda

Una plataforma de e-commerce usa SQS para bufferizar mensajes de pedidos durante Black Friday. Los pedidos llegan a 50,000/s pero el backend de procesamiento de pagos maneja 5,000/s. SQS bufferiza el pico. Funciones Lambda consumen a 5,000/s con concurrencia controlada. La cola se vacia en 10 segundos despues del burst.

RabbitMQ en Sistemas Financieros

Una plataforma de trading recibe bursts de datos de mercado al abrir. Las colas de RabbitMQ bufferizan el burst mientras el servicio de analitica procesa a ritmo constante. Sin load leveling, el servicio de analitica caeria bajo el burst de apertura.

Azure Service Bus para IoT

Una plataforma IoT recolecta telemetria de millones de dispositivos. Los mensajes llegan en bursts cuando los dispositivos se reconectan tras cortes de red. Las colas de Service Bus bufferizan los bursts mientras los servicios backend procesan a tasa controlada, previniendo sobrecarga de base de datos.

Lectura Adicional

  • Documentación oficial: consulta la referencia actualizada del framework o herramienta utilizada.
  • Guías relacionadas: explora las guías de load-balancing y pattern para profundizar.
  • Patrones complementarios: revisa los patrones de diseño aplicables a tu stack tecnológico.
  • Postmortems públicos: estudia incidentes reales de equipos que enfrentaron problemas similares en producción.

Notas de Producción

  • Despliega gradualmente usando canary o blue-green para detectar regresiones temprano.
  • Configura alertas para errores, latencia p99 y tasa de fallos antes de habilitar en producción.
  • Documenta el rollback en el runbook; prueba el procedimiento en staging al menos una vez por trimestre.
  • Revisa logs estructurados con correlation IDs para trazar requests end-to-end en incidentes.

Puntos Clave

  • Aplica patrón message queue load leveling cuando necesites una solución práctica para tu caso de uso.
  • Monitorea el rendimiento después de implementar; mide latencia, errores y uso de recursos antes y después.
  • Revisa la sección de Troubleshooting ante errores comunes; la mayoría tienen causa raíz documentada con solución.
  • Mantén dependencias actualizadas y ejecuta tests en CI para prevenir regresiones en producción.

Troubleshooting

  • Messages are lost on restart: persist messages before acknowledging.
  • Consumer lags behind producer: scale consumers, increase prefetch, and partition the topic.
  • Duplicate messages: design consumers to be idempotent.
  • Ordering is wrong after scaling: preserve partition keys and avoid rebalancing during bursts. Consider a single partition when order is mandatory.
  • Queue depth grows but consumers are idle: check network partitions, consumer health, and permission issues. Restart gracefully.

Errores Comunes en Producción

  • Aplicar el patrón donde no se necesita abstracción, agregando complejidad accidental.
  • Dejar que el patrón se filtre en módulos no relacionados y confundir los límites de responsabilidad.
  • Sobre-ingeniería en la primera implementación en lugar de comenzar simple y medir el dolor.
  • Saltar los tests de contrato, de modo que las refactorizaciones rompan consumidores en silencio.
  • Ignorar modos de fallo que el patrón no cubre.
  • Usar el patrón como opción por defecto en lugar de elegir la herramienta adecuada para la escala actual.
  • Olvidar documentar cuándo dejar de usar el patrón y qué lo reemplaza.
  • Carecer de observabilidad sobre rendimiento y propagación de errores del patrón.