跳转至

如何用 asyncio.Queue 实现生产者 消费者模式调度 Agent 任务,并支持优雅关闭和优先级调度?

面对 Agent 这种既要稳定吞吐又要弹性响应的系统,我几乎每次都把 asyncio.Queue 作为任务流的中枢神经。它天生就是为异步生产者-消费者准备的,不过要把它用到“优雅关闭”和“优先级调度”的程度,需要一些约定好的手势。


🔄 基础模式:一个队列 + 一组合格的工人

最简单的模型:主 Agent 是生产者,把任务(比如 LLM 请求、工具调用)往队列里 put;一组 worker 协程是消费者,不断 get 并执行。

image.png

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 会导致正在执行的任务中断,且队列里的任务丢失。正确的做法是先关门,再清场:

  1. 生产者先停止放新任务。

  2. 等待队列里已有的任务全部被取走处理完 (queue.join())。

  3. 向每个 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 这艘船才不会翻。