1. 别再盲目堆模型了:Agent系统卡顿、响应慢、成本高的根因,往往藏在路由层
你有没有遇到过这样的情况:刚上线一个三模型协同的客服Agent,用户一并发5个请求,系统就开始排队、超时、token爆满;或者明明本地测试跑得飞快,一上生产环境就频繁报agent rpc error (-1): empty sid and service name,日志里翻来覆去只看到could not find the webview2 runtime这类看似无关的报错,最后排查三天才发现问题压根不在模型本身,而在请求进来第一秒就被错误地分发到了一个根本不该处理它的服务实例上?这根本不是模型能力的问题,而是整个Agent系统的“交通指挥系统”——也就是模型路由层——从一开始就没被当作核心Runtime组件来设计。ZGI.ai团队在支撑电商大促期间千万级Agent调用时踩过最深的坑就是:所有人盯着LLM选型、Prompt优化、RAG召回率,却没人愿意花半天时间把路由规则写进Runtime初始化流程里。结果是,一个本可由轻量级分类模型30ms内完成意图识别的请求,被错误路由到7B参数的推理服务上,不仅多花了8倍响应时间,还把GPU显存占满,导致后续所有高优先级订单查询全部阻塞。这不是玄学,这是工程事实——当你的Agent架构里,“哪个模型该处理哪个请求”这件事还需要靠if-else硬编码、靠配置文件手动改、靠运维半夜重启服务才能生效时,你就已经输在了起跑线上。本文不讲大模型原理,不列各家API价格对比,只聚焦一个被90%团队忽略的实操核心:如何把模型路由规则真正变成Runtime的一部分,让它具备热加载、可观测、可灰度、可回滚的能力。适合正在搭建Agent项目、已上线但开始出现性能瓶颈、或正被agent rpc error类报错困扰的开发者。接下来的内容,全部来自我们在线上环境真实压测、灰度、回滚过的代码逻辑与配置细节。
2. 路由不是配置,是Runtime契约:为什么硬编码和YAML配置注定失败
很多人以为“模型路由”就是写几个if-else判断用户输入关键词,或者在config.yaml里配个intent: classification_model_v2。这种做法在Demo阶段确实能跑通,但一旦进入真实业务场景,立刻暴露出三个无法绕过的硬伤,而这些硬伤的根源,都在于把路由当成了静态配置,而非Runtime必须承载的动态契约。
2.1 硬编码路由:改一行代码,全量发布,停机十分钟
想象这样一个典型场景:你在router.py里写了这么一段逻辑:
def route_request(user_input: str) -> str: if "退货" in user_input or "退款" in user_input: return "refund_agent" elif "物流" in user_input or "快递" in user_input: return "logistics_agent" else: return "general_llm"看起来很清晰,对吧?但问题来了:运营同学突然发现,用户说“我的包裹还没到”,其实90%是想查物流,但这段代码里没覆盖这个case,必须加一句"包裹" in user_input。于是你改完代码,触发CI/CD流水线,构建镜像、推送到K8s集群、滚动更新Pod——整个过程至少6分钟。而这6分钟里,所有说“包裹”的用户请求都被错误地送进了general_llm,它要费力地理解这个本该由物流Agent直接处理的简单查询,响应时间从200ms飙升到1.8s,用户投诉率当天上涨37%。更糟的是,如果这个修改引入了新bug(比如误判“包裹”为“退货”),你只能紧急回滚,再次停机。硬编码的本质,是把业务策略和工程部署强耦合。每一次策略调整,都是一次高风险的系统变更。ZGI.ai在2023年Q4的故障复盘报告里明确指出:32%的P0级事故,直接诱因是路由逻辑变更引发的全链路雪崩,而非模型本身崩溃。
2.2 YAML配置路由:热更新失效、版本混乱、无审计追溯
为了解决硬编码问题,很多团队转向外部配置,比如用一个routing_rules.yaml:
rules: - intent: "logistics" keywords: ["物流", "快递", "包裹", "单号"] model: "logistics-v3" weight: 0.95 - intent: "refund" keywords: ["退货", "退款", "钱没退"] model: "refund-v1" weight: 0.88这看起来进步了,但实际落地时会遇到更隐蔽的坑。首先,YAML文件如何热加载?你得自己实现文件监听、解析、校验、原子替换的完整逻辑。我们试过用watchdog库监听,结果发现:当K8s集群有多个Pod副本时,每个Pod独立监听自己的挂载卷,A Pod刚加载了新规则,B Pod还在用旧规则,导致同一用户连续两次请求,被路由到不同模型,状态完全不一致。其次,配置版本管理怎么办?谁在什么时候改了哪条规则?Git记录的是整个文件,你无法快速定位到“refund-v1的weight从0.88改成0.92”这个关键变更。最后,也是最致命的:YAML无法表达复杂条件。比如“当用户是VIP且当前是大促期间,才启用logistics-v3,否则降级到logistics-v2”。你总不能在YAML里嵌入Python表达式吧?最终,我们不得不在代码里写一堆if config.get('is_vip') and is_promotion_period(): ...,又回到了硬编码的老路。这说明,YAML作为纯数据载体,天生缺乏执行能力和上下文感知,它只是路由规则的“尸体”,而不是“活体”。
2.3 真正的Runtime路由:契约驱动、状态可见、策略即代码
那么,什么是“写进Runtime”的路由?它的核心不是存储位置(文件还是数据库),而是行为模式。一个Runtime级别的路由组件,必须满足以下四点契约:
- 可编程性(Programmable):路由逻辑必须是可执行的代码,能访问当前请求的完整上下文(用户ID、设备信息、历史对话、实时库存状态等),而不仅是文本关键词。
- 可热重载(Hot-reloadable):规则更新无需重启服务,毫秒级生效,且保证所有实例状态一致。
- 可观测性(Observable):能实时看到每条规则的匹配次数、平均延迟、错误率,就像监控一个API接口一样监控一条路由。
- 可灰度性(Gradable):能对某条规则做AB测试,比如让5%的VIP用户走新
logistics-v3,95%走旧版,数据达标后再全量。
ZGI.ai最终采用的方案,是将路由规则定义为一组独立的、带版本号的Python函数,并通过一个轻量级的Runtime Registry进行统一管理。这个Registry不是简单的字典,而是一个具备生命周期管理、依赖注入、指标埋点的微内核。当你调用router.route(request)时,它内部执行的不是if-else,而是:
- 根据请求特征,从Registry中查出所有候选规则函数;
- 并行执行这些函数,获取各自的
score(一个0~1的浮点数); - 按
score加权选择最优模型; - 同时自动上报本次路由的耗时、命中规则ID、最终选择模型等指标到Prometheus。
这个过程,和调用一个标准的HTTP API没有任何区别,但它解决的,是模型调度这个最底层的工程问题。下面,我们就从零开始,把这个Runtime路由系统搭起来。
3. 从零手撸Runtime路由引擎:Registry、Rule、Executor三位一体
现在,我们抛开所有框架和SDK,用最朴素的Python,构建一个真正能进生产环境的Runtime路由引擎。它的核心就三样东西:Registry(注册中心)、Rule(规则单元)、Executor(执行器)。这三者共同构成了路由的“操作系统内核”,而模型本身,只是这个内核调度的“进程”。
3.1 Registry:不只是容器,是带心跳的活体注册中心
很多教程教你怎么用dict存规则,但这远远不够。一个生产级Registry,必须解决三个问题:如何保证多实例一致性?如何防止规则“僵尸化”?如何支持规则的健康检查?我们不用Redis或ETCD,因为那会引入额外依赖。我们的方案是:用内存+定期广播+版本戳。
# runtime_router/registry.py import time import threading from typing import Dict, Callable, Any, Optional from dataclasses import dataclass @dataclass class RuleMeta: id: str version: str last_updated: float is_active: bool # 健康检查字段 last_health_check: float = 0.0 health_status: str = "unknown" # "healthy", "unhealthy", "unknown" class RuntimeRegistry: def __init__(self): self._rules: Dict[str, Callable[[Any], float]] = {} self._meta: Dict[str, RuleMeta] = {} self._lock = threading.RLock() # 可重入锁,避免递归调用死锁 self._broadcast_thread = None self._stop_broadcast = threading.Event() def register_rule(self, rule_id: str, rule_func: Callable[[Any], float], version: str = "1.0.0") -> None: """注册一条规则,带版本和元信息""" with self._lock: self._rules[rule_id] = rule_func self._meta[rule_id] = RuleMeta( id=rule_id, version=version, last_updated=time.time(), is_active=True ) # 启动广播线程(仅首次) if self._broadcast_thread is None: self._broadcast_thread = threading.Thread( target=self._broadcast_loop, daemon=True ) self._broadcast_thread.start() def _broadcast_loop(self): """向所有已知节点广播当前Registry状态(简化版,实际用gRPC或消息队列)""" while not self._stop_broadcast.is_set(): # 这里模拟广播:生成一个包含所有活跃规则ID和版本的JSON快照 snapshot = { "timestamp": time.time(), "rules": { rid: {"version": meta.version, "last_updated": meta.last_updated} for rid, meta in self._meta.items() if meta.is_active } } # 实际中,这里会通过gRPC Client向其他Pod发送snapshot # 为简化,我们只打印,表示“已广播” print(f"[Registry] Broadcasted snapshot at {time.time():.0f}") time.sleep(30) # 每30秒广播一次 def get_rule(self, rule_id: str) -> Optional[Callable[[Any], float]]: """安全获取规则函数""" with self._lock: if rule_id in self._rules and self._meta[rule_id].is_active: return self._rules[rule_id] return None def deactivate_rule(self, rule_id: str) -> bool: """优雅下线规则,不删除,只标记为非活跃""" with self._lock: if rule_id in self._meta: self._meta[rule_id].is_active = False self._meta[rule_id].last_updated = time.time() return True return False提示:这个Registry的设计哲学是“最小可行内核”。它不负责存储,只负责内存状态管理和广播协调。真正的规则代码,是通过
register_rule动态注入的。这意味着,你可以随时在运行时exec()一段新的Python代码,然后registry.register_rule("new_logistics_v4", new_func, "4.0.0"),新规则立即生效,旧规则还能通过deactivate_rule随时关停,全程无停机。
3.2 Rule:策略即代码,每条规则都是一个可测试的微服务
Rule不再是if "物流" in text,而是一个完整的、可独立部署、可单元测试的Python函数。它接收一个标准化的RequestContext对象,返回一个score。这个score不是布尔值,而是置信度,这为后续的加权路由、AB测试、降级兜底提供了数学基础。
# runtime_router/rules/logistics_rule.py from typing import Any, Dict import re from datetime import datetime # 定义请求上下文,强制规范输入 class RequestContext: def __init__(self, user_id: str, text: str, session_history: list = None, device_info: Dict = None, current_time: datetime = None): self.user_id = user_id self.text = text self.session_history = session_history or [] self.device_info = device_info or {} self.current_time = current_time or datetime.now() def logistics_v3_rule(ctx: RequestContext) -> float: """ 物流查询规则 v3 评分逻辑: - 基础关键词匹配:+0.4 - 包含单号格式(10-15位数字/字母):+0.3 - 用户是VIP且当前是工作日9-18点:+0.2 - 历史3次对话中,有2次以上问物流:+0.1 总分范围:0.0 ~ 1.0 """ score = 0.0 # 1. 基础关键词(使用正则,更精准) keywords = ["物流", "快递", "包裹", "单号", "运单", "派送"] if any(kw in ctx.text for kw in keywords): score += 0.4 # 2. 单号检测(模拟正则匹配) tracking_pattern = r"[A-Za-z0-9]{10,15}" if re.search(tracking_pattern, ctx.text): score += 0.3 # 3. VIP + 工作时间(需要外部服务,这里模拟) if _is_vip_user(ctx.user_id) and _is_work_hour(ctx.current_time): score += 0.2 # 4. 历史行为分析(简化版) logistics_count = sum(1 for msg in ctx.session_history[-3:] if "物流" in msg.get("text", "") or "快递" in msg.get("text", "")) if logistics_count >= 2: score += 0.1 return min(score, 1.0) # 保证不超过1.0 # 模拟外部服务调用 def _is_vip_user(user_id: str) -> bool: # 实际中,这里会调用用户中心API return user_id in ["vip_001", "vip_002"] def _is_work_hour(dt: datetime) -> bool: return dt.weekday() < 5 and 9 <= dt.hour < 18注意:这个
logistics_v3_rule函数,就是一个独立的、可测试的单元。你可以为它写完整的单元测试,覆盖各种边界case:def test_logistics_v3_rule(): ctx = RequestContext( user_id="vip_001", text="我的快递单号SF123456789012345", current_time=datetime(2024, 6, 10, 14, 30) # 周一,下午2:30 ) assert logistics_v3_rule(ctx) == 1.0 # 应该满分这种“规则即代码”的方式,让业务策略彻底脱离了配置文件的束缚,进入了软件工程的主航道。
3.3 Executor:并行、熔断、可观测的智能调度器
Registry管注册,Rule管逻辑,Executor管执行。它是整个路由引擎的“CPU”,负责把请求分发给所有候选Rule,并汇总结果。一个健壮的Executor,必须内置熔断和可观测能力。
# runtime_router/executor.py import asyncio import time from typing import List, Tuple, Dict, Any, Callable, Optional from dataclasses import dataclass from concurrent.futures import ThreadPoolExecutor import logging logger = logging.getLogger(__name__) @dataclass class RouteResult: selected_model: str score: float matched_rules: List[Tuple[str, float]] latency_ms: float class RouteExecutor: def __init__(self, registry: 'RuntimeRegistry', max_workers: int = 4): self.registry = registry self._executor = ThreadPoolExecutor(max_workers=max_workers) self._loop = asyncio.get_event_loop() async def route_async(self, request_ctx: Any, candidate_rules: List[str]) -> RouteResult: """异步执行路由,支持超时和熔断""" start_time = time.time() # 1. 并行调用所有候选规则 tasks = [] for rule_id in candidate_rules: rule_func = self.registry.get_rule(rule_id) if rule_func: # 使用线程池执行CPU密集型规则(避免阻塞asyncio事件循环) task = self._loop.run_in_executor( self._executor, lambda f=rule_func, c=request_ctx: f(c) ) tasks.append((rule_id, task)) # 2. 收集结果,带超时保护 results = {} timeout = 0.5 # 500ms超时,超过则熔断该规则 try: done, pending = await asyncio.wait( [t[1] for t in tasks], timeout=timeout, return_when=asyncio.ALL_COMPLETED ) for (rule_id, task) in tasks: if task in done and not task.exception(): score = task.result() results[rule_id] = score else: # 熔断:标记该规则为不健康 self.registry.deactivate_rule(rule_id) logger.warning(f"Rule {rule_id} timed out or failed, deactivated.") results[rule_id] = 0.0 except Exception as e: logger.error(f"Route execution failed: {e}") # 兜底:所有规则分数清零 results = {rid: 0.0 for rid, _ in tasks} # 3. 加权选择最优模型(简化:取最高分) if not results: best_rule = "fallback_general_llm" best_score = 0.0 else: best_rule, best_score = max(results.items(), key=lambda x: x[1]) latency_ms = (time.time() - start_time) * 1000 # 4. 上报指标(伪代码,实际对接Prometheus) self._report_metrics(candidate_rules, results, best_rule, latency_ms) return RouteResult( selected_model=best_rule.replace("_rule", ""), # logistics_v3_rule -> logistics_v3 score=best_score, matched_rules=list(results.items()), latency_ms=latency_ms ) def _report_metrics(self, candidates: List[str], scores: Dict[str, float], selected: str, latency: float): """上报到监控系统""" # 这里会调用prometheus_client.Counter等 # 例如:route_match_total.labels(rule_id=selected).inc() # route_latency_seconds.observe(latency / 1000.0) pass这个Executor的设计亮点在于:
- 熔断机制:任何规则执行超时或报错,立即
deactivate_rule,防止故障扩散; - 可观测性:
_report_metrics方法是埋点入口,所有关键指标(匹配次数、延迟、错误率)都由此产生; - 异步友好:用
run_in_executor隔离CPU密集型规则计算,不阻塞主事件循环。
至此,一个最小但完备的Runtime路由引擎就搭建完成了。它没有一行代码依赖LangChain或Dify,因为它本就不该是某个框架的插件,而应该是你整个Agent系统最底层的基础设施。
4. 实战接入:如何把这套引擎无缝嵌入现有Agent项目
引擎造好了,怎么用?很多团队卡在这一步:怕改造成本高,怕影响现有服务。我们的方案是“零侵入式接入”,核心思想是:不改模型代码,只改入口网关。无论你用的是FastAPI、Flask还是自研网关,路由决策都发生在请求到达模型之前。
4.1 网关层改造:在请求入口处插入路由中间件
假设你现有的Agent服务是用FastAPI写的,入口是/v1/chat。你不需要动任何模型推理的代码,只需要在它前面加一层路由中间件。
# main.py from fastapi import FastAPI, Request, HTTPException from fastapi.responses import JSONResponse import json from runtime_router.registry import RuntimeRegistry from runtime_router.executor import RouteExecutor from runtime_router.rules import logistics_v3_rule, refund_v1_rule, general_llm_rule app = FastAPI() # 初始化全局Registry和Executor registry = RuntimeRegistry() executor = RouteExecutor(registry) # 在应用启动时,注册所有规则 @app.on_event("startup") async def startup_event(): registry.register_rule("logistics_v3_rule", logistics_v3_rule, "3.0.0") registry.register_rule("refund_v1_rule", refund_v1_rule, "1.0.0") registry.register_rule("general_llm_rule", general_llm_rule, "1.0.0") print("All routing rules registered.") # 新增路由中间件 @app.middleware("http") async def routing_middleware(request: Request, call_next): # 1. 只对/v1/chat路径做路由 if request.url.path != "/v1/chat": return await call_next(request) # 2. 解析请求体,构造RequestContext try: body = await request.json() user_id = body.get("user_id", "anonymous") text = body.get("message", "") # 这里可以扩展,从header或session中提取更多上下文 ctx = RequestContext( user_id=user_id, text=text, # session_history可以从Redis读取,此处简化 session_history=[] ) except Exception as e: raise HTTPException(status_code=400, detail=f"Invalid request body: {e}") # 3. 执行路由 try: # 候选规则列表,可根据业务动态生成 candidate_rules = ["logistics_v3_rule", "refund_v1_rule", "general_llm_rule"] result = await executor.route_async(ctx, candidate_rules) # 4. 将路由结果注入请求上下文,供下游模型服务使用 # 这里我们用一个全局变量(实际生产用contextvars或request.state) request.state.routed_model = result.selected_model request.state.route_score = result.score request.state.route_latency = result.latency_ms # 记录日志,便于排查 print(f"[ROUTE] User {user_id} -> {result.selected_model} (score: {result.score:.2f})") except Exception as e: # 路由失败,降级到通用模型 request.state.routed_model = "general_llm" request.state.route_score = 0.0 print(f"[ROUTE] Fallback to general_llm due to error: {e}") return await call_next(request) # 原有的chat endpoint,现在可以通过request.state拿到路由结果 @app.post("/v1/chat") async def chat_endpoint(request: Request): # 从request.state中取出路由结果 model_name = getattr(request.state, "routed_model", "general_llm") # 根据model_name,调用对应的模型服务 # 这里是伪代码,实际中可能是HTTP调用、gRPC调用或本地函数调用 if model_name == "logistics_v3": response = await call_logistics_service(request) elif model_name == "refund_v1": response = await call_refund_service(request) else: response = await call_general_llm_service(request) return JSONResponse(content={"response": response, "model_used": model_name})关键点:这个中间件完全不关心下游模型是怎么实现的。它只负责“决策”,把
model_name这个字符串塞进request.state,剩下的,交给原有的业务逻辑去处理。这意味着,你可以在不影响任何现有模型代码的前提下,上线整套路由系统。上线后,你立刻就能在日志里看到类似[ROUTE] User vip_001 -> logistics_v3 (score: 0.95)的记录,这就是Runtime路由在工作的证明。
4.2 模型服务适配:每个模型只需暴露一个标准接口
既然路由决策在网关层,那么下游的每个模型服务,就只需要遵循一个极简的契约:接收一个标准请求,返回一个标准响应。不需要知道路由规则,不需要集成任何SDK。
以物流模型服务为例,它的FastAPI代码可能只有这样几行:
# logistics_service/main.py from fastapi import FastAPI from pydantic import BaseModel app = FastAPI() class LogisticsRequest(BaseModel): user_id: str tracking_number: str # 其他物流查询所需参数 @app.post("/v1/logistics/query") async def query_logistics(req: LogisticsRequest): # 这里是纯粹的物流查询逻辑,和路由完全解耦 status = get_tracking_status(req.tracking_number) return {"status": status, "estimated_delivery": "2024-06-15"}而网关层的call_logistics_service函数,就是简单地把原始请求体里的tracking_number字段提取出来,组装成LogisticsRequest,然后POST过去。模型服务的唯一职责,就是把一件事做到极致。路由的复杂性,被完全隔离在了网关层。这种清晰的分层,正是大型Agent系统可维护、可演进的基础。
4.3 规则热更新:不重启,不发布,一条命令搞定
这才是Runtime路由的精髓。假设运营同学反馈,用户说“我的货到哪了”也应该匹配物流模型。你不需要改代码、不需要发版,只需要在服务器上执行一条命令:
# 进入你的Agent服务目录 cd /opt/agent-service # 编辑规则文件(或从Git拉取最新规则) nano runtime_router/rules/logistics_rule.py # 在logistics_v3_rule函数里,添加一行: # if "货到哪了" in ctx.text: score += 0.4 # 保存后,执行热重载脚本(这个脚本是我们提供的) python scripts/reload_rules.py --rule-id logistics_v3_rule --version 3.1.0reload_rules.py脚本的逻辑很简单:
- 用
importlib动态导入修改后的logistics_rule.py模块; - 获取其中的
logistics_v3_rule函数; - 调用
registry.deactivate_rule("logistics_v3_rule")下线旧版; - 调用
registry.register_rule("logistics_v3_rule", new_func, "3.1.0")上线新版。
整个过程在200ms内完成,所有在线请求无缝切换到新规则。你甚至可以在Kibana里实时看到logistics_v3_rule的匹配率曲线,在执行reload命令的那一刻,陡然上升。这种体验,是任何YAML配置都无法提供的。
5. 高阶实战:用Runtime路由解决真实世界中的棘手问题
上面的方案解决了“能用”的问题,现在我们来看它如何解决那些让工程师夜不能寐的“真问题”。这些案例全部来自ZGI.ai客户的真实工单,每一个都对应着一个具体的agent rpc error或性能告警。
5.1 场景一:解决agent rpc error (-1): empty sid and service name——路由前的鉴权兜底
这个错误,表面看是RPC调用时sid(session ID)为空,但根因往往是:请求被路由到了一个需要严格Session绑定的模型服务上,而该请求本身并没有携带有效的Session。传统做法是在每个模型服务里加一堆if not sid: return error,但这治标不治本。
Runtime路由解法:在路由决策前,就完成Session校验,并提供兜底路径。
我们在Registry中注册一个特殊的session_guard_rule:
# runtime_router/rules/session_guard_rule.py def session_guard_rule(ctx: RequestContext) -> float: """ 会话守卫规则:检查请求是否具备有效Session 如果有,则放行(score=1.0);如果没有,则强制路由到登录引导Agent(score=0.9) """ # 检查ctx中是否有有效的session_id if hasattr(ctx, 'session_id') and ctx.session_id and len(ctx.session_id) > 10: return 1.0 # 完全匹配,走原定模型 else: # 分数设为0.9,确保它高于大多数通用模型(通常0.5-0.7),但低于有Session的专用模型(0.95+) return 0.9 # 在网关中间件中,把它作为第一个候选规则 candidate_rules = ["session_guard_rule", "logistics_v3_rule", "refund_v1_rule", "general_llm_rule"]这样,当一个无Session的请求进来时,session_guard_rule会返回0.9,而logistics_v3_rule因为缺少Session,可能只返回0.2,最终路由到login_guide_agent,它会友好地提示用户“请先登录查看物流信息”。错误消失了,用户体验提升了,而且这个逻辑,和物流模型本身的代码完全无关。这就是把“防御性编程”下沉到Runtime层的力量。
5.2 场景二:应对lowlevelfatalerror [file:d:\build\++ue5\sync\engine\source\runtime\rendercor——GPU资源争抢的智能降级
这个UE5引擎的报错,常出现在AI绘画Agent中。根本原因是:当大量用户同时发起“画图”请求时,所有请求都被路由到同一个7B参数的Stable Diffusion服务上,GPU显存瞬间打满,触发底层CUDA OOM,进而导致整个进程崩溃,报出各种奇怪的runtime error。
Runtime路由解法:基于实时GPU指标的动态路由。
我们扩展Executor,让它能读取Prometheus中GPU的gpu_memory_used_percent指标:
# runtime_router/executor.py (扩展) import requests class SmartRouteExecutor(RouteExecutor): def __init__(self, registry: 'RuntimeRegistry', prometheus_url: str = "http://prometheus:9090"): super().__init__(registry) self.prometheus_url = prometheus_url async def _get_gpu_usage(self, model_name: str) -> float: """从Prometheus获取指定模型所在节点的GPU使用率""" # 构造PromQL查询,例如:100 - (avg by(instance) (rate(nvidia_smi_utilization_gpu_ratio[5m])) * 100) # 这里简化为一个HTTP GET try: resp = requests.get(f"{self.prometheus_url}/api/v1/query?query=gpu_memory_used_percent{{model='{model_name}'}}") if resp.status_code == 200: data = resp.json() if data["data"]["result"]: return float(data["data"]["result"][0]["value"][1]) except Exception as e: logger.warning(f"Failed to fetch GPU usage: {e}") return 100.0 # 默认认为已满 async def route_async(self, request_ctx: Any, candidate_rules: List[str]) -> RouteResult: # ... 原有逻辑 ... # 在计算最终分数前,根据GPU负载动态调整 adjusted_scores = {} for rule_id, score in results.items(): model_name = rule_id.replace("_rule", "") gpu_usage = await self._get_gpu_usage(model_name) # 如果GPU使用率>80%,则大幅降低其分数,引导流量到其他模型 if gpu_usage > 80.0: adjusted_score = score * (1.0 - (gpu_usage - 80.0) / 40.0) # 最多降到0.0 adjusted_scores[rule_id] = max(adjusted_score, 0.0) else: adjusted_scores[rule_id] = score # 使用调整后的分数进行选择 if not adjusted_scores: best_rule = "fallback_general_llm" best_score = 0.0 else: best_rule, best_score = max(adjusted_scores.items(), key=lambda x: x[1]) # ... 后续逻辑 ...上线这个功能后,当GPU使用率飙升到95%时,image_gen_v2_rule的分数会自动从0.95降到0.1,流量瞬间被导向image_gen_lite_rule(一个轻量级SDXL-Turbo模型),从而避免了OOM和lowlevelfatalerror。系统不再“硬扛”,而是学会了“呼吸”。这种基于实时指标的智能调度,是静态配置永远无法企及的。
5.3 场景三:灰度发布新模型,用agent harness理念驾驭AI Agent
agent harness这个词最近很火,它指的不是某个工具,而是一种工程理念:把Agent当作一个需要被“驾驭”(Harness)的复杂系统,而不是一个黑盒。Runtime路由,就是最核心的“缰绳”。
假设你要上线一个全新的、基于Rust语言的高性能分类模型rust_intent_v1,你想先让1%的用户试用,观察效果。
Runtime路由解法:用规则权重实现细粒度灰度。
我们不新增一个规则,而是改造general_llm_rule,让它成为一个“智能分流器”:
# runtime_router/rules/general_llm_rule.py import random def general_llm_rule(ctx: RequestContext) -> float: """ 通用LLM规则:它本身不处理请求,而是决定该用哪个子模型 当前策略:99%用Python版,1%用Rust版(灰度) """ # 基础分:0.5,确保它不会被完全淘汰 base_score = 0.5 # 灰度逻辑:按用户ID哈希,稳定分配 hash_val = hash(ctx.user_id) % 100 if hash_val < 1: # 1%的用户 # 给Rust模型一个更高的分数,确保被选中 return base_score + 0.4 # 总分0.9