这两天"ax调度"突然成了圈内热词,好几个技术群都在转发。作为一个从 prompt 工程一路折腾到多 Agent 系统的一线开发,我第一反应是:大家终于开始正视"调度"这个东西了。很多人以为 Agent 应用 = 提示词 + 模型 + 工具调用,最多再加一层记忆。可真当你把三五个 Agent 丢到生产环境、让它们共同处理源源不断的业务请求时,最先崩掉的往往不是模型,反而是任务分配和执行秩序。本文不打算给 ax调度下什么标准定义——它本来也不是某个固定开源项目,而是一类调度思想。我把它理解为 Agent eXecution Scheduling,也就是"Agent 执行调度"。下面记录的是我落地这套调度模型时的完整思路、核心代码和踩坑记录,适合正在做多智能体应用、异步任务系统或 AI 客服/自动化流程的开发者参考。
1. 从"流程编排"到"ax调度":我踩过的第一步坑
1.1 最初的串行调用:卡的不是模型,是调度
最早我用一个很朴素的写法:把所有 Agent 放进一个列表,然后for agent in agents: result = await agent.run(task)。这个写法在 demo 里完全没问题,但生产环境很快就暴露了问题。比如一个复杂任务需要意图识别 Agent、知识检索 Agent、生成 Agent、质检 Agent 多次往返协作,总共要跑 8 到 12 次模型调用,每次 2 到 5 秒,串行下来轻松超过 30 秒。前端等不了,用户也等不了。
一开始我以为加个并发就能解决,后来发现并发只是把问题从"慢"变成了"乱":请求一多,任务不知道谁先谁后,有的 Agent 忙得不可开交,有的 Agent 闲着没事干,还有大量重复调用在浪费 token。这时候我才意识到,真正缺的是一个调度层。
调度层要回答的问题很简单:任务按什么顺序执行、谁先谁后、最多同时跑多少、单个任务最久能占用多久、失败之后怎么办。没有这层设计,Agent 越多,系统越混乱。我见过不少团队在 5 个 Agent 以内还能靠手写 if-else 勉强撑住,一旦超过 10 个 Agent、日均请求量上千,调度缺失的代价就会集中爆发:要么线上超时,要么队列堆积,要么模型调用费用失控。
1.2 不要指望 Agent 自己排队:它没有全局视野
很多人会问:Agent 自己不能判断优先级吗?不能。LLM 的 API 是无状态的,每个 Agent 只能看到自己这次请求的输入和输出,它不知道队列里还有多少任务,也不知道别的 Agent 正在忙什么。让 Agent 自己协商执行顺序,在工程上等于放弃控制权。
这就像餐厅后厨:不能靠每个厨师自己喊"我先做这道菜",必须有一块订单板,把所有菜按下单时间和优先级排好。调度器就是那块订单板。所谓 ax调度,核心就是把"让谁先做、什么时候做、最多几个同时做"从 Agent 内部抽出来,放到一个可编程、可观测、可控制的中间层。
想明白这一点之后,我的设计原则就变了:Agent 只做"执行",所有"决策"交给调度器。Agent 不需要知道自己为什么被调用、前面还有多少任务、其他 Agent 在干什么,只需要接收一个明确定义的任务上下文,然后执行并返回结构化结果。这样做还有一个好处:调度器可以随时调整策略,而不需要改动 Agent 的 prompt。
1.3 流程编排和调度不是一回事
这是我在项目里反复解释的一个概念。很多人会用 LangGraph、Dify 之类的流程编排工具,把节点连成一个图,然后觉得这就是调度。但流程编排解决的是"一个任务内部按什么顺序走",调度解决的是"多个任务之间怎么排队、怎么并发、怎么限流"。前者是图,后者是队列。
我见过最接近生产真相的方案是两层配合:流程编排负责单个任务内部的 DAG,调度层负责多个任务之间的全局队列和资源控制。两者互相配合,而不是互相替代。
对比一下两者的关注点:
- 流程编排:关注单个任务内部节点顺序,典型工具是 LangGraph、自研状态机,失败后重跑子节点即可。
- ax调度:关注任务之间的排队、并发、优先级、超时、重试,典型载体是 Redis Stream、消息队列加状态机,失败后要重排队、转人工、限流降级。
如果你只做流程编排不做调度,那么单个任务跑得再顺,也无法应对突发流量。如果你只做调度不做流程编排,那么单个任务内部的复杂依赖关系会变成一团乱麻。两个都做,才算完整的 Agent 系统骨架。
2. AX 调度器的核心模型:队列、优先级和可观测性
2.1 三层结构:接入层、调度层、执行层
我在项目里把 ax调度器拆成三层,职责非常清晰。接入层负责接收外部请求、校验参数、生成任务对象,打上task_id,塞进队列。调度层负责从队列里按规则取任务、分配 worker、维护任务状态、处理超时和重试。执行层负责真正调用 Agent 和工具函数。
三层之间只通过队列和状态存储通信,不直接互相调用。这样做的直接收益是:接入方只管发任务,执行方只管干活,中间的调度逻辑可以独立升级,不影响上下游。
任务状态至少要有pending、scheduled、running、succeeded、failed、retry,如果做 DAG 依赖还要有waiting_dependency。每次状态变更都要记录时间戳和原因。因为一旦线上出问题,没有状态流转记录,就只能靠猜。我就吃过这个亏:早期任务失败后只打印了一行错误,结果排查了半天也搞不清任务到底是超时失败还是 Agent 返回了非法结构。
2.2 优先级不是简单插队
任务优先级必须分级,而且要跟超时、重试策略绑定。我一般分 P0 到 P3:
| 优先级 | 适用场景 | 超时时间 | 最大重试次数 |
|---|---|---|---|
| P0 | 线上投诉、支付失败咨询 | 30 秒 | 2 次 |
| P1 | 普通客服咨询、工单处理 | 60 秒 | 2 次 |
| P2 | 批量数据生成、周报摘要 | 180 秒 | 1 次 |
| P3 | 离线分析、知识库索引 | 600 秒 | 0 次 |
这里的关键不是分级本身,而是每一级必须配套超时和重试策略。P0 任务如果重试 2 次还是失败,应该转人工,而不是继续烧 token。P3 任务失败可以直接丢弃,记录日志即可,不值得占用调度资源。
还有一个容易忽略的坑:优先级队列会让低优先级任务长期饿死。如果 P0 和 P1 任务持续进来,P2 和 P3 可能永远得不到执行。解决方法是老化机制:任务在队列里等待超过某个阈值,自动提升一级。我通常把 P2 任务等待超过 3 分钟提升到 P1,P3 任务等待超过 5 分钟提升到 P2。配合监控队列积压数量和各级任务平均等待时间,才能保证调度策略是健康的。
2.3 一个能跑通的最小 ax 调度器示例
我用 Python 的asyncio写过一个最小实现,代码不长,但把 ax调度的骨架表达得很清楚。核心是三个东西:PriorityQueue决定顺序,Semaphore限制并发,wait_for控制超时。
import asyncio import uuid from dataclasses import dataclass, field from enum import IntEnum class Priority(IntEnum): P0 = 0 P1 = 1 P2 = 2 P3 = 3 @dataclass(order=True) class Task: priority: int task_id: str = field(compare=False) payload: dict = field(compare=False) timeout: float = field(compare=False) max_retries: int = field(compare=False) retries: int = field(default=0, compare=False) class AXScheduler: def __init__(self, max_concurrency=5): self.queue = asyncio.PriorityQueue() self.sem = asyncio.Semaphore(max_concurrency) self.workers = [] async def submit(self, priority, payload, timeout, max_retries): task = Task(priority, str(uuid.uuid4()), payload, timeout, max_retries) await self.queue.put(task) async def _process(self, task): async with self.sem: try: result = await asyncio.wait_for( execute_agent(task), timeout=task.timeout ) print(f"任务 {task.task_id} 执行成功: {result}") except asyncio.TimeoutError: print(f"任务 {task.task_id} 超时") task.retries += 1 if task.retries <= task.max_retries: await self.queue.put(task) else: # 转人工或告警 print(f"任务 {task.task_id} 超过最大重试次数") except Exception as exc: print(f"任务 {task.task_id} 执行异常: {exc}") task.retries += 1 if task.retries <= task.max_retries: await self.queue.put(task) async def worker_loop(self): while True: task = await self.queue.get() try: await self._process(task) finally: self.queue.task_done() async def start(self, worker_count=3): self.workers = [ asyncio.create_task(self.worker_loop()) for _ in range(worker_count) ] await self.queue.join() for w in self.workers: w.cancel() async def execute_agent(task: Task) -> str: # 这里是执行层,实际项目中替换为 LLM 调用或工具函数 await asyncio.sleep(1) return f"done-{task.task_id}"这个示例里execute_agent就是执行层入口,实际项目中替换成调用 LLM 或工具的逻辑。为什么用PriorityQueue而不是普通队列?因为Task这个 dataclass 上定义了priority字段,并且设置了order=True,队列会按优先级从小到大排列,P0 任务永远最先被取出。Semaphore保证同一时刻最多只有 5 个 Agent 在执行,防止外部 API 被瞬间打爆。wait_for则给每个任务设了硬超时,避免一个模型卡住就把 worker 占死。
这个进程内调度器的问题也很明显:进程重启任务会丢。所以生产环境建议用 Redis Stream 或 PostgreSQL 表做持久化队列,但调度模型本身是一样的。先把最小实现跑通,再迁移到持久化队列,是最稳妥的路径。
3. 两个真实案例:客服工单与多 Agent 协作的调度参数
3.1 案例一:客服工单自动分类与回复
第一个真实场景是客服工单处理。流程大致是:工单进来,先用一个轻量模型做意图分类和紧急程度判断,然后检索知识库,再让生成 Agent 草拟回复,最后转人工审核。这个场景对调度最敏感,因为工单里的 P0 任务(投诉、支付失败)必须立刻处理,普通咨询则可以排队。
我设计的调度参数很简单:意图分类阶段先判断紧急程度,紧急工单打上 P0 标签,超时 30 秒,重试 2 次;普通咨询打 P1 标签,超时 60 秒,重试 2 次;夜间批量整理历史工单则打 P3 标签,超时 600 秒,不重试。
上线之后效果很明显。P0 工单从进入队列到开始处理,平均等待时间从原来的"看运气"变成稳定在 5 秒以内。普通工单也没有被 P0 完全堵死,因为并发控制保证了至少有 2 个 worker 专门处理 P1 任务。还有一个容易忽略的收益:token 成本下来了。之前没有调度时,高峰期同一个工单可能被三个 Agent 重复处理,现在每个任务有明确的task_id,执行层能通过幂等键判断是否已经处理过。
这个案例给我的教训是:不要试图让意图分类 Agent 自己决定"我是不是 P0",而是让分类 Agent 输出结构化标签,由调度器决定优先级。这样即使分类结果有偏差,人工审核也能在调度层修正,而不需要改 prompt。
3.2 案例二:多 Agent 协作生成周报和代码审查
第二个场景是多 Agent 协作。比如生成周报:一个 Agent 负责汇总提交记录,一个 Agent 负责分析数据,一个 Agent 负责撰写正文,最后一个 Agent 负责格式校验。这些 Agent 之间有依赖关系,不是简单的谁先谁后,而是"B 必须等 A 完成才能开始"。
这种场景光靠优先级队列不够,必须在调度层维护任务依赖 DAG。我的做法是给每个任务加一个dependencies集合,调度器维护一个依赖计数表。只有当一个任务的所有依赖都处于succeeded状态时,它才会进入就绪队列。实现上其实不复杂:任务完成时,找到所有依赖它的下游任务,把依赖计数减一,计数归零就放入就绪队列。
这个机制在代码审查场景里更有价值。代码审查可以让一个 Agent 做静态分析,一个 Agent 检查测试覆盖率,一个 Agent 审查命名和结构,最后汇总 Agent 等待三个结果全部返回再生成综合意见。如果三个分析任务可以并行,汇总任务必须等待。调度器把并行和依赖分开处理:分析任务并发执行,汇总任务在 DAG 里等待,既保证了效率,又保证了结果完整性。
3.3 失败兜底:解析失败不要盲目重试
多 Agent 系统里最常见的失败不是模型超时,而是输出格式不合法。LLM 的稳定性再高,也不能保证每次都输出严格合法的 JSON。如果因为解析失败就自动重试,重试 3 次就等于多花 3 次模型调用的钱,而且不一定成功。
我的兜底策略是:先做 schema 校验。如果任何一个 Agent 返回的结果无法解析成预期的结构,不要直接重试,而是把这个任务标记为"转入人工队列",同时在日志里记录完整的原始输出。人工处理完以后,可以把任务重新放回队列。这样做的好处是,系统不会因为模型的一次抽风就陷入无意义的重试循环。
还有一点值得注意:依赖失败要区分对待。如果 DAG 中一个上游 Agent 失败,下游 Agent 可以做的选择有三个:继续执行(容忍部分失败)、挂起等待(重试上游)、整个任务终止转人工。这需要在调度层配置,而不是让 Agent 自己决定。我在代码审查场景里把所有失败都配置成"终止并转人工",因为一份不完整的审查报告没有意义。
4. 生产环境必须处理的三类坑:幂等、死锁与上下文污染
4.1 幂等性:重试不能双重执行
这是我在生产环境踩过最贵的坑。Agent 调用外部工具时,如果任务因为网络抖动触发重试,而工具本身没有幂等保护,就会出现重复发送短信、重复创建订单、重复扣款这类事故。
解决办法只有一个:让task_id透传到工具调用层。外部系统要根据task_id做幂等判断,比如用 RedisSETNX或者数据库唯一索引。调度器也要保证:同一个task_id在任意时刻最多只有一个执行实例。如果重试之前任务已经在执行中,新的执行请求必须被拒绝。
还要注意一个细节:LLM 的重试不是简单重放。同一个 prompt 输入两次,模型可能生成不同的内容,因为解码过程有随机性。所以幂等保护不能只靠"输入相同则输出相同"这个假设,必须在任务进入执行层之前就分配好task_id,并在所有日志、工具请求头、回调通知中携带它。
4.2 死锁:Agent 互相等待
多 Agent 协作最隐蔽的问题是死锁。比如 Agent A 在等 Agent B 的知识检索结果,Agent B 又在等 Agent A 的意图识别结果。如果调度器没有环检测,这两个任务会一直卡在waiting_dependency状态,worker 白白占着,队列永不推进。
我在项目里做了三道防线。第一,任务创建时做依赖环检测:如果一个任务直接或间接依赖自己,直接拒绝创建。第二,每个依赖关系设置最大等待时间,超过时间就把下游任务从"等待"改为"转人工"。第三,对重试队列做长度限制,防止因为无限重试导致新任务被长期堵在队列外。这三道防线同时生效之后,死锁问题基本没有再出现过。
还要警惕另一种"假死":不是真正的依赖环,而是重试次数太多,导致任务在队列里反复进出。看起来系统在运行,实际上有效处理能力趋近于零。最好给整个队列设置积压告警,一旦某个优先级的任务等待时间超过阈值,就触发降级策略。
4.3 上下文污染:并发写同一个 context
这是并发场景下的经典问题。很多人在早期会用一个全局 dict 来保存会话上下文,多个任务并发执行时,后写入的 context 会覆盖前面的,导致 Agent 回答张冠李戴。
解决方案非常明确:每个任务必须有独立的 context 实例,用task_id作为唯一 key。Agent 之间传递数据只通过 task payload,不通过任何全局状态。如果你用了记忆模块,记忆键也要带上task_id或用户会话 ID,不能做成一个全局共享的"大脑"。
我习惯把 context 设计成不可变对象:每次更新都生成新版本,而不是在原对象上修改。这样并发执行时不会互相干扰,而且方便追溯每个任务的上下文演变过程。调度器日志里要能够按task_id检索到完整的事件链路,包括入队时间、开始执行时间、每次重试原因、最终结果。用 OpenTelemetry 做 span 的话,把task_id和priority作为 attribute 打进去,排查问题的速度会快很多。
4.4 上线前检查清单
我把经验总结成一张清单,每次新接入一个 Agent 任务类型都会过一遍:
- 任务是否携带全局唯一的
task_id - 所有外部副作用操作是否有幂等键保护
- 是否配置了超时时间和最大重试次数
- 重试队列是否有长度限制和积压告警
- 任务之间的依赖关系是否做了环检测
- 每个任务是否使用独立 context,不读写全局状态
- 日志是否可按
task_id检索完整事件链路
如果全部满足,基本上不会出现调度层面的重大事故。如果有一项不满足,建议先补齐再上线。我见过太多团队在模型效果上反复调优,最后线上事故却出在"任务被重复执行"这种低级问题上。
5. 从调度走向治理:上线后的演进方向
5.1 先有指标,再有优化
ax调度器上线后的第一件事,不是优化模型,而是监控队列指标。我通常重点看四个:队列长度、各级任务平均等待时间、worker 利用率、失败任务分布。没有这些指标,任何优化都是拍脑袋。
队列长度能反映瞬时压力。等待时间能反映调度策略是否公平。worker 利用率能反映资源配置是否合理,如果 5 个 worker 长期只有 1 个在忙,说明并发上限设置得太保守;如果全部打满而且队列持续增长,就该扩容或者优化 Agent 响应时间。失败任务分布则能告诉你,是某个 Agent 特别容易超时,还是某种输入格式特别容易触发解析失败。
我一般用 Prometheus 加 Grafana 做指标可视化,但如果你不想引入新组件,先把结构化日志做好,配合日志检索也能解决大部分问题。调度器的日志一定要和业务日志打进同一个链路,否则线上排查时要两头对,非常痛苦。
5.2 审计与回放
当业务上需要回答"为什么这个任务被处理成这个样子"时,完整的状态流转记录就是审计依据。我要求每个任务从submit开始记录所有事件:进入队列时间、被取出的时间、开始执行时间、每次重试的原因和间隔、最终终态。这些记录存在一张独立的 task_event 表里,保留至少 30 天。
有了这些数据之后,你甚至可以做一个简单的回放工具:输入task_id,系统按时间线展示这个任务经历过的所有调度决策。哪个环节耗时最长、哪次重试浪费了 token、哪个 Agent 返回了非预期结构,一清二楚。这对于跨团队协作特别有用,算法团队可以拿着事件记录跟模型问题做关联分析。
5.3 渐进式引入:不要一开始就上重型框架
最后说一点我对工具选型的看法。团队在做 Agent 调度时,容易犯的错是一上来就引入分布式工作流引擎,结果被复杂的配置、部署和运维问题拖垮,连核心任务都没跑通。
我的建议是分三步走。第一步,用进程内队列加状态机,实现 2.3 节那样的最小调度器,跑通核心业务。第二步,把队列迁移到 Redis Stream,解决持久化和多实例问题。第三步,等任务类型和团队规模都上来了,再考虑独立的调度服务或成熟的工作流引擎。80% 的收益来自优先级、并发控制、超时重试这三件事,而不是框架本身。把这三件事做扎实,远比追求技术栈的先进性重要。
我个人在实际操作中的体会是,ax调度最难的其实不是写调度器,而是让团队成员统一承认"任务不能由 Agent 自己说了算"。一旦接受了这个前提,很多争论自然就消失了。最后再分享一个小技巧:把调度器的日志和业务日志打到同一个链路里,线上排查问题的速度至少快一半;如果再加一个按task_id维度的检索入口,效果更明显。希望这些实战内容对正在搭建 Agent 系统的你有帮助。