单次模型调用更快,不代表任务完成得更快:低质量输出引发的重试、降级和人工修正都会进入端到端耗时。本文用 Node.js 实现一个可运行的动态路由器,并通过决策日志、影子流量和离线回放验证质量、延迟与成本。所有策略都保留静态路由开关,出现异常时可以快速回滚。

先定义真正需要优化的指标

模型选型经常只比较首字延迟 TTFT 或每秒输出 token 数,但业务关心的是“请求到达后,多久得到可接受的结果”。可以把任务耗时写成:

任务耗时 = 路由耗时 + 模型调用耗时 + 重试耗时 + 质量检查耗时 + 必要的人工修正时间

如果快速模型用 2 秒生成结果,却因为缺少字段而重试一次,最终可能比一次调用较强模型更慢。路由器因此不能只保存模型延迟,还要记录最终是否通过质量门槛、尝试次数和 token 用量。

指标用途注意事项
客户端端到端耗时判断用户实际等待时间应包含重试和回退
单次模型耗时定位具体模型性能不能单独代表任务完成时间
质量通过率衡量结果能否直接使用优先采用可确定执行的规则
输入、输出 token估算成本价格应通过独立配置维护
回退率发现路由是否过于激进高回退率会抵消快速模型优势

质量门槛应尽量贴近任务契约。例如信息抽取可检查必需字段,代码生成可运行测试,分类任务可验证枚举值。对于开放式写作,简单字符串规则只能检查格式,不能冒充完整的语义评价。

按任务、上下文与质量门槛路由

示例采用两类模型:FAST_MODEL 负责短文本摘要和抽取,STRONG_MODEL 负责高质量、代码及长上下文任务。规则不是永久结论,而是需要被真实流量持续验证的初始假设。

条件首选模型原因
quality=high强模型避免低质量结果触发二次调用
代码任务强模型通常需要更严格的正确性
估算上下文超过 8000 token强模型为长上下文保留独立策略
短摘要或抽取快模型任务约束明确,容易检查
其他任务强模型默认采取保守策略

上下文 token 数最好使用对应模型的 tokenizer。为了让示例不依赖特定 SDK,下面只用字符数除以 4 做路由估算;API 返回的真实 usage 才用于后续统计,不能把估算值当作计费数据。

实现可运行的 Node.js 路由器

下面程序基于 Node.js 20 的原生 fetch,调用真实存在的 OpenAI Chat Completions API。模型名称通过环境变量传入,避免把可能变化的模型列表写死。保存为 server.mjs

import http from "node:http";
import { appendFile, mkdir } from "node:fs/promises";
import { createHash, randomUUID } from "node:crypto";
import { performance } from "node:perf_hooks";

const API_KEY = process.env.OPENAI_API_KEY;
const API_BASE = process.env.OPENAI_API_BASE || "https://api.openai.com/v1";
const FAST = process.env.FAST_MODEL;
const STRONG = process.env.STRONG_MODEL;
const STATIC = process.env.STATIC_MODEL || STRONG;
const MODE = process.env.ROUTER_MODE || "dynamic";
const SHADOW = process.env.SHADOW_MODEL;
const SHADOW_RATE = Number(process.env.SHADOW_RATE || 0);
const CONFIG_VERSION = process.env.CONFIG_VERSION || "v1";
if (!API_KEY || !FAST || !STRONG) throw new Error("缺少 OPENAI_API_KEY、FAST_MODEL 或 STRONG_MODEL");
await mkdir("data", { recursive: true });

const writeLog = event => appendFile("data/decisions.jsonl", JSON.stringify({
  time: new Date().toISOString(), ...event
}) + "\n");
const estimateTokens = messages => Math.ceil(JSON.stringify(messages).length / 4);
const hash = value => createHash("sha256").update(value).digest("hex");
const sampled = id => parseInt(hash(id).slice(0, 8), 16) / 0xffffffff < SHADOW_RATE;
const acceptable = (text, terms = []) => typeof text === "string" && text.trim().length > 0 &&
  terms.every(term => text.toLowerCase().includes(String(term).toLowerCase()));

function choose(input) {
  if (MODE === "static") return { model: STATIC, reason: "static-kill-switch" };
  const tokens = estimateTokens(input.messages);
  if (input.quality === "high") return { model: STRONG, reason: "high-quality" };
  if (input.task === "code") return { model: STRONG, reason: "code-task" };
  if (tokens > 8000) return { model: STRONG, reason: "long-context" };
  if (["summarize", "extract"].includes(input.task) && tokens < 4000)
    return { model: FAST, reason: "short-structured-task" };
  return { model: STRONG, reason: "conservative-default" };
}

async function callModel(model, messages) {
  const started = performance.now();
  const controller = new AbortController();
  const timer = setTimeout(() => controller.abort(), 60000);
  try {
    const response = await fetch(`${API_BASE}/chat/completions`, {
      method: "POST",
      headers: { "Authorization": `Bearer ${API_KEY}`, "Content-Type": "application/json" },
      body: JSON.stringify({ model, messages, temperature: 0 }),
      signal: controller.signal
    });
    if (!response.ok) throw new Error(`API ${response.status}: ${await response.text()}`);
    const data = await response.json();
    return {
      model, text: data.choices?.[0]?.message?.content,
      usage: data.usage || {}, latency_ms: Math.round(performance.now() - started)
    };
  } finally {
    clearTimeout(timer);
  }
}

async function runShadow(id, input, primaryModel) {
  if (!SHADOW || SHADOW === primaryModel || !sampled(id)) return;
  try {
    const result = await callModel(SHADOW, input.messages);
    await writeLog({ type: "shadow", id, model: SHADOW,
      latency_ms: result.latency_ms, usage: result.usage,
      quality_passed: acceptable(result.text, input.mustInclude) });
  } catch (error) {
    await writeLog({ type: "shadow-error", id, model: SHADOW, error: error.message });
  }
}

async function readJson(req) {
  let raw = "";
  for await (const chunk of req) {
    raw += chunk;
    if (raw.length > 1_000_000) throw new Error("请求体过大");
  }
  return JSON.parse(raw);
}

const server = http.createServer(async (req, res) => {
  if (req.method !== "POST" || req.url !== "/chat") {
    res.writeHead(404).end(); return;
  }
  const started = performance.now();
  const id = randomUUID();
  try {
    const input = await readJson(req);
    if (!Array.isArray(input.messages)) throw new Error("messages 必须是数组");
    const decision = choose(input);
    await writeLog({ type: "decision", id, config: CONFIG_VERSION,
      task: input.task, quality: input.quality, model: decision.model,
      reason: decision.reason, estimated_tokens: estimateTokens(input.messages),
      prompt_hash: hash(JSON.stringify(input.messages)) });

    if (process.env.RECORD_REPLAY === "1") {
      await appendFile("data/replay.jsonl", JSON.stringify({ id, ...input }) + "\n");
    }

    const attempts = [];
    let result;
    try {
      result = await callModel(decision.model, input.messages);
      attempts.push(result);
    } catch (error) {
      attempts.push({ model: decision.model, error: error.message });
    }

    let passed = result && acceptable(result.text, input.mustInclude);
    if ((!result || !passed) && decision.model !== STRONG) {
      result = await callModel(STRONG, input.messages);
      attempts.push(result);
      passed = acceptable(result.text, input.mustInclude);
    }
    if (!result) throw new Error("所有模型调用均失败");

    const total = Math.round(performance.now() - started);
    await writeLog({ type: "completion", id, final_model: result.model,
      total_ms: total, quality_passed: passed,
      attempts: attempts.map(x => ({ model: x.model, latency_ms: x.latency_ms,
        usage: x.usage, error: x.error })) });
    void runShadow(id, input, decision.model);

    res.writeHead(200, { "Content-Type": "application/json; charset=utf-8" });
    res.end(JSON.stringify({ id, model: result.model, output: result.text,
      quality_passed: passed, total_ms: total, attempts: attempts.length }));
  } catch (error) {
    await writeLog({ type: "request-error", id, error: error.message });
    res.writeHead(400, { "Content-Type": "application/json; charset=utf-8" });
    res.end(JSON.stringify({ id, error: error.message }));
  }
});
server.listen(3000, () => console.log("router listening on :3000"));

启动与请求示例:

OPENAI_API_KEY=你的密钥 \
FAST_MODEL=你的快速模型 \
STRONG_MODEL=你的强模型 \
SHADOW_MODEL=待验证模型 \
SHADOW_RATE=0.05 node server.mjs

curl http://localhost:3000/chat -H 'content-type: application/json' -d '{
  "task":"extract",
  "quality":"normal",
  "mustInclude":["Node.js"],
  "messages":[{"role":"user","content":"只用一句话说明本文使用什么运行时:Node.js。"}]
}'

mustInclude 只是一个可执行的演示门槛。生产系统应按任务注册不同验证器,而不是让调用方随意决定全部质量标准。

用影子流量、回放和回滚验证策略

影子请求不影响主响应,却会产生真实调用成本。示例用请求 ID 的哈希稳定采样,并记录候选模型延迟、token 和质量检查结果。进程短暂退出可能丢失未完成的异步任务;生产环境应把影子任务投递到持久化队列,由独立 worker 执行。

开启 RECORD_REPLAY=1 后,程序会保存可回放请求。它可能包含敏感提示词,因此只能用于获得授权的流量,并应增加脱敏、访问控制、保留期限和加密。离线回放时,对相同数据集调用候选模型,再按任务聚合端到端耗时、通过率、回退率和 token。价格会变化,建议用独立价格配置把 usage.prompt_tokensusage.completion_tokens 换算为成本,不要把价格硬编码进路由器。

验证时应先确定不可退让的质量下限,再比较满足下限的方案。开放式任务可采用盲审或经过校准的评审流程,但不要把另一个模型的单次打分直接视为客观真值。还要避免只回放成功请求,否则会产生幸存者偏差。

上线顺序可以是:离线回放、1% 影子流量、小比例真实路由、逐步放量。每次修改 CONFIG_VERSION,这样日志才能回答“哪一版策略导致变化”。如果质量通过率下降、回退率上升或总耗时恶化,设置 ROUTER_MODE=static 并令 STATIC_MODEL 指向上一稳定模型,重启实例即可绕过全部动态规则。配置、代码和数据集版本也应一起保留,才能复现结论。

总结

  • 优化目标应是获得可接受结果的端到端任务耗时,而不是孤立的首字延迟或输出速度。
  • 路由决策至少要考虑任务类型、上下文长度和质量门槛,并把回退调用计入延迟与成本。
  • 决策日志应记录配置版本、路由原因、尝试次数、真实 token 用量和质量结果,同时避免直接记录敏感正文。
  • 影子流量适合观察线上分布,离线回放适合可重复比较;二者都需要一致的任务级质量验证器。
  • 动态路由必须保留静态模型开关。先定义质量下限,分阶段放量,并确保任何策略都能安全回滚。