如何用 asyncio.Queue 实现生产者 消费者模式调度 Agent 任务,并支持优雅关闭和优先级调度?
面对 Agent 这种既要稳定吞吐又要弹性响应的系统,我几乎每次都把 asyncio.Queue 作为任务流的中枢神经。它天生就是为异步生产者-消费者准备的,不过要把它用到“优雅关闭”和“优先级调度”的程度,需要一些约定好的手势。
🔄 基础模式:一个队列 + 一组合格的工人¶
最简单的模型:主 Agent 是生产者,把任务(比如 LLM 请求、工具调用)往队列里 put;一组 worker 协程是消费者,不断 get 并执行。

import asyncio
async def producer(queue: asyncio.Queue):
for task in incoming_tasks:
await queue.put(task)
async def worker(queue: asyncio.Queue, name: str):
while True:
task = await queue.get()
try:
await process(task)
except Exception as e:
log_error(e)
finally:
queue.task_done() # 通知队列本任务完成
启动时:
queue = asyncio.Queue(maxsize=100) # 有背压上限
workers = [asyncio.create_task(worker(queue, f"w{i}")) for i in range(5)]
producer_task = asyncio.create_task(producer(queue))
🛑 优雅关闭:让队列说“打烊了”,别直接拔电源¶
直接 cancel() worker 会导致正在执行的任务中断,且队列里的任务丢失。正确的做法是先关门,再清场:
-
生产者先停止放新任务。
-
等待队列里已有的任务全部被取走处理完 (
queue.join())。 -
向每个 worker 发送 sentinel 信号(例如
None),让它们自行退出。
async def graceful_shutdown(queue, workers, producer_task):
# 1. 等待生产者完成
await producer_task
# 2. 等待队列清空
await queue.join()
# 3. 发送退出信号(每个 worker 一个)
for _ in workers:
await queue.put(None) # None 作为哨兵
# 4. 等待所有 worker 退出
await asyncio.gather(*workers)
# worker 中处理:
async def worker(queue, name):
while True:
task = await queue.get()
if task is None: # 收到打烊通知
queue.task_done()
break
try:
await process(task)
finally:
queue.task_done()
这里的关键是:sentinel 的数量必须等于 worker 的数量,而且要在 queue.join() 之后再放,否则可能被某个 worker 提前拿到而其他 worker 拿不到。
🚦 优先级调度:双队列 + “先看贵宾室”¶
asyncio.PriorityQueue 要求元素可排序,且数字越小优先级越高,但容易遇到任务对象不可比、优先级相同导致比较失败等问题。
我更倾向于用双队列,把高低优先级任务物理隔离开。
Worker 每次循环时,优先尝试从高优队列取(非阻塞),取不到才从低优队列取(可阻塞)。这样保证了高优任务几乎无延迟。
high_q = asyncio.Queue()
low_q = asyncio.Queue()
async def priority_worker(high_q, low_q, name):
while True:
# 优先拿高优,不阻塞
try:
task = high_q.get_nowait()
except asyncio.QueueEmpty:
# 高优空,才从低优拿,可短暂阻塞
try:
task = await asyncio.wait_for(low_q.get(), timeout=0.5)
except asyncio.TimeoutError:
continue # 没任务就继续循环
if task is None: # 统一用 None 做哨兵
# 确保 task_done 后 break
...
防止低优任务饿死:给低优任务打时间戳,如果等待超过一定阈值(比如 10 秒),可以自动把它升级到高优队列,或者统计低优队列的等待时间动态调整。
🧩 加上并发控制:Semaphore 是队列的好搭档¶
只靠 worker 数量控制并发还不够精细,我们可以在 process(task) 前加一层 asyncio.Semaphore,限制同时执行的总任务数,无论队列怎么塞,下游资源都稳定。
sem = asyncio.Semaphore(10)
async def bounded_worker(...):
while True:
task = await get_task()
async with sem:
await process(task)
这样队列长度只决定排队任务的堆积量,不影响系统稳定性。
最后想说的是,一个 Agent 的调度系统,优雅关闭和优先级处理不是附加题,而是必答题。线上谁都不想看到发版时丢任务,或者用户请求被批处理卡死。用队列把这些职责凝固成代码,比你每次手写一堆 if 要稳妥得多。我习惯于把 Queue 当成一条河,Semaphore 是水闸,优先级是分流渠——把流控制好了,Agent 这艘船才不会翻。