StackPractices
intermediate Por Mathias Paulenko

Patrón de Cola de Prioridad: Programación por Urgencia

Usa el patrón de cola de prioridad para procesar primero el trabajo crítico. Ejemplos en Python, Java y JavaScript con heaps y Redis.

Visión General

El Patrón de Cola de Prioridad organiza tareas o mensajes de modo que los elementos de mayor prioridad se procesen antes que los de menor prioridad, independientemente del orden de llegada. En lugar de la cola tradicional FIFO donde las tareas se manejan en orden de envío, una cola de prioridad ordena las tareas por importancia, urgencia o valor de negocio.

Este patrón es esencial cuando los recursos son limitados y no todas las tareas pueden procesarse inmediatamente. Garantiza que las operaciones críticas (detección de fraude, solicitudes de clientes VIP, alertas del sistema) reciban atención inmediata mientras el trabajo de fondo rutinario espera.

Cuándo Usar

  • Capacidad de procesamiento limitada con importancia heterogénea de tareas
  • Experiencias de clientes VIP o por niveles donde los usuarios premium obtienen servicio más rápido
  • Sistemas de respuesta a incidentes donde las alertas críticas deben preceder a las advertencias
  • Programación de trabajos donde los plazos o SLAs determinan el orden de ejecución
  • Procesadores de tareas de fondo con cargas mixtas (email, reportes, exportaciones)
  • Sistemas multi-tenant donde los inquilinos de mayor pago obtienen prioridad

Para estrategias relacionadas, consulta el patrón Queue-Based Load Leveling y el patrón Throttling.

Cuándo Evitar

  • Todas las tareas tienen igual importancia: una cola FIFO regular es más simple y justa
  • El hambre de tareas de baja prioridad es inaceptable: considerar envejecimiento o scheduling justo
  • El costo de determinar prioridad excede el costo de procesar la tarea
  • El orden FIFO estricto es un requisito de negocio
  • Volúmenes muy pequeños donde el ordenamiento no aporta beneficio

Solución

flowchart diagram: Llega tarea

Python (Cola de Prioridad basada en Heap)

import heapq
import time
from dataclasses import dataclass, field
from typing import List, Callable
from enum import Enum
import threading

class Priority(Enum):
    CRITICAL = 1
    HIGH = 2
    NORMAL = 3
    LOW = 4
    BACKGROUND = 5

@dataclass(order=True)
class Task:
    priority: int
    timestamp: float = field(compare=True)
    task_id: str = field(compare=False)
    payload: dict = field(compare=False)
    handler: Callable = field(compare=False, default=None)

class PriorityQueueProcessor:
    """Procesa tareas por prioridad con equidad dentro de cada nivel"""

    def __init__(self, num_workers=4):
        self.heap = []
        self.lock = threading.Lock()
        self.workers = []
        self.running = False
        self.num_workers = num_workers

    def submit(self, task_id: str, payload: dict,
               priority: Priority = Priority.NORMAL,
               handler: Callable = None):
        """Envía una tarea con una prioridad dada"""
        task = Task(
            priority=priority.value,
            timestamp=time.time(),
            task_id=task_id,
            payload=payload,
            handler=handler
        )
        with self.lock:
            heapq.heappush(self.heap, task)

    def _process_next(self):
        """Worker: saca la tarea de mayor prioridad y la procesa"""
        with self.lock:
            if not self.heap:
                return None
            task = heapq.heappop(self.heap)

        try:
            if task.handler:
                task.handler(task.payload)
            else:
                self._default_handler(task)
            print(f"Completado: {task.task_id} (prioridad {task.priority})")
        except Exception as e:
            print(f"Falló {task.task_id}: {e}")

    def _default_handler(self, task: Task):
        """Lógica de procesamiento default"""
        print(f"Procesando {task.task_id}: {task.payload}")
        time.sleep(0.1)  # Simular trabajo

    def _worker_loop(self):
        while self.running:
            self._process_next()
            time.sleep(0.01)

    def start(self):
        self.running = True
        for _ in range(self.num_workers):
            t = threading.Thread(target=self._worker_loop, daemon=True)
            t.start()
            self.workers.append(t)

    def stop(self):
        self.running = False
        for w in self.workers:
            w.join(timeout=2)

# Uso
processor = PriorityQueueProcessor(num_workers=2)
processor.start()

processor.submit("email-batch", {"type": "newsletter"}, Priority.LOW)
processor.submit("fraud-alert", {"user_id": 12345, "risk_score": 0.95}, Priority.CRITICAL)
processor.submit("report-gen", {"format": "pdf"}, Priority.NORMAL)
processor.submit("vip-onboarding", {"customer_id": "VIP-001"}, Priority.HIGH)

time.sleep(2)
processor.stop()

Java (PriorityBlockingQueue con Thread Pool)

import java.util.Comparator;
import java.util.concurrent.*;

public class PriorityQueueScheduler {

    private final PriorityBlockingQueue<PriorityTask> queue;
    private final ExecutorService executor;

    public PriorityQueueScheduler(int numWorkers) {
        // Comparator: menor valor de prioridad = mayor prioridad, luego FIFO dentro de la misma
        this.queue = new PriorityBlockingQueue<>(1000, Comparator
            .comparingInt(PriorityTask::getPriority)
            .thenComparingLong(PriorityTask::getTimestamp));

        this.executor = Executors.newFixedThreadPool(numWorkers);
        startWorkers(numWorkers);
    }

    private void startWorkers(int numWorkers) {
        for (int i = 0; i < numWorkers; i++) {
            executor.submit(this::workerLoop);
        }
    }

    private void workerLoop() {
        while (!Thread.currentThread().isInterrupted()) {
            try {
                PriorityTask task = queue.take();
                processTask(task);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                break;
            }
        }
    }

    private void processTask(PriorityTask task) {
        System.out.printf("Procesando [%s] prioridad=%d: %s%n",
            task.getTaskId(), task.getPriority(), task.getPayload());

        try {
            task.getHandler().run();
        } catch (Exception e) {
            System.err.println("Task falló: " + task.getTaskId() + " - " + e.getMessage());
        }
    }

    public void submit(String taskId, Runnable handler, Priority priority) {
        queue.offer(new PriorityTask(taskId, priority.value, handler));
    }

    public void shutdown() {
        executor.shutdown();
    }

    enum Priority {
        CRITICAL(1), HIGH(2), NORMAL(3), LOW(4), BACKGROUND(5);
        final int value;
        Priority(int value) { this.value = value; }
    }

    static class PriorityTask {
        private final String taskId;
        private final int priority;
        private final Runnable handler;
        private final long timestamp = System.currentTimeMillis();

        PriorityTask(String taskId, int priority, Runnable handler) {
            this.taskId = taskId;
            this.priority = priority;
            this.handler = handler;
        }

        public int getPriority() { return priority; }
        public long getTimestamp() { return timestamp; }
        public String getTaskId() { return taskId; }
        public Runnable getHandler() { return handler; }
    }

    public static void main(String[] args) {
        PriorityQueueScheduler scheduler = new PriorityQueueScheduler(2);

        scheduler.submit("report-gen", () -> System.out.println("Generando reporte..."), Priority.NORMAL);
        scheduler.submit("fraud-check", () -> System.out.println("Chequeando fraude..."), Priority.CRITICAL);
        scheduler.submit("data-cleanup", () -> System.out.println("Limpiando..."), Priority.BACKGROUND);
        scheduler.submit("vip-request", () -> System.out.println("VIP request..."), Priority.HIGH);

        try {
            Thread.sleep(2000);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
        scheduler.shutdown();
    }
}

JavaScript (Cola de Prioridad con Redis Sorted Set)

const Redis = require('ioredis');

class RedisPriorityQueue {
    constructor(redis, queueName) {
        this.redis = redis;
        this.queueName = queueName;
    }

    async enqueue(task, priority = 3) {
        // Menor score = mayor prioridad; timestamp rompe empates para FIFO dentro de prioridad
        const score = priority * 1000000000 + Date.now();
        const taskJson = JSON.stringify(task);

        await this.redis.zadd(this.queueName, score, taskJson);
    }

    async dequeue() {
        // Saca y elimina el item de menor score
        const result = await this.redis.zpopmin(this.queueName, 1);
        if (result.length === 0) return null;

        const [taskJson, score] = result;
        return {
            task: JSON.parse(taskJson),
            score: parseFloat(score)
        };
    }

    async peek() {
        const result = await this.redis.zrange(this.queueName, 0, 0, 'WITHSCORES');
        if (result.length === 0) return null;
        return { task: JSON.parse(result[0]), score: parseFloat(result[1]) };
    }

    async size() {
        return await this.redis.zcard(this.queueName);
    }
}

// Implementación del worker
class PriorityWorker {
    constructor(redis, queueName, options = {}) {
        this.queue = new RedisPriorityQueue(redis, queueName);
        this.handlers = new Map();
        this.running = false;
        this.pollInterval = options.pollInterval || 100;
        this.concurrency = options.concurrency || 1;
    }

    registerHandler(taskType, handler) {
        this.handlers.set(taskType, handler);
    }

    async start() {
        this.running = true;
        const workers = Array(this.concurrency).fill().map(() => this.workerLoop());
        await Promise.all(workers);
    }

    async workerLoop() {
        while (this.running) {
            const item = await this.queue.dequeue();
            if (!item) {
                await this.sleep(this.pollInterval);
                continue;
            }

            const { task } = item;
            const handler = this.handlers.get(task.type);

            if (handler) {
                try {
                    await handler(task.payload);
                } catch (err) {
                    console.error(`Task ${task.id} falló:`, err);
                }
            }
        }
    }

    sleep(ms) {
        return new Promise(resolve => setTimeout(resolve, ms));
    }

    stop() {
        this.running = false;
    }
}

// Uso
const redis = new Redis();
const worker = new PriorityWorker(redis, 'task-queue', { concurrency: 2 });

worker.registerHandler('email', async (payload) => {
    console.log(`Enviando email a ${payload.to}`);
});

worker.registerHandler('process-payment', async (payload) => {
    console.log(`Procesando pago ${payload.orderId}`);
});

// Enviar tareas con prioridades (1 = mayor)
async function submitTasks() {
    const queue = new RedisPriorityQueue(redis, 'task-queue');
    await queue.enqueue({ id: '1', type: 'email', payload: { to: 'user@example.com' } }, 4);
    await queue.enqueue({ id: '2', type: 'process-payment', payload: { orderId: 'ORD-123' } }, 1);
    await queue.enqueue({ id: '3', type: 'email', payload: { to: 'vip@example.com' } }, 2);
}

submitTasks().then(() => worker.start());

TypeScript (Cola de Prioridad Genérica con Heap)

// Priority Queue: elementos con mayor prioridad se procesan primero
class PriorityQueue<T> {
  private heap: { priority: number; data: T }[] = [];

  enqueue(data: T, priority: number): void {
    this.heap.push({ priority, data });
    this.bubbleUp(this.heap.length - 1);
  }

  dequeue(): T | null {
    if (this.heap.length === 0) return null;
    const top = this.heap[0];
    const last = this.heap.pop()!;
    if (this.heap.length > 0) {
      this.heap[0] = last;
      this.bubbleDown(0);
    }
    return top.data;
  }

  peek(): T | null { return this.heap.length > 0 ? this.heap[0].data : null; }
  size(): number { return this.heap.length; }
  isEmpty(): boolean { return this.heap.length === 0; }

  private bubbleUp(idx: number): void {
    while (idx > 0) {
      const parent = Math.floor((idx - 1) / 2);
      if (this.heap[idx].priority <= this.heap[parent].priority) break;
      [this.heap[idx], this.heap[parent]] = [this.heap[parent], this.heap[idx]];
      idx = parent;
    }
  }

  private bubbleDown(idx: number): void {
    while (true) {
      const left = 2 * idx + 1;
      const right = 2 * idx + 2;
      let largest = idx;
      if (left < this.heap.length && this.heap[left].priority > this.heap[largest].priority) largest = left;
      if (right < this.heap.length && this.heap[right].priority > this.heap[largest].priority) largest = right;
      if (largest === idx) break;
      [this.heap[idx], this.heap[largest]] = [this.heap[largest], this.heap[idx]];
      idx = largest;
    }
  }
}

// Uso: sistema de tickets de soporte
interface Ticket { id: string; subject: string; }

const ticketQueue = new PriorityQueue<Ticket>();
ticketQueue.enqueue({ id: "T1", subject: "Question" }, 1);   // Baja
ticketQueue.enqueue({ id: "T2", subject: "Bug" }, 3);         // Alta
ticketQueue.enqueue({ id: "T3", subject: "Feature" }, 2);     // Media
ticketQueue.enqueue({ id: "T4", subject: "Outage" }, 5);      // Crítica

console.log(ticketQueue.dequeue()?.id); // T4 (Outage, prioridad 5)
console.log(ticketQueue.dequeue()?.id); // T2 (Bug, prioridad 3)
console.log(ticketQueue.dequeue()?.id); // T3 (Feature, prioridad 2)
console.log(ticketQueue.dequeue()?.id); // T1 (Question, prioridad 1)

Explicación

Las colas de prioridad usan una estructura de datos heap (o sorted set) para mantener el ordenamiento:

  • Inserción: Las tareas llegan con un valor de prioridad asignado. Se colocan en el heap segun la prioridad, no el tiempo de llegada.
  • Extracción: El worker siempre saca el elemento en la cima del heap. Ese es el de mayor prioridad. Si varios elementos comparten la misma prioridad, el orden secundario (timestamp) asegura equidad.
  • Equidad dentro de la prioridad: Las tareas con la misma prioridad se procesan en orden FIFO.

Variantes

VarianteMecanismoIdeal Para
Heap binarioHeap en array en memoriaProgramación de tareas de un solo proceso, alto throughput
Sorted sets de RedisEstructura ordenada externaWorkers distribuidos, cola persistente
Fair queuing ponderadoAsignación de ancho de banda proporcionalControl de tráfico de red, rate limiting de APIs
Cola de retroalimentacion multinivelAjuste live de prioridadProgramación de procesos de sistema operativo
Basado en plazosPrimero el plazo más cercanoSistemas en tiempo real, procesamiento guiado por SLA

Mejores Prácticas

  • Prevenir el hambre de tareas de baja prioridad. Estas tareas deben ejecutarse eventualmente: implementar envejecimiento (aumentar la prioridad con el tiempo) o una cuota mínima.
  • Mantener los niveles de prioridad limitados. Demasiados niveles (más de 20) hacen que el sistema sea difícil de razonar y no ayudan al throughput. Con 3-5 niveles alcanza.
  • Documentar las asignaciones de prioridad. Dejar claro qué se considera CRÍTICO versus ALTO para que los equipos no pongan todo en la máxima prioridad.
  • Monitorear la profundidad de la cola por prioridad. Un backlog creciente de tareas de ALTA prioridad señala un problema de capacidad, no solo descuido de las de BAJA.
  • Considerar la apropiación. Si llega una tarea CRÍTICA mientras se ejecuta una de BAJA prioridad, ¿debería pausarse la de BAJA? Para estrategias de throttling, ver el patrón Throttling.

Errores Comunes

  • Todo es de ALTA prioridad. Cuando todo es de alta prioridad, la cola degenera en FIFO y el sistema pierde su valor.
  • Ignorar el hambre de tareas. Una cola llena de tareas de ALTA y CRÍTICA prioridad puede nunca procesar las de BACKGROUND. Usar envejecimiento o cupos de tiempo.
  • Cálculos de prioridad complejos. Si calcular la prioridad toma más que procesar la tarea, se agregó más carga que cualquier beneficio.
  • Falta de visibilidad. Sin métricas que muestren la profundidad de la cola por prioridad, los operadores no tienen forma de saber si el sistema se comporta como se espera.
  • Prioridades codificadas en duro. Las prioridades de negocio cambian: hacer la asignación de prioridad configurable.

Ejemplos del Mundo Real

Kubernetes

Kubernetes usa una cola de prioridad cuando programa pods. Los pods con mayor priorityClassName se programan antes que los de menor prioridad. Si no se puede programar un pod de mayor prioridad, el scheduler puede apropiarse (evictar) pods de menor prioridad para hacer lugar.

RabbitMQ Priority Queue

RabbitMQ soporta colas de prioridad mediante el argumento x-max-priority, así los mensajes pueden adelantarse. Los mensajes de mayor prioridad se entregan antes que los de menor prioridad dentro de la misma cola, hasta el nivel máximo configurado.

AWS Lambda

Los mapeos de fuentes de eventos de SQS respetan la prioridad mediante colas separadas. Las organizaciones usan varias colas (crítica, normal, background) con diferentes asignaciones de concurrencia de Lambda para lograr procesamiento basado en prioridad.

Resumen

  • Usá una cola de prioridad cuando las tareas tienen distinta urgencia y no podés procesarlas todas a la vez.
  • Los heaps dan O(log n) de inserción y extracción; los sorted sets de Redis dan persistencia distribuida.
  • Mantené los niveles de prioridad chicos (3-5) para que el sistema sea predecible y debuggeable.
  • Prevení el hambre con envejecimiento o cupos mínimos para las tareas de baja prioridad.
  • Monitoreá la profundidad de la cola por prioridad. Un backlog creciente de ALTA significa que tenés un problema de capacidad.
  • Documentá qué se considera CRÍTICO versus ALTO para que los equipos no pongan todo al máximo.
  • Una sola cola es más simple; varias colas (una por prioridad) escalan y aíslan mejor.

See Also

Preguntas frecuentes

¿Cuál es la diferencia entre una cola de prioridad y fair queuing ponderado?

Una cola de prioridad siempre toma el ítem de mayor prioridad primero. El fair queuing ponderado, en cambio, asigna una porción proporcional de recursos a cada clase de prioridad, así las de menor prioridad no mueren de hambre.

¿Cómo evito que las tareas de baja prioridad mueran de hambre?

Usa envejecimiento de tareas: aumenta la prioridad a medida que esperan. También puedes asignar franjas de tiempo fijas a cada nivel o pasar a fair queuing en lugar de prioridad estricta.

¿Puedo cambiar la prioridad de una tarea después del envío?

Sí, pero tenés que quitarla primero. Actualiza la prioridad y vuélvela a insertar. En Redis es un zrem seguido de un zadd. En PriorityBlockingQueue de Java, remuévela y vuélvela a ofrecer; la cola no se reordena sola.

¿Las colas de prioridad son justas?

Las colas de prioridad estrictas no son justas para las tareas de baja prioridad. Si importa la equidad, agrega envejecimiento, limita la apropiación o pasa a un modelo de asignación proporcional.

¿Debo usar una cola de prioridad o varias?

Una sola cola es más simple, pero puede convertirse en un cuello de botella. Varias colas (una por prioridad con workers separados) escalan y aíslan mejor, pero agregan complejidad operativa.

¿Cuándo elijo una cola de prioridad en lugar de una FIFO?

Usa una cola de prioridad cuando los elementos tengan distinta urgencia: tickets críticos antes que preguntas, jobs de alto valor antes que batch. Usa una cola FIFO cuando el orden de llegada importe: pedidos, mensajes, logs de transacciones. La cola de prioridad reordena por urgencia; FIFO conserva el orden de llegada. Para sistemas de soporte, una cola de prioridad. Para procesamiento de transacciones, FIFO. Para el scheduling de un SO, una cola de prioridad por prioridad de procesos.