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
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
| Variante | Mecanismo | Ideal Para |
|---|---|---|
| Heap binario | Heap en array en memoria | Programación de tareas de un solo proceso, alto throughput |
| Sorted sets de Redis | Estructura ordenada externa | Workers distribuidos, cola persistente |
| Fair queuing ponderado | Asignación de ancho de banda proporcional | Control de tráfico de red, rate limiting de APIs |
| Cola de retroalimentacion multinivel | Ajuste live de prioridad | Programación de procesos de sistema operativo |
| Basado en plazos | Primero el plazo más cercano | Sistemas 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
- Documentación de heapq en Python: referencia oficial para colas de prioridad basadas en heaps.
- Java PriorityBlockingQueue: cola de prioridad thread-safe para workers concurrentes.
- Redis sorted sets: cómo funcionan
zaddyzpopminpara colas de prioridad distribuidas. - Colas de prioridad en RabbitMQ: cómo
x-max-prioritydeja que los mensajes se adelanten. - Priority classes en Kubernetes: cómo los pods se programan por prioridad.
- Para load leveling con colas, ver el patrón Queue-Based Load Leveling.
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.
Recursos Relacionados
Nivelación de Carga con Colas: Suaviza Picos de Tráfico
Usá Queue-Based Load Leveling para desacoplar productores y consumidores, absorber picos y procesar trabajo constante. Ejemplos en Python, Java y JavaScript.
PatternPatron de Programador-Agente-Supervisor
Coordina la programacion de trabajos resilientes con un supervisor que monitorea agentes, reinicia fallas y gestiona el ciclo de vida.
PatternPatron de Throttling
Limita la tasa a la que un sistema procesa solicitudes o consume recursos para prevenir sobrecarga, asegurar uso justo y mantener rendimiento predecible.
PatternPatrón Lock-Free Queue
Construir colas de alto throughput usando operaciones atomicas en lugar de locks. Multiples threads pueden encolar y desencolar concurrentemente sin bloqueo ni overhead de context-switch.
PatternPatró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.
PatternPatron Serverless Throttling
Maneja backpressure en serverless usando SQS, token buckets y limites de concurrencia para proteger servicios downstream de trafico burst.