摘要:OpenRouter Fusion 证明了多模型融合能以更低成本达到接近前沿的质量。但企业如何在自有环境中实现类似的多模型融合?本文提供从零搭建企业级多模型融合路由系统的完整教程,涵盖并行调度、裁判模型融合、自动故障切换、成本监控四个核心模块,所有代码可直接运行。
目录
- 一、系统架构设计
- 二、环境准备
- 三、核心模块实现
- 四、融合策略配置
- 五、监控与成本追踪
- 六、生产环境部署
- 七、常见问题与排错
一、系统架构设计
1.1 整体架构
┌─────────────────────────────────────────────┐
│ 应用层 (Your App) │
├─────────────────────────────────────────────┤
│ FusionRouter (融合路由器) │
│ ┌───────────┐ ┌──────────┐ ┌────────────┐ │
│ │ 任务分类器 │ │ 模型选择器│ │ 融合策略器 │ │
│ └───────────┘ └──────────┘ └────────────┘ │
├─────────────────────────────────────────────┤
│ ParallelDispatcher (并行调度器) │
│ ┌──────┐ ┌──────┐ ┌──────┐ ┌──────┐ │
│ │模型A │ │模型B │ │模型C │ │模型D │ │
│ └──┬───┘ └──┬───┘ └──┬───┘ └──┬───┘ │
│ └────────┼────────┼────────┘ │
│ ▼ │
│ ┌──────────────────────┐ │
│ │ JudgeModel (裁判) │ │
│ └──────────────────────┘ │
├─────────────────────────────────────────────┤
│ 统一 API 层 (微元算力 weytoken) │
│ api.weytoken.com/v1 │
└─────────────────────────────────────────────┘
1.2 核心设计原则
- 模型无关:不绑定任何单一模型,模型切换不影响业务代码
- 故障自愈:单个模型不可用时自动降级,不影响整体输出
- 成本可控:按任务复杂度自动选择融合策略(预算/前沿/单模型)
- 审计可追溯:每次调用的模型组合、耗时、成本全程记录
二、环境准备
2.1 依赖安装
pip install openai asyncio python-dotenv
2.2 环境变量配置
# .env
WEYTOKEN_API_KEY=wt-your-api-key
WEYTOKEN_BASE_URL=https://api.weytoken.com/v1
2.3 模型清单
# models.py – 可用模型清单
AVAILABLE_MODELS = {
# 前沿模型
"fable5": "claude-fable-5",
"gpt55": "gpt-5.5",
"opus48": "claude-opus-4-8",
# 中端模型
"sonnet4": "claude-sonnet-4-20250514",
"deepseek_v4": "deepseek-v4-pro",
"kimi_k26": "kimi-k2.6",
"glm52": "glm-5.2",
# 预算模型
"gemini3_flash": "gemini-3-flash",
}
# 模型健康状态(模拟)
MODEL_HEALTH = {
"claude-fable-5": "unknown", # 可能被禁
"gpt-5.5": "healthy",
"gemini-3-flash": "healthy",
"kimi-k2.6": "healthy",
"deepseek-v4-pro": "healthy",
"glm-5.2": "healthy",
}
三、核心模块实现
3.1 并行调度器
# dispatcher.py
import asyncio
from openai import AsyncOpenAI
from typing import Optional
import os
from dotenv import load_dotenv
load_dotenv()
class ParallelDispatcher:
"""多模型并行调度器"""
def __init__(self):
self.client = AsyncOpenAI(
api_key=os.getenv("WEYTOKEN_API_KEY"),
base_url=os.getenv("WEYTOKEN_BASE_URL"),
timeout=300,
max_retries=1,
)
async def dispatch(self,
prompt: str,
models: list[str],
temperature: float = 0.7) –> list[dict]:
"""
并行分发给所有参团模型
Returns:
[{"model": str, "content": str, "status": "success"|"failed",
"latency": float, "error": str|None}, …]
"""
async def call_single(model: str) –> dict:
import time
start = time.time()
try:
response = await self.client.chat.completions.create(
model=model,
messages=[{"role": "user", "content": prompt}],
temperature=temperature,
)
return {
"model": model,
"content": response.choices[0].message.content,
"status": "success",
"latency": time.time() – start,
"usage": {
"input_tokens": response.usage.prompt_tokens,
"output_tokens": response.usage.completion_tokens,
},
"error": None,
}
except Exception as e:
return {
"model": model,
"content": None,
"status": "failed",
"latency": time.time() – start,
"usage": None,
"error": str(e),
}
tasks = [call_single(m) for m in models]
return await asyncio.gather(*tasks)
3.2 裁判模型融合器
# fusion.py
import json
from typing import Optional
class FusionJudge:
"""裁判模型:结构化分析 + 融合输出"""
FUSION_SYSTEM_PROMPT = """你是一个多模型融合裁判。你会收到多个模型对同一问题的回答。
请执行以下分析步骤:
1. 共识识别:列出所有模型达成共识的结论
2. 矛盾检测:标注互相矛盾的观点及各自的论据
3. 独到见解:找出某个模型独有的、其他模型未提及的洞见
4. 盲区标注:指出所有模型共同遗漏的重要维度
5. 综合输出:基于以上分析,撰写一份融合了所有模型优点的最终答案
输出格式(JSON):
{
"consensus": ["共识点1", "共识点2", …],
"contradictions": [
{"topic": "矛盾主题", "view_a": "观点A", "view_b": "观点B", "resolution": "推荐方案"}
],
"unique_insights": [
{"model": "模型名", "insight": "独到见解"}
],
"blind_spots": ["盲区1", "盲区2", …],
"final_answer": "综合最终答案"
}"""
def __init__(self, dispatcher):
self.dispatcher = dispatcher
async def analyze_and_fuse(self,
original_prompt: str,
worker_responses: list[dict],
judge_model: str) –> dict:
"""
分析所有 worker 回答并融合输出
Args:
original_prompt: 原始问题
worker_responses: 并行调度的结果列表
judge_model: 裁判模型名称
Returns:
融合结果字典
"""
successful = [r for r in worker_responses if r["status"] == "success"]
failed = [r for r in worker_responses if r["status"] == "failed"]
if not successful:
return {
"status": "all_failed",
"failed_models": [f["model"] for f in failed],
"final_answer": "所有模型均调用失败,请稍后重试。",
}
# 构建裁判 prompt
judge_prompt = self._build_judge_prompt(
original_prompt, successful, failed
)
# 调用裁判模型
from openai import OpenAI
client = OpenAI(
api_key=os.getenv("WEYTOKEN_API_KEY"),
base_url=os.getenv("WEYTOKEN_BASE_URL"),
)
response = client.chat.completions.create(
model=judge_model,
messages=[
{"role": "system", "content": self.FUSION_SYSTEM_PROMPT},
{"role": "user", "content": judge_prompt},
],
response_format={"type": "json_object"},
)
try:
result = json.loads(response.choices[0].message.content)
except json.JSONDecodeError:
result = {
"final_answer": response.choices[0].message.content,
"consensus": [],
"contradictions": [],
"unique_insights": [],
"blind_spots": [],
}
return {
"status": "success",
"worker_count": len(successful),
"failed_models": [f["model"] for f in failed] if failed else [],
"judge_model": judge_model,
**result,
}
def _build_judge_prompt(self, original: str,
successful: list[dict],
failed: list[dict]) –> str:
"""构建裁判模型的输入 prompt"""
responses_text = "\\n\\n—\\n\\n".join([
f"[模型: {r['model']}]\\n{r['content']}"
for r in successful
])
failed_text = ""
if failed:
failed_text = f"\\n注意:以下模型调用失败,不参与分析:{', '.join([f['model'] for f in failed])}"
return f"""原始问题:
{original}
以下是 {len(successful)} 个模型的回答:
{responses_text}
{failed_text}
请按照 JSON 格式输出分析结果。"""
3.3 融合路由器(主入口)
# router.py
import asyncio
import time
from typing import Optional
from dataclasses import dataclass, field
from enum import Enum
class TaskComplexity(Enum):
SIMPLE = "simple" # 单模型
MEDIUM = "medium" # 预算级融合
COMPLEX = "complex" # 前沿级融合
CRITICAL = "critical" # 最强融合
class FusionRouter:
"""企业级多模型融合路由器"""
# 预设融合策略
STRATEGIES = {
TaskComplexity.SIMPLE: {
"mode": "single",
"models": ["claude-sonnet-4-20250514"],
},
TaskComplexity.MEDIUM: {
"mode": "fusion",
"workers": ["gemini-3-flash", "kimi-k2.6", "deepseek-v4-pro"],
"judge": "claude-sonnet-4-20250514",
},
TaskComplexity.COMPLEX: {
"mode": "fusion",
"workers": ["claude-fable-5", "gpt-5.5"],
"judge": "claude-fable-5",
},
TaskComplexity.CRITICAL: {
"mode": "fusion",
"workers": ["claude-fable-5", "gpt-5.5", "glm-5.2"],
"judge": "claude-fable-5",
},
}
# 模型不可用时的自动降级策略
FALLBACK_CHAINS = {
"claude-fable-5": ["claude-opus-4-8", "claude-sonnet-4-20250514"],
"gpt-5.5": ["gpt-4o", "claude-sonnet-4-20250514"],
}
def __init__(self):
self.dispatcher = ParallelDispatcher()
self.judge = FusionJudge(self.dispatcher)
self.call_log = []
def classify_complexity(self, prompt: str) –> TaskComplexity:
"""简单的任务复杂度分类(可替换为 ML 分类器)"""
prompt_lower = prompt.lower()
# 关键任务关键词
critical_keywords = [
"安全审计", "安全漏洞", "安全审查", "security audit",
"合规", "compliance", "法律", "legal",
]
if any(kw in prompt_lower for kw in critical_keywords):
return TaskComplexity.CRITICAL
# 复杂任务关键词
complex_keywords = [
"架构", "architecture", "重构", "refactor",
"分析报告", "深度研究", "代码审查",
]
if any(kw in prompt_lower for kw in complex_keywords):
return TaskComplexity.COMPLEX
# 中等任务关键词
medium_keywords = [
"文档", "document", "解释", "explain",
"总结", "summary", "对比", "compare",
]
if any(kw in prompt_lower for kw in medium_keywords):
return TaskComplexity.MEDIUM
return TaskComplexity.SIMPLE
async def execute(self,
prompt: str,
strategy: Optional[TaskComplexity] = None,
force_single: Optional[str] = None) –> dict:
"""
执行融合路由
Args:
prompt: 用户请求
strategy: 融合策略(None = 自动分类)
force_single: 强制单模型(跳过融合)
"""
start_time = time.time()
# 强制单模型模式
if force_single:
result = await self.dispatcher.dispatch(
prompt, [force_single]
)
self._log(start_time, "single", [force_single], result)
return {
"mode": "single",
"model": force_single,
"content": result[0]["content"],
"latency": time.time() – start_time,
}
# 自动分类
if strategy is None:
strategy = self.classify_complexity(prompt)
config = self.STRATEGIES[strategy]
if config["mode"] == "single":
result = await self.dispatcher.dispatch(
prompt, config["models"]
)
self._log(start_time, "single", config["models"], result)
return {
"mode": "single",
"model": config["models"][0],
"content": result[0]["content"],
"latency": time.time() – start_time,
}
# 融合模式
worker_responses = await self.dispatcher.dispatch(
prompt, config["workers"]
)
fusion_result = await self.judge.analyze_and_fuse(
prompt, worker_responses, config["judge"]
)
self._log(start_time, "fusion",
config["workers"] + [config["judge"]],
worker_responses)
return {
"mode": "fusion",
"strategy": strategy.value,
"workers": config["workers"],
"judge": config["judge"],
"worker_count": fusion_result.get("worker_count", 0),
"failed_models": fusion_result.get("failed_models", []),
"content": fusion_result.get("final_answer", ""),
"analysis": {
"consensus": fusion_result.get("consensus", []),
"contradictions": fusion_result.get("contradictions", []),
"unique_insights": fusion_result.get("unique_insights", []),
"blind_spots": fusion_result.get("blind_spots", []),
},
"latency": time.time() – start_time,
}
def _log(self, start_time, mode, models, results):
self.call_log.append({
"timestamp": time.time(),
"mode": mode,
"models": models,
"latency": time.time() – start_time,
})
四、融合策略配置
4.1 国产优先策略
# 数据合规敏感场景的国产模型融合策略
CN_FIRST_STRATEGY = {
"mode": "fusion",
"workers": [
"glm-5.2", # MIT 开源,可自部署
"kimi-k2.6", # 长上下文优势
"deepseek-v4-pro", # 推理能力强
],
"judge": "glm-5.2",
"description": "全链路国产,数据不出境,合规零风险",
}
4.2 使用示例
# main.py
import asyncio
async def main():
router = FusionRouter()
# 示例 1:自动分类 – 简单任务
result = await router.execute(
"写一个 Python 函数,将 CSV 转换为 JSON"
)
print(f"[简单任务] 模式: {result['mode']}, 延迟: {result['latency']:.1f}s")
# 示例 2:自动分类 – 复杂任务(自动触发融合)
result = await router.execute(
"分析这个微服务架构的安全漏洞,给出修复方案"
)
print(f"[复杂任务] 模式: {result['mode']}, Worker: {result.get('worker_count', 0)} 个")
# 示例 3:手动指定策略 – 预算级融合
result = await router.execute(
"总结这份 50 页的行业报告",
strategy=TaskComplexity.MEDIUM
)
print(f"[预算融合] 延迟: {result['latency']:.1f}s")
print(f"盲区: {result['analysis']['blind_spots']}")
# 示例 4:强制单模型
result = await router.execute(
"写一个冒泡排序",
force_single="claude-sonnet-4-20250514"
)
print(f"[单模型] 模型: {result['model']}")
if __name__ == "__main__":
asyncio.run(main())
五、监控与成本追踪
5.1 调用日志与成本估算
# monitor.py
from dataclasses import dataclass
from typing import list
# 简化的成本模型(美元/百万 token)
PRICING = {
"claude-fable-5": {"input": 10, "output": 50},
"gpt-5.5": {"input": 10, "output": 50},
"claude-sonnet-4-20250514": {"input": 3, "output": 15},
"gemini-3-flash": {"input": 0.15, "output": 0.60},
"kimi-k2.6": {"input": 1, "output": 4},
"deepseek-v4-pro": {"input": 0.5, "output": 2},
"glm-5.2": {"input": 1, "output": 4},
}
class CostTracker:
"""成本追踪器"""
def __init__(self):
self.records = []
def record(self, mode: str, models: list[str],
usage: dict, latency: float):
"""记录一次调用的成本"""
cost = 0
for model in models:
if model in PRICING:
price = PRICING[model]
input_cost = usage.get("input_tokens", 0) * price["input"] / 1_000_000
output_cost = usage.get("output_tokens", 0) * price["output"] / 1_000_000
cost += input_cost + output_cost
self.records.append({
"mode": mode,
"models": models,
"cost": cost,
"latency": latency,
})
def summary(self):
"""打印成本摘要"""
total_cost = sum(r["cost"] for r in self.records)
avg_latency = sum(r["latency"] for r in self.records) / len(self.records)
fusion_calls = [r for r in self.records if r["mode"] == "fusion"]
single_calls = [r for r in self.records if r["mode"] == "single"]
print(f"""
=== 成本摘要 ===
总调用次数: {len(self.records)}
融合调用: {len(fusion_calls)} 次
单模型调用: {len(single_calls)} 次
总成本: ${total_cost:.4f}
平均延迟: {avg_latency:.1f}s
""")
六、生产环境部署
6.1 Docker 部署
# Dockerfile
FROM python:3.12-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install –no-cache-dir -r requirements.txt
COPY . .
CMD ["python", "main.py"]
6.2 环境变量
# docker-compose.yml
version: '3.8'
services:
fusion-router:
build: .
environment:
– WEYTOKEN_API_KEY=${WEYTOKEN_API_KEY}
– WEYTOKEN_BASE_URL=https://api.weytoken.com/v1
restart: unless–stopped
6.3 关于统一 API 层
本教程使用微元算力(weytoken)聚合平台 作为统一 API 层,原因是:
- 一个 Key 打通所有模型:无需为每个模型单独申请 Key、学习不同 API 格式
- OpenAI 兼容格式:所有模型统一调用方式,代码零改动
- 企业级基础设施:全链路审计日志、增值税专票、多租户隔离
- 数据安全:API 调用在企业自有服务端完成,数据不出管控范围
七、常见问题与排错
Q1:某个模型调用失败怎么办?
融合路由器会自动处理——失败的模型不影响整体输出,裁判模型基于成功的回答进行融合。failed_models 字段记录了失败详情。
Q2:融合延迟太高怎么办?
- 降低 worker 数量(2-3 个即可)
- 对于简单任务,使用 force_single 跳过融合
- 考虑异步调用,不阻塞主线程
Q3:成本如何控制?
- 简单任务自动走单模型(成本最低)
- 中等任务走预算级融合(成本约前沿模型的 50%)
- 设置每日成本上限
Q4:裁判模型的选择有什么建议?
- 预算优先:Sonnet 4(成本低,裁判能力足够)
- 质量优先:Fable 5(裁判能力最强)
- 国产优先:GLM-5.2(全链路国产,MIT 开源)

