asyncio.Queue.shutdown ofrece una forma explícita de cerrar colas asíncronas, liberar productores y consumidores bloqueados y evitar tareas colgadas al finalizar una aplicación. En sistemas con workers, crawlers, integraciones, pipelines y procesamiento por lotes, el cierre suele ser más difícil que la ejecución normal. Esta guía explica el estado de cierre, el tratamiento de QueueShutDown, el modo inmediato y las pruebas necesarias para proteger el trabajo pendiente.
Por qué una cola necesita cierre
Una cola conecta productores que llaman a put() con consumidores que llaman a get(). Durante la ejecución normal el contrato es sencillo. Al detener el servicio, un consumidor puede esperar un elemento que nunca llegará, un productor puede quedar bloqueado porque la cola limitada está llena y el coordinador puede esperar join() indefinidamente si falta un task_done().
Las versiones anteriores suelen usar valores centinela como None. El patrón sigue siendo válido, pero exige un centinela por worker, un valor que no pueda confundirse con datos reales y una estrategia para productores bloqueados. shutdown() centraliza ese estado en la propia cola.
Cómo funciona asyncio.Queue.shutdown
Después de queue.shutdown(), la cola deja de aceptar elementos. Las llamadas futuras a put() lanzan QueueShutDown y los productores ya bloqueados se liberan con la misma excepción. Los consumidores aún pueden retirar los elementos existentes. Cuando la cola queda vacía, las llamadas posteriores a get() lanzan QueueShutDown.
Así se obtiene un cierre gradual: detener la entrada, procesar lo pendiente, esperar join() y permitir que los workers salgan al detectar la cola cerrada. Consulta la documentación oficial de colas de asyncio y la guía de tareas y cancelación.
Ejemplo de cierre gradual
import asyncio
async def worker(nombre: str, cola: asyncio.Queue[int]) -> None:
while True:
try:
elemento = await cola.get()
except asyncio.QueueShutDown:
print(f"{nombre}: cola cerrada")
return
try:
await asyncio.sleep(0.1)
print(f"{nombre}: procesó {elemento}")
finally:
cola.task_done()
async def main() -> None:
cola: asyncio.Queue[int] = asyncio.Queue(maxsize=10)
workers = [asyncio.create_task(worker(f"worker-{i}", cola)) for i in range(3)]
for elemento in range(20):
await cola.put(elemento)
cola.shutdown()
await cola.join()
await asyncio.gather(*workers)
asyncio.run(main())El orden es esencial. Primero termina la producción, luego se cierra la entrada, después join() espera que cada elemento tenga su task_done() y finalmente los workers salen al solicitar otro elemento.
task_done y join
Por cada elemento devuelto por get() debe existir exactamente un task_done(). Colocarlo en finally evita que una excepción de procesamiento deje el contador interno inconsistente. Si se omite, join() no termina; si se llama de más, la cola genera un error.
El cierre no reemplaza el seguimiento de tareas. Solo controla la entrada de trabajo y la liberación de operaciones bloqueadas. Un pipeline robusto combina shutdown, join y task_done.
Cuándo usar el modo inmediato
queue.shutdown(immediate=True) prioriza detenerse ahora en lugar de completar todo. La cola se vacía y las operaciones bloqueadas se liberan. Esto puede romper la expectativa habitual de join(), que puede regresar aunque los elementos no se hayan procesado realmente.
Úsalo ante errores fatales, apagado forzado, trabajo invalidado o destrucción inminente del proceso. No debe ser el valor predeterminado. Pagos, mensajes, importaciones y escrituras pueden perderse. Registra cuántos elementos fueron abandonados y persiste el trabajo importante fuera de memoria.
Productores bloqueados y backpressure
Una cola con maxsize aplica backpressure. Cuando está llena, put() espera. Durante el cierre, esos productores deben liberarse. QueueShutDown les proporciona una salida clara.
async def productor(cola: asyncio.Queue[str]) -> None:
for valor in fuente_de_datos():
try:
await cola.put(valor)
except asyncio.QueueShutDown:
guardar_checkpoint(valor)
returnNo captures una excepción genérica para continuar el bucle, porque convertirías el cierre en una secuencia de fallos. Trata QueueShutDown de forma específica, libera recursos y termina.
Integración con TaskGroup
asyncio.TaskGroup combina bien con las colas porque define una vida útil estructurada para los workers. Un coordinador puede crear tareas, producir elementos, iniciar el cierre y esperar el vaciado en un único ámbito. Consulta el artículo de Academify sobre TaskGroup y concurrencia estructurada.
También son útiles asyncio.Runner en Python, asyncio.Barrier en Python, queue.SimpleQueue en Python y asyncio.eager_task_factory. Estos contenidos amplían el ciclo de vida, la sincronización y el coste de planificación.
Cancelación y cierre no son iguales
Cancelar todos los workers puede ser necesario, pero no equivale a cerrar la cola. La cancelación interrumpe tareas; shutdown cambia el contrato de la cola. Para un cierre gradual, rechaza trabajo nuevo, procesa lo pendiente y cancela solo las tareas que superen un timeout.
cola.shutdown()
try:
async with asyncio.timeout(30):
await cola.join()
except TimeoutError:
cola.shutdown(immediate=True)
for tarea in workers:
tarea.cancel()Este patrón ofrece una ventana de finalización y luego aplica una política de emergencia. Si capturas CancelledError, limpia los recursos y normalmente vuelve a lanzar la excepción.
Errores comunes
Un error es cerrar la cola mientras los productores siguen activos sin tratar QueueShutDown. Otro es olvidar task_done. También es frecuente usar el modo inmediato y afirmar que todo fue procesado. Guardar referencias de las tareas es imprescindible para poder esperarlas o cancelarlas. Mezclar centinelas y shutdown sin reglas claras crea caminos de salida duplicados.
Cómo probar el ciclo completo
Prueba una cola vacía, una cola con elementos, una cola llena con productor bloqueado y un consumidor bloqueado en get(). Comprueba que los nuevos puts fallan después del cierre, que el modo gradual procesa todo y que el inmediato termina sin deadlock. Añade timeouts a las pruebas.
Inyecta fallos dentro del worker. Si el elemento ya fue retirado, task_done debe ejecutarse. Define si el trabajo fallido se reintenta, se persiste o se envía a una dead-letter queue. La política debe estar cubierta por pruebas.
Compatibilidad entre versiones
La API depende de la versión de Python. Las bibliotecas compatibles con versiones anteriores pueden usar un adaptador: llamar a shutdown cuando exista y usar centinelas cuando no. Declara la versión mínima y evita asumir que todos los entornos ya incluyen el método.
Observabilidad operativa
Registra tamaño de cola, tareas pendientes, productores bloqueados, latencia, duración del cierre gradual y elementos abandonados por cierre inmediato. Los logs estructurados deben indicar la causa del cierre y si el pipeline se vació correctamente. Estas métricas diferencian un despliegue limpio de una pérdida de datos o un bloqueo.
Conclusión
asyncio.Queue.shutdown convierte la finalización de pipelines asíncronos en un estado explícito. La estrategia segura consiste en detener la entrada, procesar lo pendiente, mantener el contador correcto, esperar a los workers y reservar el modo inmediato para emergencias. Con timeouts, tratamiento específico de QueueShutDown y pruebas de operaciones bloqueadas, el servicio termina de forma predecible.







