El módulo queue ofrece colas sincronizadas para pasar trabajo entre threads. Evita implementar bloqueos alrededor de una lista compartida y hace explícito el backpressure del productor.
from queue import Queue
from threading import Thread
fila: Queue[int | None] = Queue(maxsize=100)
def consumidor() -> None:
while True:
item = fila.get()
try:
if item is None:
return
print(item * 2)
finally:
fila.task_done()
worker = Thread(target=consumidor)
worker.start()
for valor in range(5):
fila.put(valor)
fila.put(None)
fila.join()
worker.join()
put() espera cuando una cola limitada está llena y get() espera un elemento. Cada elemento retirado requiere task_done(), también si el trabajo falla, por eso conviene usar finally. join() espera hasta que no queden tareas pendientes.
Buenas prácticas
Elige Queue para FIFO, LifoQueue como pila y PriorityQueue para tuplas con prioridad. No uses qsize() para tomar una decisión concurrente: el tamaño puede cambiar antes de la siguiente instrucción. Define el cierre con un centinela por consumidor o con la API de shutdown disponible en tu versión.
Compara modelos en multithreading y multiprocessing en Python. Para código con event loop, revisa la guía de async y await.
La documentación oficial del módulo queue, consultada el 22 de julio de 2026, detalla la API y sus garantías.
Por qué una cola es mejor que una lista compartida
Una lista protegida con Lock puede funcionar, pero la espera, la notificación, la capacidad y el cierre pronto se convierten en código propio. Consultar repetidamente if elementos: desperdicia CPU o agrega latencia. Queue reúne almacenamiento, bloqueo y señalización en una abstracción probada. Productores y consumidores operan a la vez sin corromper el contenedor.
La garantía cubre el intercambio, no todos los efectos posteriores. Si varios consumidores modifican el mismo diccionario, archivo u objeto mutable, ese recurso necesita su propia sincronización.
Backpressure con maxsize
Una cola ilimitada permite que productores rápidos consuman memoria mientras el destino está lento. Con Queue(maxsize=20), el siguiente put() espera hasta que haya espacio. Así la presión vuelve al productor y se limita el trabajo pendiente. Elige la capacidad según el tamaño de cada elemento, la memoria aceptable, la duración de los picos y la latencia.
put_nowait() y get_nowait() lanzan Full o Empty. Son útiles con una política explícita, como descartar telemetría poco importante. No consultes antes full() o empty() suponiendo que el resultado seguirá vigente: otra thread puede cambiar la cola.
from queue import Full, Queue
eventos: Queue[str] = Queue(maxsize=2)
def publicar(evento: str) -> bool:
try:
eventos.put(evento, timeout=0.5)
return True
except Full:
return False
El timeout evita una espera infinita y permite registrar, reintentar o cancelar.
Varios consumidores y cierre
Con centinelas, envía uno por consumidor. Un único centinela detiene solamente la thread que lo recibe. Escoge un valor imposible de confundir con trabajo válido. None basta si los elementos reales nunca son nulos; un object() privado es más seguro en código genérico.
from queue import Queue
from threading import Thread
DETENER = object()
trabajos: Queue[object] = Queue()
def ejecutar() -> None:
while True:
trabajo = trabajos.get()
try:
if trabajo is DETENER:
return
procesar(trabajo)
except Exception:
registrar_fallo(trabajo)
finally:
trabajos.task_done()
workers = [Thread(target=ejecutar, name=f"worker-{i}") for i in range(3)]
for worker in workers:
worker.start()
for elemento in cargar_trabajos():
trabajos.put(elemento)
for _ in workers:
trabajos.put(DETENER)
trabajos.join()
for worker in workers:
worker.join()
Queue.join() espera que todos los elementos reciban task_done(). Thread.join() espera que una thread termine. Usar ambos confirma que el trabajo quedó contabilizado y que ningún worker sigue vivo.
Errores y tareas pendientes
Cada put() incrementa el contador de tareas sin finalizar. task_done() reduce ese contador, no elimina un elemento. Demasiadas llamadas producen ValueError; olvidar una puede bloquear join() para siempre. Colócala en finally, pero solo después de un get() exitoso.
Define qué ocurre con los fallos. Una excepción no capturada termina el consumidor y puede detener el flujo. El código real puede registrar contexto, enviar el elemento a una cola de fallos o reintentar con un límite estricto. Reinsertarlo indefinidamente crea un bucle e impide el cierre.
FIFO, LIFO y prioridad
Queue ofrece FIFO y suele ser la opción inicial. LifoQueue devuelve primero lo reciente, útil al explorar árboles, aunque puede postergar trabajos antiguos. PriorityQueue retira el valor menor. Las tuplas (prioridad, payload) fallan si un empate obliga a comparar payloads no ordenables.
from dataclasses import dataclass, field
from queue import PriorityQueue
from typing import Any
@dataclass(order=True)
class Trabajo:
prioridad: int
secuencia: int
payload: Any = field(compare=False)
pendientes: PriorityQueue[Trabajo] = PriorityQueue()
pendientes.put(Trabajo(2, 0, {"tipo": "informe"}))
pendientes.put(Trabajo(1, 1, {"tipo": "alerta"}))
La secuencia desempata sin comparar diccionarios. Documenta si un número menor representa mayor urgencia.
Queue, SimpleQueue, asyncio y procesos
SimpleQueue es una FIFO ilimitada con API reducida. Úsala si no necesitas capacidad, task_done() ni join() de tareas. asyncio.Queue coordina corutinas de un event loop y sus operaciones se esperan con await; no sirve como mecanismo general entre threads. multiprocessing.Queue serializa valores entre procesos y tiene otros costes y reglas de cierre.
Las threads suelen ayudar en trabajo de entrada y salida. El código Python intensivo en CPU puede quedar limitado por el GIL, por lo que los procesos pueden encajar mejor. La cola organiza el flujo, pero no elige el modelo de concurrencia.
Pruebas y observabilidad
Usa capacidades pequeñas en pruebas para forzar backpressure. Añade timeouts para que un defecto genere un fallo claro en vez de bloquear la suite. Cubre trabajo exitoso, excepción del consumidor, cola llena y cierre con varios workers. No afirmes qué consumidor recibirá cada elemento, pues la planificación no es determinista.
En producción mide elementos procesados, fallos, duración y espera. qsize() puede servir como métrica aproximada, aunque no como decisión de sincronización. Un crecimiento sostenido indica que la entrada supera la capacidad. Reduce la producción, aumenta workers respetando el destino u optimiza el procesamiento.
Lista de comprobación
Limita la cola cuando los productores puedan superar a los consumidores. Empareja cada get() exitoso con un único task_done(). Diseña cancelación y cierre antes de iniciar threads. Captura fallos con contexto suficiente. Finalmente, documenta la semántica de entrega: una Queue en memoria coordina, pero no aporta persistencia, transacciones ni recuperación tras finalizar el proceso.