跳转至

Agent 需要同时处理用户实时请求(高优先级)和后台批量任务(低优先级),如何设计任务调度?

面试时遇到这个场景题,说明对方在考察你在资源有限的情况下,如何做任务的“阶级划分”和“优雅降级”。我通常会先用一个贴地的比喻切入,再拆解成具体可落地的调度架构。


🚦 先理解本质:这不是功能需求,是资源分配的“政治问题”

你的 Agent 就像一个只有 3 个窗口的办事大厅:

  • VIP 窗口(用户实时请求):立刻响应,等久了会骂人。

  • 普通窗口(后台批量任务):可以排队,但不能永远不被叫号。

如果来一个 VIP 就插队,普通窗口的人会“饿死”;如果完全公平,VIP 的体验就崩了。所以我们要设计一套带优先级的、可抢占的、防饿死的调度规则。


🧩 核心方案:双队列 + 优先级调度 + 动态并发隔离

别把所有任务丢进一个池子里搅和。分成两个逻辑队列,配合 asyncio.PriorityQueue 或双队列协调:

⚡ 为什么不用单一 PriorityQueue?

单一优先级队列在纯协程模型里有坑:不能真正抢占。一个已经在运行的低优任务,不会因为队列后来插入高优任务而自动让出资源,必须等它主动 await。 所以我们用双队列 + 独立的 Worker 组,既能保证高优快速进入,又能通过动态调整 Worker 数量实现“软抢占”。


🔧 代码骨架:用 asyncio 实现分层调度

import asyncio
from asyncio import Queue, Task

class PriorityScheduler:
    def __init__(self, max_concurrency: int = 10):
        self.max_concurrency = max_concurrency
        self.high_queue = Queue()       # 高优任务
        self.low_queue = Queue()        # 低优任务
        self._running = False

    async def _worker(self, worker_id: int):
        """每个 worker 循环:优先取高优,无高优再取低优"""
        while self._running:
            # 优先尝试高优队列,用 get_nowait 非阻塞检查
            try:
                task = self.high_queue.get_nowait()
            except asyncio.QueueEmpty:
                # 高优空,取低优(可阻塞等待,给低优任务保底)
                try:
                    task = await asyncio.wait_for(
                        self.low_queue.get(), timeout=0.5
                    )
                except asyncio.TimeoutError:
                    continue   # 没任务就继续循环

            try:
                await task()
            except Exception as e:
                # 错误处理 + 日志
                pass
            finally:
                # 标记任务完成
                if task in self.high_queue._queue:
                    self.high_queue.task_done()
                else:
                    self.low_queue.task_done()

    async def run(self):
        self._running = True
        # 启动固定数量的 workers
        workers = [asyncio.create_task(self._worker(i))
                   for i in range(self.max_concurrency)]
        # 动态调整高优槽位:用单独的协程监控
        asyncio.create_task(self._adaptive_reserve())
        await asyncio.gather(*workers)

    async def _adaptive_reserve(self):
        """定期调整,保证低优任务不被饿死"""
        while self._running:
            # 这里可以根据低优队列长度和等待时间动态提升低优 worker 数
            await asyncio.sleep(1)

在实际业务中,high_queuelow_queue 可以直接封装成异步生成器,供路由层 put


🎯 如何给低优任务“兜底”并防止饿死?

  1. 保底槽位:比如总并发 10,设置至少 2 个 worker 只从低优队列取任务(可以启动两个专用 worker 跑 _low_worker)。

  2. 老化升级:给低优任务打时间戳,等待超过一定阈值(如 5 秒),自动升级为高优,避免永远轮不到。

  3. 动态槽位:如果高优队列持续为空,低优任务可临时占用全部槽位;一旦高优到来,部分 worker 立刻切回去。

async def _worker_with_priority_boost(self):
    while True:
        # 尝试高优
        task = await self._get_high_or_aged_low()
        ...

📊 方案对比:避免只会一种模式

查看内嵌表格

Agent 场景大多是 IO 密集型,所以异步 + 双队列是最均衡的方案。


🧪 压测时会发现的问题 & 我踩过的坑

  • Worker 阻塞:如果一个低优任务里写了大文件的同步 IO,哪怕只有一个,整个 worker 就被卡死,高优也进不来。解决方案是所有任务都强制异步,或者丢给线程池执行器并严格控制超时。

  • 高优风暴:突然涌来 500 个高优请求,如果让它们占满所有槽位,低优任务会完全停滞。我加了一个高优容量上限,比如高优最多同时占用 80% 槽位,剩下的留给低优喘气。

  • 监控信号:必须暴露 high_queue.qsize()low_queue.qsize()avg_wait_time 给 Prometheus,不然线上怎么死的都不知道。


💡 把这个调度想象成餐厅的出餐口

高优是堂食订单,低优是外卖预约单。高峰期,厨师长会喊:“堂食单优先,但外卖单不能停超过 10 分钟!”

你设计的调度器,就是那个懂得平衡的厨师长。面试时能说出“我会给低优任务加一个最大等待时间,超时就临时升级优先级”,面试官就知道你在真实环境里扛过流量,不是在搭乐高。