做AI应用这一年多我最大的体会不是模型进步太快而是“底层执行链先崩了”。去年我们把AI客服灰度上线顺手打通了订单侧的数据链路结果有一天凌晨三点任务在没有人工确认的情况下被连续执行了几十次内部代号AX。产品在群里问“为什么一条指令会被派生这么多次”我盯着日志半天说不出话。事后复盘我们缺的既不是prompt也不是模型能力而是一层能把“规划、执行、复核、兜底”全部串起来的执行基础设施。这个东西我们内部叫AX调度现在团队和几个做Agent的朋友群里“ax调度”几乎成了Agent Execution调度的通用缩写谁家要是没个调度层开会都不好意思提。这篇文章不聊大模型原理就聊怎么把模型产出的“意图”变成可控、可观测、可回滚的任务流。内容包括ax调度的整体设计、状态机、超时与幂等、并发控制、一段能直接改的Python执行器骨架以及我调生产问题时遇到的高频坑。适合正在做多Agent编排、工具调用、或者准备把Agent接入业务系统的朋友尤其是那些已经被“工具重复执行”“任务卡死”“上下文丢失”折磨过的人。1. 起因为什么我会自己动手写一套ax调度先说一个容易被忽视的事实大模型本身不是不可控的真正不可控的是“模型调用工具之后工具到底怎么执行”。LLM输出一段JSON里面写“调用库存服务扣减10件”这一步只是生成了一句话。从这句话到数据库里的UPDATE真正生效中间隔着HTTP超时、权限校验、并发冲突、重复请求。任何一个环节出问题后果都直接落在真实业务上。我当时面临的问题主要有三个。第一个是工具失控。模型有时候会重复生成同一动作比如用户说“帮我把订单状态改成已退款”模型可能规划了两次“修改订单状态”的动作结果退款被提交两次。这不是模型故意使坏而是多轮上下文里它已经忘了上一步做过什么。第二个是执行时间不可控。模型思考很快但工具调用链路很长如果同步等待整条链跑完再返回用户早就不耐烦了。第三个是流程不可观测。一旦出问题日志是散的排查靠猜谁也不知道当前Agent到底执行到哪一步、还有几个子动作没做。这三个问题指向同一个答案不能把模型的“决策”直接当成“执行”必须把执行动作从主流程里拆出去交给一个专门组件来管。这个组件就是ax调度。它的职责可以概括成一句话模型只负责决定“做什么”ax调度负责决定“什么时候做、怎么做、做失败怎么办”。1.1 AI工具调用真正的问题不是“模型笨”而是“没人管”我见过团队在模型层面较劲反复调temperature、调few-shot示例就为了让模型不要重复调用工具。说实话效果有限。因为重复调用很多时候不是模型推理错而是执行层没有做幂等。模型生成两次相同参数的工具调用在它的“认知”里是两次独立的尝试在业务系统里却是两次扣款。这种时候靠prompt是压不住的。你得有一个执行层能在运行时识别“这个工具调用和上一次是同一个语义意图”直接把第二次请求拦下来或者返回上一次的结果。这就是幂等控制属于调度层的基本功而不是模型能力问题。还有一个更微妙的点模型调用工具是“大概率正确”也就是它有可能调用一个并不存在的工具或者传一个格式错误的参数。如果主流程是同步的只要有一环出错整个请求就失败用户只能重新发起。更麻烦的是某些工具有副作用失败了但系统部分状态已经改掉了你又不能简单重发。这些都需要调度层通过状态、重试、补偿机制来处理不是一个try-catch能解决的。1.2 ax调度的定位和普通任务调度有什么区别不少人第一次听到“ax调度”会想这不是咱公司那套定时任务平台能干的事吗确实像但细看差很多。普通的任务调度比如crontab或者xxl-job这类核心模型是“时间触发、执行任务、失败重试”。任务是预先定义好的参数是确定的执行完就结束不太需要考虑上下文。ax调度不一样。它面对的任务是模型动态生成的任务之间还有先后依赖关系。一个根任务可以拆出多个子任务某些子任务必须等前一个成功之后才能执行某些子任务可以并行。子任务失败后不是简单重跑而是可能要根据失败原因决定是重试、跳过、还是回滚整个根任务。我用两张表对比一下普通任务调度和ax调度的差异。维度普通任务调度ax调度触发方式定时触发或外部调用模型意图触发动态拆解任务内容固定脚本或接口模型生成的工具调用参数不确定执行链单体任务无关联父子任务、依赖关系、条件分支失败处理固定重试次数按错误类型分级重试、跳过、人工介入上下文不需要必须保留完整执行轨迹幂等要求低极高涉及真实业务副作用可观测性简单状态即可需要全链路trace、状态机、审计日志我常说一句话普通调度的核心是“准时”ax调度的核心是“靠谱”。前者关心任务有没有在正确时间跑后者关心任务跑完之后系统是不是还在一个确定的状态里。2. 整体设计逻辑把一次Agent执行变成可推进的流程我们设计ax调度时最先讨论的不是技术选型而是“一次Agent执行到底应该被抽象成什么”。最后定下来的模型非常朴素一次执行就是一棵树。树的根节点是用户发起的原始意图比如“帮我调一下A仓库和B仓库的库存差异”。模型接收到这个意图后可能拆出几个子节点“查询A仓库存”“查询B仓库存”“对比计算”“生成调拨单”“推送审核”。其中前三个可以并行后面两个必须按顺序执行。每一个节点我们称之为一个执行单元对应一次具体的工具调用。所有执行单元都挂在同一个执行上下文里。这个上下文记录了父任务是谁、子任务有哪些、当前节点状态、执行输入输出、重试次数、取消标志。说白了就是把一个Agent的“脑子”里想过的事完整落成一张可查询的表。这样用户问“刚才那个任务到底跑完没有”你不需要看日志直接查这张表就行。2.1 核心抽象根任务、子任务与执行树在ax调度的数据模型里一张任务表是核心。字段粗略包括任务ID全局唯一ID用UUID或者雪花ID都行。根任务ID标记这棵树从哪来。父任务ID标记当前节点属于谁。任务类型对应具体工具名。任务状态见下面的状态机。输入参数模型生成的JSON参数。执行结果成功后的返回值或者失败后的错误码。重试次数当前已经重试了几次最大允许几次。创建时间、执行时间、完成时间。这张表不仅是存储还是整个系统的“真相来源”。调度器推进状态时读它执行器写它人工后台也看它。所有操作都以任务ID为线索能把一棵树的执行历史完整捞出来。有了执行树这个抽象很多麻烦事就变得清晰了。比如取消操作用户说“等一下我不要调拨了”那就不需要遍历代码去停进程只需要把根任务状态置为CANCELLED所有还在排队、尚未执行的子任务看到这个标志后直接跳过。已经在执行的子任务通过协程或上下文取消信号来终止具体机制下面会讲。2.2 状态机怎么设计才不会跑飞状态机是整个ax调度里我最重视的部分因为状态一旦乱了整个系统就变成“薛定谔的执行中”。我的经验是状态不能太少也不能太多。太少了区分不了“超时”和“失败”太多了团队自己都记不住。我们最终使用的状态集如下状态含义可流转到WAITING任务已创建等待执行RUNNING, CANCELLED, DEADRUNNING任务正在执行SUCCESS, FAILED, RETRY_WAIT, CANCELLED, DEADRETRY_WAIT执行失败等待重试RUNNING, DEADSUCCESS执行成功无FAILED执行失败且不再重试DEAD, 人工修复后重跑CANCELLED任务被取消无DEAD死信需要人工介入可重新RUNNING这几个状态之间有一条铁律只有RUNNING状态下的执行器才允许产生业务副作用。其他状态都只是计划或结果。这点很重要如果没有这条约束调度器可能会在RETRY_WAIT状态下重复下发导致业务重复执行。另外我特别加了一个DEAD状态。很多人不喜欢死信觉得意味着系统失败。但实际上在AI工具调用场景里有些错误是你预设的重试规则处理不了的比如参数完全不符合业务规则。这种情况下最危险的不是报错而是“静默失败”或者“盲目重试”。所以我宁可把这类任务丢进DEAD让运营人员人工看一眼也不想让它自己无脑跑到天荒地老。2.3 为什么用异步事件驱动而不是同步顺序执行最开始我们做过一版同步编排逻辑很简单Agent调调度器调度器直接调工具拿到结果再让模型判断下一步。这套方案在demo里跑得很流畅一上生产就傻了。原因倒不复杂同步方式下整个链路的总计时间等于模型思考时间加所有工具调用时间相加。模型思考3秒工具调用2秒再来几个分支用户那边就是长达十几秒的无响应白屏。更致命的是同步模式下任何一环超时整个请求就挂在那边既不能取消也不能重试中间断点。有一次某个下游接口slow了客户端直接断开调度器这边还傻傻等着返回值整个执行链就悬在半空。后来我们改成异步事件驱动。核心思想是调度器只负责把任务推入队列并更新状态执行器从队列里取任务、执行、把结果作为事件发回。调度器收到事件后再决定是推进下一个子任务、重试当前节点还是标记整个根任务失败。这样一个模型会话的“思考”和“执行”不再强绑定用户体验提升一大截。异步化带来的额外好处是天然削峰。比如早上九点大量用户提问每个问题拆出三五个子任务队列会自然堆积执行器按自己的消费速度慢慢跑。不会因为瞬时并发太高直接把下游系统打挂。3. 关键细节超时、幂等、锁与单飞这一部分是ax调度能不能在生产环境活下去的分水岭。前两章讲的是“架子”这章讲的是“里子”。3.1 超时到底设多少才算既有体验又防失控给工具调用设超时是所有人都知道要做的事但设多少合适很多人拍脑袋。我的建议是分三层来看。第一层是模型调用超时。模型不参与工具执行但参与工具链的下一步规划。这一层我一般留30秒到60秒因为模型偶尔会在长上下文下思考很久。第二层是单次工具调用超时。如果是内部API一般10秒到15秒足够如果是外部第三方接口我建议不要超过20秒因为外部依赖更容易拖尾。第三层是根任务整体超时。这个取决于业务场景比如库存同步这种重操作我设30分钟如果是查天气这种轻量查询2分钟就顶天了。这里有个容易踩坑的点只在HTTP客户端加超时是不够的。如果工具内部发起了一条数据库查询或者调用另一个内部服务HTTP层超时只能断掉当前连接的等待但下游可能还在继续执行。更安全的方式是使用带超时传递的调用链比如在Python里用asyncio.timeout包住整个业务逻辑同时要求下游数据库或服务自身也设置语句超时。给个具体数字参考吧最近我在一个Go项目里习惯用context.WithTimeoutDB层设置statement_timeout5s外部HTTP调用timeout15s模型服务timeout45s。这样每一层都有限制不会出现“HTTP超时了但数据库还在跑”的尴尬局面。3.2 幂等性不是“加个唯一索引”就完事先讲一个真实事故。我们用唯一索引防重复给任务表加了“请求ID”唯一键以为万事大吉。结果线上还是出现了两次扣款排查到最后发现扣款服务是在事务提交之前先发了消息通知而通知里的业务操作没走幂等判断等于同一个请求在系统里走了两条不同的通道。在ax调度里幂等性至少要做好三张表。第一张是请求ID表。每次模型决定调用一个工具时调度器就生成一个全局唯一的request_id比如toolname-uuid。这个ID会在整个调用链中透传。第二张是幂等键表。request_id是唯一键表的字段包括幂等键、工具名、参数哈希、状态、首次执行时间、最后执行时间、返回结果。执行器在真正调用工具之前先往这张表写一条“PROCESSING”的记录工具调用成功后把结果写进这条记录并置为“SUCCESS”工具调用失败置为“FAILED”。如果同一个request_id再次闯进来直接查这张表状态是SUCCESS就返回老结果状态是FAILED就按失败逻辑处理。第三张是变更审计表。每次工具执行产生的业务变更、状态流转、返回值都追加记录。这张表不是直接用于防止重复而是让事后排查有据可查。特别是模型可能出现幻觉你总要知道上一次它到底是拿什么参数去调用工具的。核心原则是幂等控制必须发生在真正产生业务副作用之前的“最后一道门”而不是在入口处草草判断一下就觉得安全了。3.3 并发锁与请求合并同一会话不允许重复执行链有了幂等表能挡住“同一次重试”造成的重复。但还有另一种情况用户刷新页面或连点两次按钮同一个Agent会话被创建了两个执行链它们可能拿着不同的request_id但业务意图是同一个“改订单地址”。这时候幂等表就拦不住了因为参数可能一样但request_id不同。解决这个问题的办法是会话级锁。我习惯在Redis里建一把锁键名为agent_lock:{session_id}值存当前根任务ID加过期时间并带上一个随机token。当新任务准备执行时先尝试加锁。如果锁不存在说明当前没有执行链新任务获得执行权如果锁已存在后到的任务直接返回“当前会话已有任务在处理中”而不是排队等待。可能有人问为什么不直接让后到的任务排队等前一个跑完再跑我之前也这么设计过上线后发现体验很怪用户明明在追问一个新问题系统却在排队等上一个任务跑完给人一种“机器人卡住了”的感觉。后来改成“后到的任务直接拒绝并提示重试”反而更符合用户预期。加锁之后还要处理锁过期问题。模型工具链条一旦很长执行时间可能超过锁的过期时间锁自动消失新的执行链趁虚而入。我的做法是给锁加一个看门狗协程在任务运行期间自动续期续期到根任务最大执行时间后结束。另外锁的值要带一个随机token释放锁时校验token避免误删别人的锁。4. 实操从队列到执行器的最小闭环讲完设计来点能直接落地的。下面的方案我用了很久组件都是常见的东西不挑框架。4.1 最小闭环需要哪些零件一个能跑起来的ax调度至少需要五个零件任务存储、消息队列、调度器、执行器、管理后台。任务存储我用MySQL因为事务支持好任务状态流转需要强的数据一致性。消息队列我用Redis Stream它比普通List更适合这个场景因为支持消费者组和Ack机制。执行器是独立的Python服务也可以是多语言混合的这取决于工具生态。调度器是一个比较轻的服务它负责把根任务拆成子任务、把子任务推入队列、监听队列结果事件、更新任务状态。管理后台可以先用一套简单的Web界面用来查任务状态、手动取消、重新执行死信。我的建议是千万别一上来就上Kafka加多个微服务复杂度会把你的精力耗尽。先用一个进程把所有逻辑跑通等任务量真的上来了再拆分也不迟。4.2 一个可复现的执行器代码骨架这是我最想分享的部分。执行器的核心逻辑不复杂难的是把“超时、重试、心跳、结果上报”这些边界问题处理干净。我用Python写了一个最小骨架读起来更直观。# ax_worker.py import asyncio import json import time import uuid import redis.asyncio as redis STATUS_WAITING WAITING STATUS_RUNNING RUNNING STATUS_SUCCESS SUCCESS STATUS_FAILED FAILED STATUS_RETRY_WAIT RETRY_WAIT class AXWorker: def __init__(self, stream_nameax:queue): self.redis_client redis.from_url(redis://localhost:6379) self.stream_name stream_name async def heartbeat(self, task_id: str): # 每5秒上报一次心跳避免调度器误判为卡死 while True: await self.redis_client.hset( fax:heartbeat:{task_id}, last_beat, time.time() ) await asyncio.sleep(5) async def process_task(self, task): task_id task[task_id] tool_name task[tool_name] payload json.loads(task[payload]) await self.redis_client.hset( fax:task:{task_id}, status, STATUS_RUNNING ) # 启动心跳 hb_task asyncio.create_task(self.heartbeat(task_id)) try: # 这里替换成真实的工具调用注意超时 result await asyncio.wait_for( self.call_tool(tool_name, payload), timeout15 ) await self.redis_client.hset( fax:task:{task_id}, status, STATUS_SUCCESS ) await self.redis_client.hset( fax:task:{task_id}, result, json.dumps(result) ) # 上报结果事件让调度器推进下一个节点 await self.redis_client.xadd( ax:result, { task_id: task_id, status: STATUS_SUCCESS, result: json.dumps(result), } ) return True except asyncio.TimeoutError: await self.redis_client.hset( fax:task:{task_id}, status, STATUS_RETRY_WAIT ) await self.redis_client.xadd( ax:result, { task_id: task_id, status: STATUS_RETRY_WAIT, error: timeout, } ) return False except Exception as exc: await self.redis_client.hset( fax:task:{task_id}, status, STATUS_FAILED ) await self.redis_client.xadd( ax:result, { task_id: task_id, status: STATUS_FAILED, error: str(exc), } ) return False finally: hb_task.cancel() async def call_tool(self, tool_name: str, payload: dict): # 工具注册表按需扩展 if tool_name order_status: # 模拟一次内部接口调用 await asyncio.sleep(2) return {order_id: payload[order_id], status: refunded} raise ValueError(ftool {tool_name} not found) async def consume(self): # 从Redis Stream读取任务 while True: res await self.redis_client.xreadgroup( groupax_workers, consumerfworker-{uuid.uuid4().hex[:6]}, streams{self.stream_name: }, count1, block5000, ) if not res: continue for stream, messages in res: for msg_id, msg in messages: task msg ok await self.process_task(task) if ok: await self.redis_client.xack( self.stream_name, ax_workers, msg_id ) else: # 失败的不ACK让它进入Pending列表 pass这个骨架里有一个容易被忽略的细节失败的任务不ACK让它留在Redis Stream的Pending队列里。这样调度器可以通过排查Pending列表找出那些“执行了但没成功确认”的任务。有些团队在这里图省事直接ACK然后重新入队处理不好就会导致任务丢失。再说一下为什么不直接用redis里的BLPOP这种简单列表。因为Redis Stream的消费者组支持Ack机制任务被某个worker取走后调度器能明确知道它有没有处理完。普通List拿一个消息就少一个没有任何确认机制万一worker处理到一半挂了这条消息就彻底丢了这在业务工具调用场景是不可接受的。4.3 如何接进LangChain这类Agent框架很多人的Agent已经基于LangChain或LlamaIndex这类框架开发了。要接入ax调度不需要重写框架只需要把工具调用那一层替换成异步下发。简单说就是在原来工具函数内部不再直接执行业务代码而是创建一个调度任务把参数发给ax返回一个task_id。框架继续往下走的时候通过这个task_id去查状态。原来的同步返回值变成“轮询状态后的最终结果”。我给你一个变形示意原来的工具函数可能是这样def check_stock(warehouse_id: str): return stock_service.check(warehouse_id)接入ax调度后变成def check_stock(warehouse_id: str): task_id ax_scheduler.submit( tool_namecheck_stock, payload{warehouse_id: warehouse_id} ) return ax_scheduler.wait_for_result(task_id, timeout30)这段代码看起来变化不大但背后本质变了业务不再直接跟stock_service打交道而是经过调度层。调度层可以做超时、重试、幂等、并发控制而业务代码保持简单。如果你是纯用LangChain的AgentExecutor直接改写Tool的_run方法就好了。特别提醒一点改造之后工具调用耗时可能会有一定增加因为多了一层队列和状态查询。对于耗时敏感的场景不要把wait_for_result的等待时间设得太短否则Agent会误判工具失败。5. 问题清单调试ax调度遇到的5个真实坑这部分是我觉得全文最有价值的地方。以下问题全部来自生产环境每一个都让我加过班。5.1 任务一直停在RUNNING积压不消费有一天我看到任务表里几十个任务都停在RUNNING但队列里明明没有积压。排查后确认是某个worker进程卡死了但没有发心跳。当时的执行器没有做心跳上报调度器只能看到任务一直没完成却不知道worker是不是还活着。后来加上心跳机制调度器定期检查ax:heartbeat:{task_id}这个键。如果某个worker过期超过30秒没心跳就把任务重新投回队列并由调度器发告警。这里注意不能简单把任务重新投回就完事因为原本的worker可能还没死透只是卡在某个慢调用上。如果重新投递的worker和旧worker同时执行就会造成重复。所以重新投递之前必须通过幂等表确认旧worker没有真正产生业务变更。5.2 工具被调了两次幂等没拦住的怪异现场这是最让我头疼的一个问题。幂等表明明加了唯一索引结果还是出现了重复执行。排查后定位到原因是第一次请求已经写入幂等表但没有提交事务第二次请求看到的是旧快照判断“这个键不存在”于是又插入了一条。最终两条都提交了。解决办法是使用“条件插入 检查影响行数”的方式。先执行INSERT ... ON DUPLICATE KEY UPDATE如果影响行数为0说明这条幂等记录已经存在就不再执行工具。如果涉及事务隔离级别问题可以考虑把幂等表查询改成SELECT ... FOR UPDATE但这对性能有影响。更务实的方案是把幂等记录的写入和真正业务操作的执行放在同一个事务里确保“先占坑、再执行、同生共死”。5.3 模型幻觉式传参调度层应该怎么办LLM不是百分百按schema出参的它可能传一个不存在的仓库编号或者把单位“件”写成“箱”。如果工具直接执行轻则返回错误重则写错业务数据。在ax调度里我在执行器前面加了一层“参数校验网关”。每个工具注册时附带一份JSON Schema执行器拿到参数后先做schema校验不合格的不会调用真实工具而是把任务置为“QUARANTINE”状态相当于一个等待人工处理的特殊状态池。针对那些有破坏性的工具比如批量删除、修改价格、退款我额外加了“二次确认”规则。模型生成的调用会先进到确认队列由运营或用户点确认后才真正执行。这个机制看起来“慢”了一点但在真实业务里能救大命。5.4 任务风暴与队列积压的常见原因某个大促期间我们突然发现队列积压了上万条任务最夸张的一次延迟了20分钟。看了下罪魁祸首是某条根任务拆子任务时写了一个“死循环”模型没有判停反复调用同一个查询工具。一个用户会话就能生成几十个工具调用几十个用户就能把队列打满。要防这种问题只靠调度器限流是不够的。我给每个根任务设置了子任务数量上限比如最多50个子任务还给每个模型的“决策轮次”设置了上限。一旦超出调度器会强制终止这条执行链并且把根任务标记为DEAD。另一招是给执行器设置最大消费速率比如每秒最多消费20个任务防止下游系统被打爆。5.5 上下文用完就丢重试后恢复不了最后说一个比较隐蔽的坑。一个根任务里如果有个子任务失败了我们常规操作是重试这个子任务。但重试时需要重新组装模型上下文而原上下文里包含了前面几个子任务的执行结果。如果这些结果没有被持久化重试的时候就拿不到必要的上下文模型只能瞎猜。我的解决方案是把每个子任务的输入输出都持久化到任务表或者专门的trace表重试时直接从数据库读回之前的执行结果而不是指望内存里还留着。这个改动简单但对重试成功率提升非常明显。如果团队用的是内存队列一定得注意worker重启后任务状态还在不在不然一重启所有上下文都没了。说实话ax调度这层我踩过最大的坑不是锁写得不严谨也不是状态机设计失误而是“半吊子调度”。任务发出去没人管看着很轻量一出事就得跪着道歉。后来我把每条执行链都当成一条要收口的线段来看待必须有状态、有心跳、有兜底才算真正把AI工具调用放进了生产环境。我的个人体会是别太指望模型提供商未来会在SDK里帮你把调度、幂等、重试都内置好。模型迭代速度再快调度层的“可靠”依然是要自己动手垒的。调用一次模型很便宜出一次数据事故很贵这笔账越早算明白越好。希望这篇ax调度的拆解能帮正在做Agent工程化的你少走点弯路。