上个月凌晨两点,我被值班告警电话叫醒。客户那边的第四个业务方接入了我们私有化部署的AI Agent平台,晚高峰还没来得及扩容,大模型推理服务的超时率直接飙到百分之十几,前端的人工助手一直在转圈,用户那边已经开始投诉。那晚排查到凌晨四点,翻日志看得头疼,最后定位到的根因其实特别朴素:我们的任务调度是单一队列直连模型,所有请求排一条队,所有任务挤同一个推理资源池。
那之后我重新设计了整套调度链路,核心思路就是标题里说的“两级任务分流架构”。第一级在入口网关按请求类别拆通道,第二级在调度层按模型路径和资源池分派任务。这篇文章把整个设计思路、落地细节、参数配置和排障经验完整写出来,给正在做企业级AI Agent私有化部署、或者已经接了多个业务方但稳定性扛不住的团队做个参考。不画大饼,全部是可落地的做法。
1. 单队列直连模型,为什么上线三周就开始崩
先说说我们原来的架构有多简单。客户端请求进来,网关直接丢给一个全局的任务队列,队列后面挂着几个推理实例,谁空闲谁取任务,任务执行完原路返回。这套架构在内部测试、单业务方联调阶段一点问题没有,一天几千次调用稳稳的。但一旦接入多个业务方,请求类型混杂起来,稳定性立刻雪崩。我复盘那次事故,根因可以归纳成三个词:队列饥荒、故障放大、资源争用。
1.1 从一次凌晨的事故复盘说起
当天晚上其实是有预兆的。下午五点左右,新增业务方上线了批量文档摘要功能,单批任务就要跑几十个文档,每个文档生成摘要前还要先做文本切分和重新排序。这类长任务单个耗时接近四十秒。到了六点半,前台客服问答的短任务也进高峰了,但短任务排在长任务后面,FIFO队列里一动不能动。九点左右,其中一个推理实例因长上下文任务持续占用显存出现OOM,实例直接崩了,但我们的负载均衡并不知道,后续请求还是往那个实例上发,超时一个接一个。
更惨的是,超时触发了应用层的自动重试。重试请求再次进入队列尾端排队,整个链路进入恶性循环:长任务越积越多,短任务越等越长,重试请求又进一步加大了队列深度。到晚上十一点,P95响应时间已经超过了三千毫秒,客户后台的超时告警刷了满屏。
1.2 三个根因,得分开看
第一个是队列饥荒。单一先进先出队列天然对短任务不友好,长任务只要排进队列,后面所有短任务都必须等它执行完,哪怕短任务只需要一秒。这就好比你先在超市排队窗口里放了一个买两百件商品的人,你买瓶水也得跟在他后面慢慢挪。现实中我们大量的问答型交互请求——比如查个工单状态、问个知识库条目——理想响应时间应该在一两秒内,却被后面排队的文档分析任务拖到不可用。
第二个是故障放大。单队列直连模式下,消费端对推理实例的“健康认识”是模糊的。一个实例是否OOM、是否卡在长序列生成上、是否半死不活地接受连接但响应极慢,负载均衡都感知不到。只要没有熔断机制,请求仍按轮询分发到故障实例,本来一个实例宕机只应影响该实例的并发量,却因为重试和超时扩散成了整体可用性归零。
第三个是资源争用。大模型推理有典型的显存和计算资源非线性使用特征。一个一万token上下文的生成任务占用的KVCache和计算量,可能是普通问答任务的好几十倍。把这两种任务放进同一个资源池,会发生极端情况:某实例正好在处理一个大上下文任务,同时又有十来个普通问答任务打过来,显存被打满,算力被挤占,小任务全部排队。这本质上是把不同执行特征的任务强行混部了。
当时也想过第二个方向:加机器。但加了机器后问题只是从“一台机器承载全部流量”变成“三台机器承载全部流量”,一旦流量再涨,同样的雪崩还会重演。所以真正该改的不是机器数量,而是任务分流的逻辑。也就有了后面的“两级分流”。
2. 第一级分流:入口网关先按请求类别拆开通道
两级分流架构里,第一级负责的是“分类”,也就是把请求拆进不同通道;第二级负责的是“调度”,也就是决定任务真正落在哪个模型路径和资源池上。第一级的设计原则是:只管归类和隔离,不做调度、不做重试、不做复杂逻辑。
2.1 分流维度怎么选:按来源、按类型、还是按紧急度
真正落地时,你会发现分类维度有非常多的选法。按业务来源分最简单:客户A的流量进A队列,客户B的流量进B队列,互不干扰。这个维度的好处是故障隔离最干净——某客户发起大批量任务,最坏情况只是他自己那组队列堵住,不会影响其他客户。按任务类型分则更贴近执行特征:问答类短请求、批量生成类中等请求、长文档分析类重请求,分别走不同队列。按紧急度分则照顾SLA:实时交互请求优先级最高,离线任务最低。
我的实践结论是:这三个维度都要用到,但不需要堆出十几个队列。队列数量太多,运维和监控成本会指数增加,故障排查也更难。我最终落地的方案是5个通道:实时交互通道、业务任务通道、批量离线通道、内部测试通道、高优重试通道。每个通道在Redis Stream里对应一个独立的Stream Key,入口网关根据请求头里的标签决定写入哪个Key。
分类标签从哪里来?我要求各业务方在接入时统一上报三个字段:source(业务来源)、bizType(实时/任务/离线)、priority(高/普通/低)。网关侧做白名单校验,不在白名单里的请求一律落到普通通道。这样可以防止某个业务方乱传标签、把高优通道塞满。
2.2 入口网关的实现:标签打点加Stream写入
我们的网关是自研的轻量Java服务,基于Netty实现,没有挂Spring Cloud Gateway全家桶,因为内部很多存量系统还是走自研RPC框架。核心逻辑其实就两件事:解析请求标签,然后按标签把任务原样写入对应队列。这里给出一个简化版的伪代码:
public void onRequest(AgentRequest request) { // 1. 解析并校验标签 RequestTag tag = request.getTag(); if (!tagValidator.isValid(tag)) { tag = RequestTag.createDefault(request.getSource()); } // 2. 根据标签挑选队列Key String queueKey = queueRouter.route(tag); // 实时交互 -> agent:stream:realtime // 业务任务 -> agent:stream:task // 批量离线 -> agent:stream:batch // 内部测试 -> agent:stream:intra // 高优重试 -> agent:stream:retry // 3. 生成幂等ID并写入Stream String requestId = IdGenerator.generate(request); redisTemplate.opsForStream().add( queueKey, Map.of("requestId", requestId, "payload", request.getBody(), "source", tag.getSource(), "bizType", tag.getBizType()) ); // 4. 先返回“已受理”,真正的执行由调度器异步完成 response.ok(requestId); }这里有一个容易踩的坑:网关层不要自动重试。早期版本里我在网关层加了失败重试逻辑,结果第一级重试和第二级重试叠加,线上出现大量重复执行。后来我把重试全部收敛到调度器,网关只负责写入队列,队列写入失败才允许重试一次且要求幂等ID不变。这样语义就清晰了:网关保证“任务受理成功”,调度器保证“任务最终被执行”。
第一级分流还有一个容易被忽略的收益:它可以做入口流控。比如实时交互通道的队列长度一旦超过5000,网关直接拒绝新增写入并返回503和Retry-After头。这个处理放在第一级比放在模型侧好得多,因为它可以在请求还没到推理资源之前就拦住,避免无谓的排队等待和资源占用。
3. 第二级分流:调度器决定任务去哪条模型路径和资源池
如果说第一级分流是“把不同性质的水引到不同管道”,第二级分流就是“决定每一股水进哪个水轮机”。调度器是这套两级架构的核心大脑,负责从各通道取出任务,结合任务的元数据、模型能力、资源池负载和优先级,做出分派决策。
3.1 调度器怎么判断任务该走哪条模型路径
模型路径这个概念是实际部署中很快会遇到的。私有化环境里不太可能只部署一个大模型服务,通常会有2到3种规格的模型服务同时在线。比如:轻量模型服务处理简单的意图识别、抽取式问答;中等模型服务处理函数调用、代码生成、结构化输出;重量模型服务处理长文档分析、复杂推理、多轮规划。即便只有一种模型,推理服务通常也支持不同的并发配置和超时参数,需要在调度层区分对待。
调度器拿到任务后,先看通道来源:实时通道的任务默认走轻量路径,批量通道的任务默认走重量路径。然后看任务自身的元数据:是否带图片输入、是否声明了长上下文、是否需要外部工具调用。大概的执行流如下:
def dispatch_task(task): path = None if task.channel == "realtime": if task.need_tool_call: path = "medium" elif task.estimated_cost < 0.3: path = "fast" else: path = "medium" elif task.channel == "batch": path = "heavy" else: path = "medium" pool = resource_manager.get_pool(path) if pool.is_full(): # 非高优任务等待,高优任务可以挤占预留容量 pool.wait_or_preempt(task) else: pool.submit(task)这层逻辑听起来简单,但落地时比较容易出问题的是“任务预估成本”。一开始我们简单用channel来定路径,但后来发现同一通道下,任务内容差异极大。一个小问答可能只有200 token,一个复杂文档问答却有2万 token。后来我们引入了一个前置的轻量分析器,用规则引擎估算任务输入长度、输出长度、是否包含长文本列表,把任务打上small/medium/large三个档次,再决定模型路径。这个分析器本身不能太重,否则预估的时间比执行还长,那就得不偿失了。
3.2 资源池划分和并发闸门:给每个池设一道硬上限
资源池的本质是一组具有相似执行特征的推理实例集合。也就是说不是单个模型实例,而是K8s Deployment下的一组副本,或者一组vLLM服务实例。每个资源池必须有独立的并发上限,不能让所有任务同时涌入同一批实例。
我当时的划分方式是这样的:
| 资源池 | 适用任务 | 副本数 | 单池最大并发 | 超时阈值 |
|---|---|---|---|---|
| fast-pool | 短问答、简单抽取 | 3 | 30 | 3s |
| medium-pool | 函数调用、结构化输出 | 4 | 20 | 15s |
| heavy-pool | 长文档、复杂推理 | 2 | 8 | 60s |
| fallback-pool | 降级兜底、规则回复 | 1 | 50 | 0.5s |
并发闸门用信号量实现。调度器向资源池提交任务前先获取信号量,获取不到就等待或拒绝。这一步非常重要,因为它隔离了不同任务之间的相互影响。重活池即使全满,fast-pool依然能正常处理短问答,互不干扰。
这里我踩过一个很深的坑:并发上限设得太高。第一次上线时fast-pool并发上限拍脑袋设成了50,觉得反正实例多、负载低。结果50路并发同时请求轻量模型,显存和GPU利用率双双打满,单次推理时间从500毫秒涨到2秒。后来我把并发上限压到30,P95反而降了40%。这个经验是:并发上限设置要参考实测指标,不是越大越好,要根据模型服务的单实例并发能力和显存占用反推。比较稳妥的做法是先按“单实例并发可容忍上限”的70%配置,观察几天后再微调。
3.3 短任务插队和长任务隔离:优先级不能是无限优先级
第二级分流除了分派路径,还要处理优先级。实时交互通道的任务应该被优先执行,但不能无限制插队,否则离线任务永远饿死。我们这里用了一个相对简单的权重调度策略:实时通道任务和批量通道任务的权重配比默认是5:1,且实时通道每连续插队20次后,强制调度至少1个批量任务。
代码上可以这样实现:
realtime_count = 0 while True: if realtime_count < 20: task = get_task("realtime") if task: dispatch(task) realtime_count += 1 continue task = get_task("batch") or get_task("realtime") if task: dispatch(task) realtime_count = 0这个策略保证了批量任务永远不会被完全饿死,也保证了实时任务的延迟不会因为批量任务堆积而失控。实测下来比严格的优先级队列好维护得多,因为你不必处理复杂的优先级反转问题。
另一个容易犯的错误是“一条逻辑链路里的多个任务被拆散到不同资源池”。比如一次多轮对话,用户连续问了五个问题,我们生成了五个独立的任务,调度器把三个分到fast-pool、两个分到medium-pool,最终回复顺序可能乱序,上下文衔接也容易出错。解决方法是给每个逻辑链路分配一个routeKey,调度器遇到routeKey相同的多个任务时,优先把它们分到同一个资源池的同一批实例上,以保持上下文局部性。
4. 两级之间的衔接:队列、状态机和幂等控制
两级分流架构跑起来之后,新增的复杂度就是“链路变长了”。请求从网关进队列,调度器取出后再入执行队列,执行器最终调用模型并回写结果。这个链路里面任何一个环节都可能失败,所以两级之间的衔接设计就显得格外重要。
4.1 队列选型:Redis Stream、RabbitMQ 还是 Kafka
任务队列的选型,我实测对比过三种方案。Redis Stream最简单,部署零成本,天然支持消费者组,任务量在每天几百万级别以内完全够用。RabbitMQ路由灵活,成熟稳定,但集群维护成本略高,而且业务量大了之后管理界面会比较乱。Kafka吞吐量最大,适合做离线分析和事件溯源,但Agent调度场景更偏在线实时,Kafka的一大问题是任务积压时offset滞后很难观测,不如Redis Stream的Pending List直观。
我的选型结论是:单机房场景,优先Redis Stream,理由有三。第一,私有化部署环境里Redis几乎是标配,不用额外引入中间件;第二,Redis Stream天然支持消费者组的负载均衡,调度器多副本部署时可以直接复用;第三,Stream的Pending List可以很方便地看到哪些任务长时间未处理,这对排查“任务卡住”类问题非常关键。
队列Key的设计上,每个通道独立一个Key,可以配合Redis的过期策略做队列积压告警。我用INFO命令每分钟采集一次队列长度,超过阈值就触发告警。这个指标在设计监控大盘时非常重要。
4.2 任务状态机:为什么任务状态必须独立于业务状态
任务被调度器取出后,会进入一个明确的执行状态机。我们用了六个状态:pending(待调度)、dispatched(已分派)、running(执行中)、success(成功)、failed(失败)、timeout(超时)。状态保存到MySQL里,任务执行的全过程都围绕状态机流转。
为什么非要引入状态机?因为一个任务可能经过网关、调度器、执行器三个组件,任何一个组件崩溃,任务都可能处于一个“没有归属方”的状态。状态机独立于业务状态存在,就是为了让恢复逻辑有据可依。调度器每隔一段时间扫描一次dispatched状态且超过N分钟没有推进的任务,把它们重新入队并更新状态为pending,这样可以自愈。
状态机的表结构设计也很简单:
CREATE TABLE agent_task ( request_id VARCHAR(64) PRIMARY KEY, route_key VARCHAR(32), channel VARCHAR(16), status VARCHAR(16), model_path VARCHAR(16), retry_count INT DEFAULT 0, created_at DATETIME, updated_at DATETIME );我强烈建议在更新状态时使用乐观锁,用version字段防止两个调度实例同时更新同一个任务导致的状态覆盖。这一步在测试环境可能看不出来,但在生产环境一旦遇到竞态就是很隐蔽的线上bug。
4.3 超时、重试和幂等控制的边界到底在哪
超时这个参数,不同模型路径应该不同。fast-pool的任务3秒超时,heavy-pool的任务60秒超时,这个阈值一定不能全局统一。早期我犯过的错是把所有任务都设成30秒超时,结果fast-pool的任务等待重试比实际执行还慢。后来改成按模型路径设置阈值,同时增加“快速失败”机制:如果任务在队列里等待的时间已经超过它的总超时预算,就直接标记为timeout,不再执行。
重试要收敛。网关层只处理队列写入失败的重试,调度器层处理分派失败的重试,执行器层不重试。每层重试次数上限为1到2次,并且使用指数退避。否则多个组件同时重试,故障恢复后会出现成倍的压力冲击模型服务。我见过一次重试风暴,故障结束后三个组件同时向模型服务发送重试请求,直接让实例再次OOM,这个教训非常深刻。
幂等控制用的是requestId加routeKey的组合,在状态机的request_id字段上建唯一索引。调用外部工具(比如发送短信、创建工单)之前,执行器检查一次request_id是否已经执行过工具调用,避免重复副作用。
5. 稳定性兜底:熔断、降级和过载保护怎么配
有了两级分流,稳定性的地基算是打牢了,但它只能解决“流量混合”的问题,不能解决“模型实例故障”的问题。真正把稳定性撑住的,还得是一套完整的兜底组合拳。
5.1 推理实例健康检查和熔断阈值怎么定
每个推理实例必须有自己的健康探活机制。我们用的是主动探测:每10秒向模型服务发送一个极小的ping请求,同时记录错误率。如果连续三个探测周期内失败率超过20%,调度器就把该实例从路由表里摘除。摘除后的实例进入冷却期,2分钟后再放回,放回时先分配少量流量观察是否恢复,也就是半开探测。
熔断阈值上,我建议以“单位时间失败率”作为主要判定指标,而不是原始错误次数。比如设定最近20个请求中失败超过10个就熔断,这样对小流量实例更公平。大流量实例可能出现瞬时抖动,所以可以加一个熔断最小请求数门槛,比如至少累计30个请求后才允许触发熔断,防止只有几个请求时误熔断。
5.2 降级链:从大模型到规则兜底
私有化部署环境里,模型服务不可用的时候,不能直接给用户报错,需要有降级方案。我们的降级链从重到轻分了四层:模型服务正常执行,模型服务降级到备用小模型,小模型也不可用时降级到预置的规则回复,最后才是返回明确错误。每一层降级触发条件都写入监控系统,保证降级动作可回溯。
比如实时交互通道的任务,如果大模型超时,调度器会把请求路由到fallback-pool的规则引擎。规则引擎可能不能回答很复杂的问题,但能给出“系统繁忙,请稍后再试”的兜底响应,用户感知是“延迟”,而不是“不可用”。批量离线通道的任务更适合进入重试队列,等模型服务恢复后再补跑。
5.3 优雅降载和队列积压保护
过载保护的核心是“尽早拒绝”,而不是“扛到最后”。我们给每个队列设置长度上限,比如实时通道最多5000条,批量通道最多20000条,超过上限的请求在网关层直接返回503和Retry-After。这样做的原因是,先进队列后端的压力只会随积压越来越大,模型服务的负载率一旦接近100%,单任务延迟会指数级恶化,最终谁的服务都保不住。
另外一件值得做的事是给推理服务本身加保护。vLLM类的服务可以在初始化是设置最大并发请求数和最大输入序列长度。我把最大输入序列长度压到模型原生支持长度的70%,事实证明这个操作能显著降低显存开销和OOM概率。队列积压保护则是对流量做削峰填谷,配合离线通道的任务进入夜间补跑队列,可以在资源空闲时段消化积压。
6. 上线实测数据,以及三个差点让我放弃的疑难杂症
架构改完后,系统逐渐稳定。但上线过程中踩的坑也不少,我挑几个典型的记录一下,希望对后来的人有参考价值。
6.1 监控指标里看到的真实收益
这套两级分流架构上线稳定运行一个月后,我拉了一次数据对比。峰值时段超时率从之前的11.8%降到了1.1%,P95响应时间从3200毫秒降到了860毫秒。更重要的是,故障恢复时间从分钟级降低到了秒级。之前单实例OOM需要人工介入才能恢复,现在熔断机制会在几十秒内自动摘除故障实例,客户端无感。
再到最核心的一个指标:在线率。分流量之前,某个大客户一个月内经历了4次可用性波动,其中两次是排队拥塞导致的间接不可用。分流量之后,连续两个月可用性都在99.9%以上。这个数字,在私有化部署环境里确实来之不易。
6.2 疑难杂症一:任务卡在“已分派”状态,排查了两天
上线后第一次遇到严重故障,是任务大量卡在dispatched状态不动。开始以为是调度器问题,但调度器日志显示任务已提交给执行器。接着查执行器,执行器日志却显示从未收到新任务。两边日志对不上,出了问题。
排查发现是Redis Stream的消费者组负载均衡选项在作怪。调度器有多个副本,某个副本重启后,Redis把尚未ACK的消息重新分配给存活消费者,这个重平衡过程配合消费者日志报错,导致任务实际是在调度器重启过程中丢失了。最终解决方案是在调度器启动时增加恢复流程:扫描所有dispatched状态超过3分钟的任务,把它们重新标记为pending并重新入队。同时,调度器不再在本地缓存任何任务上下文,所有信息都从Redis和MySQL读取,保证重启后状态可恢复。
6.3 疑难杂症二:推理实例内存缓慢增长,直到OOM
另一个疑难杂症是某个模型实例的显存使用率随运行时间缓慢增长,最终在运行6小时后OOM。排查发现这和任务调度没有直接关系,而是模型服务本身对长上下文任务的缓存没有及时清理。长任务留下的KVCache占用了显存,短任务无法复用,时间一长就爆了。
这个问题的处理方式比较直接:给长任务单独划分资源池,也就是前面表格里的heavy-pool。短任务永远不进入heavy-pool,heavy-pool的实例定期重启清理缓存。这样长任务的显存碎片只在heavy-pool内积累,不影响fast-pool的稳定性。
6.4 疑难杂症三:重试风暴让故障恢复变成二次故障
前面提到过重试风暴,这次是真的遇上了。某个晚上模型服务出现了两分钟的网络抖动,网关层、调度器层、执行器层同时触发了重试,抖动恢复后,三股重试流量叠加在一起,直接把模型服务再次压垮。
修复措施有两步:第一,把重试权全部收归调度器,网关层和调度器层分别只允许一次失败重试,执行器层彻底不重试;第二,给重试任务加一个“重试优先级”标签,重试任务只能从高优重试通道消费,并且该通道的单实例并发被压到很低,防止重试流量对正常流量造成冲击。
最后说点个人感受。这套架构不是一步到位的,最初我也试图一步到位把网关、调度、队列全部重新设计,结果漏洞百出。后来老实了,先把单队列改成通道分离,再逐步引入调度器和熔断降级。每改动一个小点,就在监控上观察一到两周。这样迭代下来,整个系统的稳定性是稳步上升的,而不是推倒重来的“爆炸式提升”。如果正在看这篇文章的你也打算做类似改造,我的建议是:先把流量分类看清楚,再动手写代码;分类对了,后面的调度设计就顺了。