El módulo queue ofrece colas sincronizadas para comunicación segura entre threads. Implementa FIFO, LIFO, prioridad y una cola simple, con operaciones bloqueantes, timeouts, capacidad máxima y seguimiento de tareas. El patrón productor-consumidor resulta más claro porque los workers intercambian mensajes en vez de modificar estructuras compartidas directamente.
Una cola no vuelve correcto el procesamiento automáticamente. La aplicación todavía necesita ownership, esquema de mensajes, backpressure, manejo de excepciones, shutdown e idempotencia. También debes distinguir queue.Queue para threads de multiprocessing.Queue y asyncio.Queue.
Crea una cola FIFO
Queue entrega elementos en orden de entrada.
from queue import Queue
cola = Queue()
cola.put("primero")
cola.put("segundo")
print(cola.get())
La cola sincroniza su estado interno entre threads.
Define capacidad
maxsize limita aproximadamente la cantidad de elementos pendientes.
cola = Queue(maxsize=100)
Una cola limitada crea backpressure cuando los consumidores no acompañan.
put bloqueante
Por defecto, put() espera hasta existir espacio.
cola.put(item, timeout=5)
Con timeout, lanza queue.Full si no puede insertar.
put_nowait
put_nowait() intenta insertar inmediatamente.
from queue import Full
try:
cola.put_nowait(item)
except Full:
registrar_descarte(item)
Define una política: esperar, rechazar, persistir, reducir producción o devolver error.
get bloqueante
get() espera un elemento.
item = cola.get(timeout=2)
Al vencer el plazo, lanza queue.Empty.
No uses empty para decidir get
empty(), full() y qsize() son snapshots aproximados en concurrencia.
Otra thread puede cambiar la cola antes de la siguiente operación. Usa métodos no bloqueantes o timeout y trata la excepción.
Productor y consumidor
from threading import Thread
from queue import Queue
cola = Queue(maxsize=20)
def consumidor():
while True:
item = cola.get()
try:
procesar(item)
finally:
cola.task_done()
thread = Thread(target=consumidor, daemon=True)
thread.start()
El finally actualiza el contador incluso después de un error.
task_done
Cada elemento retirado con get() debe recibir exactamente una llamada a task_done() cuando termina su trabajo.
Llamar de más genera ValueError; olvidarlo puede bloquear join() indefinidamente.
join
cola.join() espera que todas las tareas insertadas sean marcadas como concluidas.
for item in items:
cola.put(item)
cola.join()
No detiene workers; solo espera el contador de tareas.
Sentinelas para shutdown
Un objeto único puede indicar que el consumidor debe parar.
PARAR = object()
def consumidor():
while True:
item = cola.get()
try:
if item is PARAR:
return
procesar(item)
finally:
cola.task_done()
Envía una sentinela por consumidor independiente.
Orden de shutdown
Detén productores, termina la entrada normal, envía sentinelas, espera cola.join() y luego une las threads.
Un orden incorrecto puede dejar trabajo después de las sentinelas o bloquear productores en una cola llena.
APIs modernas de shutdown
Versiones recientes ofrecen shutdown explícito para colas sincronizadas. Estas APIs impiden nuevas inserciones y despiertan operaciones bloqueadas según la política.
Comprueba la versión mínima y maneja la excepción específica. Las sentinelas siguen siendo útiles para compatibilidad amplia.
LifoQueue
LifoQueue entrega primero el elemento más reciente.
from queue import LifoQueue
pila = LifoQueue()
pila.put("a")
pila.put("b")
print(pila.get())
Puede favorecer trabajo reciente, pero los elementos antiguos pueden sufrir starvation.
PriorityQueue
PriorityQueue entrega el menor valor primero.
from queue import PriorityQueue
cola = PriorityQueue()
cola.put((10, "normal"))
cola.put((1, "urgente"))
Guarda prioridad, secuencia y payload para evitar comparar tareas.
Empates estables
from itertools import count
contador = count()
cola.put((prioridad, next(contador), tarea))
La secuencia conserva orden de inserción e impide comparar payloads incompatibles.
SimpleQueue
SimpleQueue es una FIFO no limitada con API menor.
from queue import SimpleQueue
cola = SimpleQueue()
cola.put(item)
item = cola.get()
Úsala cuando no necesites capacidad, task tracking o backpressure interno.
Elige la cola correcta
Usa Queue para FIFO limitada y seguimiento, LifoQueue para pila, PriorityQueue para prioridad y SimpleQueue para comunicación simple sin límite.
Documenta la elección porque afecta fairness y memoria.
Excepciones en consumidores
Una falla no debería matar silenciosamente todos los workers.
try:
procesar(item)
except Exception as error:
registrar_fallo(item, error)
finally:
cola.task_done()
Decide retry, dead-letter, shutdown o continuación.
Retries
No reinsertes indefinidamente el mismo item.
Controla intentos, backoff y deadline y separa fallos terminales. Las operaciones con side effects deben ser idempotentes.
Varios consumidores
Más threads pueden aumentar throughput de I/O, pero también presionan bases, APIs y filesystem.
Dimensiona workers según la capacidad downstream.
GIL y CPU
Threads normalmente no aceleran código Python puramente CPU-bound en builds tradicionales.
Para CPU, evalúa procesos o código nativo que libere el GIL. Consulta multiprocessing en Python.
Payloads mutables
La cola transfiere una referencia, no una copia profunda.
Después de put(), el productor debería tratar el objeto como propiedad del consumidor o enviar datos inmutables.
Esquemas de mensajes
Usa dataclasses, NamedTuple u objetos pequeños con tipo, payload, ID, intento y deadline.
Evita tuplas largas y diccionarios sin esquema.
Backpressure y locks
Una cola limitada evita crecimiento ilimitado, pero productores bloqueados pueden crear deadlock.
No llames put() bloqueante mientras mantienes un lock que necesita el consumidor.
Fairness
El orden de wakeup entre threads no debe considerarse justicia perfecta.
Si fairness es requisito, modela cuotas, clases de servicio o colas separadas.
Integración con selectors
Workers pueden enviar comandos por una cola y un loop de I/O despertarse mediante socketpair o pipe.
Consulta selectors en Python.
Integración con ThreadPoolExecutor
Los executors ya mantienen una cola interna. Evita otra capa sin necesidad.
Una cola propia sirve para capacidad, prioridad, mensajes explícitos o workers personalizados.
Timeouts y cancelación
Un timeout en get() permite revisar una flag de parada.
while not detener.is_set():
try:
item = cola.get(timeout=0.5)
except Empty:
continue
Cancelar trabajo ya iniciado requiere cooperación.
Daemon threads
Las daemon threads no mantienen vivo el proceso, pero pueden interrumpirse sin cleanup.
Usa threads normales y shutdown explícito para datos importantes.
Observabilidad
Mide tamaño aproximado, espera, edad del item, tasa de entrada, throughput, fallos y retries.
qsize() sirve como métrica aproximada, no garantía lógica.
Seguridad
Limita cantidad y tamaño de mensajes. No encoles callbacks arbitrarios de usuarios.
Redacta secretos en logs y autoriza trabajo antes de producir tareas.
Pruebas
Prueba cola vacía y llena, timeout, varios productores y consumidores, excepción, sentinela, retry, shutdown, task_done() y deadlocks por locks.
Usa eventos de sincronización en vez de sleeps exactos.
Errores comunes
Los fallos frecuentes son usar empty() como garantía, olvidar task_done(), confundir join de cola con join de thread, enviar pocas sentinelas, permitir crecimiento ilimitado, bloquear con un lock y modificar payloads después de put().
Conclusión
queue ofrece comunicación sincronizada entre threads. Usa capacidad para backpressure, task_done()/join() para seguimiento y sentinelas o shutdown explícito para terminar.
Consulta la documentación oficial de queue, heapq en Python y selectors en Python.







