构建具有双上游自动容灾的 OpenAI 兼容网关
- 作者

- 姓名
- Nino
- 职业
- Senior Tech Editor
在生产环境中,过度依赖单一的大语言模型(LLM)服务商是极其危险的。在业务高峰期,API 接口经常会遇到延迟飙升、频次受限(Rate Limits)或直接宕机的情况。例如,当 DeepSeek 官方 API 在凌晨 3 点发生短时波动时,后台运行的批量处理任务可能会大面积报错。虽然我们可以在应用层代码中编写复杂的捕获重试逻辑,但这会使业务代码变得臃肿且难以维护。
最优雅的解决方案是将容灾逻辑与业务逻辑彻底解耦。通过构建一个轻量级的、兼容 OpenAI 协议的 API 网关,我们可以将所有的大模型请求先发送到本地网关。一旦首选的主渠道发生故障,网关会自动将请求转发给备用渠道。对于上层应用而言,这个切换过程是完全无感知的,只需要在 OpenAI SDK 中修改 base_url 即可完成对接,无需改动任何核心业务代码。
尽管自行开发网关是一个极佳的技术实践,但在实际生产中,维护高可用架构、全球延迟优化以及多模型计费是一项不小的挑战。对于追求极致稳定且不希望承担运维成本的企业,像 n1n.ai 这样的托管式 API 聚合器提供了开箱即用的多通道容灾路由、统一账单管理以及毫秒级的故障自动切换服务。
自动容灾网关的架构设计
我们的自研网关主要围绕以下核心需求进行设计:
- OpenAI 协议兼容:提供标准
/v1/chat/completions接口,支持标准的请求体与请求头。 - 双上游自动容灾:当主渠道(如 SiliconFlow 接口)响应超时(超过 8 秒)或返回 5xx 服务端错误时,自动无缝切换至备用渠道(如 DeepSeek 官方 API)。
- 流式传输支持:完美支持 SSE(Server-Sent Events)协议,确保流式输出不会因为中转而产生卡顿或内容截断。
- 本地配额与日志:利用轻量级 SQLite 数据库,记录每次请求消耗的 Token 数量,并对用户密钥进行额度扣减与校验。
以下是该网关的请求生命周期示意图:
[客户端应用]
│ (OpenAI SDK / v1/chat/completions)
▼
[FastAPI 容灾网关]
│
├─► [主上游通道: SiliconFlow] (超时时间 < 8s?)
│ │ (若发生故障或超时)
│ ▼
└─► [备用上游通道: DeepSeek 官方]
│
▼
[SQLite 数据库] (记录 Token 消耗并校验可用额度)
200 行 Python 代码的完整实现
下面是基于 FastAPI、httpx(异步 HTTP 客户端)以及 aiosqlite(异步 SQLite 库)实现的完整网关代码。您可以直接将其部署在本地或云服务器中:
import json
import time
import asyncio
import logging
from typing import AsyncGenerator
from fastapi import FastAPI, Request, HTTPException, Depends
from fastapi.responses import StreamingResponse
import httpx
import aiosqlite
# 配置日志输出
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("llm-gateway")
app = FastAPI()
# 核心配置参数
DATABASE_PATH = "gateway.db"
PRIMARY_URL = "https://api.siliconflow.cn/v1/chat/completions"
PRIMARY_KEY = "sk-siliconflow-key-here"
FALLBACK_URL = "https://api.deepseek.com/v1/chat/completions"
FALLBACK_KEY = "sk-deepseek-key-here"
TIMEOUT_LIMIT = 8.0 # 超时阈值(秒)
# 初始化数据库表结构
async def init_db():
async with aiosqlite.connect(DATABASE_PATH) as db:
await db.execute("""
CREATE TABLE IF NOT EXISTS users (
api_key TEXT PRIMARY KEY,
token_quota INTEGER,
token_used INTEGER DEFAULT 0
)
""")
await db.execute("""
CREATE TABLE IF NOT EXISTS usage_logs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
api_key TEXT,
upstream TEXT,
prompt_tokens INTEGER,
completion_tokens INTEGER,
timestamp REAL
)
""")
# 预先插入一个本地测试 Key,赠送 50 万 Token 额度
await db.execute(
"INSERT OR IGNORE INTO users (api_key, token_quota) VALUES (?, ?)",
("sk-local-test-key", 500000)
)
await db.commit()
@app.on_event("startup")
async def startup_event():
await init_db()
# 校验 API Key 及其剩余额度
async def verify_user(request: Request) -> str:
auth_header = request.headers.get("Authorization")
if not auth_header or not auth_header.startswith("Bearer "):
raise HTTPException(status_code=401, detail="API Key 格式错误或缺失")
api_key = auth_header.split(" ")[1]
async with aiosqlite.connect(DATABASE_PATH) as db:
async with db.execute(
"SELECT token_quota, token_used FROM users WHERE api_key = ?",
(api_key,)
) as cursor:
row = await cursor.fetchone()
if not row:
raise HTTPException(status_code=401, detail="未授权的 API Key")
quota, used = row
if used >= quota:
raise HTTPException(status_code=429, detail="额度已耗尽")
return api_key
# 记录 Token 消耗并更新用户表
async def log_usage(api_key: str, upstream: str, prompt: int, completion: int):
async with aiosqlite.connect(DATABASE_PATH) as db:
await db.execute(
"""
INSERT INTO usage_logs (api_key, upstream, prompt_tokens, completion_tokens, timestamp)
VALUES (?, ?, ?, ?, ?)
""",
(api_key, upstream, prompt, completion, time.time())
)
await db.execute(
"UPDATE users SET token_used = token_used + ? WHERE api_key = ?",
(prompt + completion, api_key)
)
await db.commit()
# 异步迭代器:处理流式响应并实时解析 Token 数量
async def forward_stream(response: httpx.Response, api_key: str, upstream_name: str) -> AsyncGenerator[str, None]:
total_prompt = 0
total_completion = 0
try:
async for line in response.aiter_lines():
if not line.strip():
continue
yield f"{line}\n\n"
# 解析 SSE 数据流中的 usage 字段
if line.startswith("data: "):
data_str = line[6:]
if data_str.strip() == "[DONE]":
continue
try:
data_json = json.loads(data_str)
usage = data_json.get("usage")
if usage:
total_prompt = usage.get("prompt_tokens", 0)
total_completion = usage.get("completion_tokens", 0)
except json.JSONDecodeError:
pass
finally:
await response.aclose()
# 传输结束后写入日志
if total_prompt > 0 or total_completion > 0:
await log_usage(api_key, upstream_name, total_prompt, total_completion)
# 发送 HTTP 请求并捕获 5xx 错误
async def attempt_request(client: httpx.AsyncClient, url: str, key: str, payload: dict) -> httpx.Response:
headers = {
"Authorization": f"Bearer {key}",
"Content-Type": "application/json"
}
response = await client.post(url, json=payload, headers=headers, timeout=TIMEOUT_LIMIT)
if response.status_code >= 500:
raise httpx.HTTPStatusError("上游服务器返回 5xx 错误", request=response.request, response=response)
return response
@app.post("/v1/chat/completions")
async def chat_completions(request: Request, api_key: str = Depends(verify_user)):
payload = await request.json()
stream = payload.get("stream", False)
# 配置主备上游通道
upstreams = [
("SiliconFlow", PRIMARY_URL, PRIMARY_KEY),
("DeepSeek Official", FALLBACK_URL, FALLBACK_KEY)
]
async with httpx.AsyncClient() as client:
for name, url, key in upstreams:
try:
logger.info(f"正在尝试通过上游通道发送请求: {name}")
response = await attempt_request(client, url, key, payload)
if stream:
return StreamingResponse(
forward_stream(response, api_key, name),
media_type="text/event-stream"
)
else:
data = response.json()
usage = data.get("usage", {})
prompt = usage.get("prompt_tokens", 0)
completion = usage.get("completion_tokens", 0)
await log_usage(api_key, name, prompt, completion)
return data
except (httpx.RequestError, httpx.HTTPStatusError, asyncio.TimeoutError) as e:
logger.warning(f"上游通道 {name} 异常: {str(e)}。正在切换至下一个通道...")
continue
# 若所有上游通道均不可用,返回 502 错误
raise HTTPException(status_code=502, detail="所有配置的上游 LLM 通道均无法访问或已超时。")
生产环境落地必须要解决的痛点
虽然构建一个简单的代理网关非常容易,但在实际的高并发生产场景中,有几个非常微妙的坑需要特别注意。
1. 不同服务商 Token 统计字段不一致
不同的 LLM 服务商(如 SiliconFlow、DeepSeek 官方或 OpenAI)在流式输出(Streaming)时,返回 Token 统计信息的方式大相径庭。有些服务商会在最后一个 data: [DONE] 之前的包中携带 usage 信息;而有些服务商则需要显式在请求体中传入 stream_options: {"include_usage": true} 才会返回统计。如果代码没有做好兼容,网关记录的 Token 消耗就会始终为 0。
技术建议:在网关层建立严格的数据结构清洗机制。如果上游接口未返回 usage 字段,应在本地引入 tiktoken 库或对特定的开源分词器(Tokenizer)进行离线计算,从而估算 prompt 与 completion 的实际大小,确保计费与限流数据的准确。
2. 流式传输中途断开的尴尬局势
如果主上游在刚开始建立连接时就超时,我们的 try-except 捕获逻辑可以完美处理并切换到备用渠道。但是,如果主渠道已经成功响应,并且已经向客户端输出了前 50 个 Token,此时连接突然中断,网关就无法再切换到备用渠道了。因为客户端已经接收了部分 HTTP 响应体,此时切换会导致客户端收到格式错乱的拼接数据。
针对这种中途断流的情况,网关只能选择立刻关闭连接,将错误抛给客户端,由客户端发起整体重试。高级的网关系统会在内存中建立一个微型的“首包缓冲区”,在确认连接稳定输出若干个 Token 后,才开始向客户端下发数据。
3. 并发冲突与数据库死锁
在上述代码中,我们通过 SQLite 进行了简单的 Token 额度扣减。在低并发场景下这完全没有问题。但在高并发场景下,多个异步协程同时读写同一个 SQLite 数据库,极易引发数据库锁死或由于竞态条件导致用户额度超支(例如,两个并发请求同时读取到剩余额度为 10,然后各自扣减了 8,导致最终额度变为负数)。
为了避免此问题:
- 所有的额度扣减必须采用原子操作,例如使用
UPDATE users SET token_used = token_used + ?配合事务,而不是在 Python 代码中计算完再写回。 - 如果并发量极大,建议将 SQLite 替换为 Redis,并使用 Redis 的
DECRBY或 Lua 脚本来实现无锁的高并发配额管理。
自建网关 vs. 托管聚合服务
自建网关的自由度极高,但随之而来的是服务器运维、多区域网络延迟优化等琐碎工作。通过下表,我们可以直观地对比自建网关与直接使用像 n1n.ai 这样的统一 API 聚合服务的优劣:
| 评估维度 | 自建 FastAPI 网关 | n1n.ai 聚合服务 |
|---|---|---|
| 运维成本 | 高(需要管理 VPS 服务器、SSL 证书、进程守护及监控) | 零(即用即付,Serverless 架构) |
| 网络延迟优化 | 受限于网关服务器的物理位置 | 依托全球边缘网络(Edge Network),就近接入 |
| 多模型支持 | 需要手动维护每个模型的 API Key 和请求格式 | 统一接入 DeepSeek-V3, Claude 3.5 Sonnet, OpenAI o3 等主流模型 |
| 并发性能 | 需自行解决高并发下的数据库读写瓶颈 | 工业级高并发支撑与智能路由调度 |
| 计费对账 | 需自行编写复杂的数据库统计与对账后台 | 提供可视化面板,统一管理各模型消费明细 |
总结
通过部署一个兼容 OpenAI 协议的本地容灾网关,我们可以极大地提升大模型应用在生产环境中的稳定性。这种“主备双通道”的设计思路,能够帮我们平稳渡过各 LLM 服务商的临时故障期。
当然,如果您希望免去服务器配置、网络调优以及复杂的 Token 兼容性开发,选择像 n1n.ai 这样成熟的聚合服务是更为高效的选择。它能够帮您在不改动任何核心业务逻辑的前提下,瞬间获得全网主流大模型的高可用路由能力。
Get a free API key at n1n.ai