小厂 AI 网关自研之路:基于 Go 实现模型动态路由、权重分配与熔断器
·
小厂 AI 网关自研之路:基于 Go 实现模型动态路由、权重分配与熔断器

随着团队内大模型业务的高速扩展,我们逐步接入了多家公有云大模型供应商(阿里云百炼、腾讯混元、百度千帆、OpenAI、DeepSeek)以及团队本地自建的 vLLM 私有推理节点。
然而在多模型共存的生产环境下,各种痛点接踵而至:
- 流量分配僵硬:无法在不发版的情况下,动态调整各供应商的流量比例(例如给新模型切 10% 灰度流量试跑);
- 故障无法自愈:主用厂商偶尔发生机房网络中断或 429 报错,全线上游业务直接超时瘫痪,必须人工半夜爬起来改配置重启服务;
- 缺乏多租户优先级调度:核心客户的实时在线对话与后台异步长报表生成混在同一个队列,导致核心用户卡顿。
针对这些痛点,我们用 Go 语言自研了一套轻量级、高吞吐的 AI 动态路由网关(LLM Dynamic Gateway)。单进程支持每秒 2 万并发路由,内存占用不足 80MB。
今天我们把这套网关的加权动态路由算法、基于滑动窗口的自愈熔断器(Circuit Breaker)与故障秒级转移核心代码完整开源分享。
一、架构设计全景与核心组件
flowchart TD
ClientReq[客户端统一 API 请求] --> Router[动态加权路由器 WeightedRouter]
Router --> CheckBreaker{熔断器状态判定 CircuitBreaker}
CheckBreaker -- 节点熔断打开 (Open) --> Fallback[秒级降级至备用供应商]
CheckBreaker -- 节点正常 (Closed) / 半开 (Half-Open) --> SendPrimary[转发至目标模型渠道]
SendPrimary --> ResponseStatus{响应判定}
ResponseStatus -- 成功 200 --> RecordOk[上报成功,重置错误窗口]
ResponseStatus -- 连续 5xx/429/超时 --> TripBreaker[熔断器触发,切断该渠道流量]
二、生产级核心代码实现
1. 带平滑加权轮询的动态路由器(Smooth Weighted Round-Robin)
直接随机分配权重容易在微小时间窗口内产生局部倾斜。我们采用 Nginx 同款的 平滑加权轮询算法(Smooth Weighted Round-Robin),让流量分布在时间轴上极其均匀。
package gateway
import (
"errors"
"sync"
)
type UpstreamNode struct {
Name string
BaseURL string
APIKey string
Weight int // 配置的基础权重
CurrentWeight int // 运行时动态权重
Breaker *CircuitBreaker
}
type DynamicRouter struct {
mu sync.RWMutex
nodes []*UpstreamNode
}
func NewDynamicRouter(nodes []*UpstreamNode) *DynamicRouter {
return &DynamicRouter{nodes: nodes}
}
// 平滑加权选择健康节点
func (r *DynamicRouter) SelectNode() (*UpstreamNode, error) {
r.mu.Lock()
defer r.mu.Unlock()
var bestNode *UpstreamNode
totalWeight := 0
for _, node := range r.nodes {
// 过滤掉当前已被熔断器拉黑的不可用节点
if !node.Breaker.AllowRequest() {
continue
}
node.CurrentWeight += node.Weight
totalWeight += node.Weight
if bestNode == nil || node.CurrentWeight > bestNode.CurrentWeight {
bestNode = node
}
}
if bestNode == nil {
return nil, errors.New("no available healthy upstream node")
}
// 扣除总权重,实现平滑轮询
bestNode.CurrentWeight -= totalWeight
return bestNode, nil
}
2. 基于滑动窗口的自愈熔断器(Circuit Breaker)
package gateway
import (
"sync"
"time"
)
type State int
const (
StateClosed State = iota // 正常通信
StateOpen // 熔断打开,全量阻断
StateHalfOpen // 半开探测,尝试放行部分请求
)
type CircuitBreaker struct {
mu sync.Mutex
state State
failureCount int
maxFailures int // 连续失败阈值 (如 5 次)
openTimeout time.Duration // 熔断冷却期 (如 30 秒后尝试自愈)
lastStateTime time.Time
}
func NewCircuitBreaker(maxFailures int, openTimeout time.Duration) *CircuitBreaker {
return &CircuitBreaker{
state: StateClosed,
maxFailures: maxFailures,
openTimeout: openTimeout,
}
}
func (cb *CircuitBreaker) AllowRequest() bool {
cb.mu.Lock()
defer cb.mu.Unlock()
now := time.Now()
if cb.state == StateOpen {
// 检查是否已过冷却期,若是则进入半开状态
if now.Sub(cb.lastStateTime) > cb.openTimeout {
cb.state = StateHalfOpen
cb.lastStateTime = now
return true
}
return false // 仍在熔断期,坚决拒绝
}
return true
}
func (cb *CircuitBreaker) RecordSuccess() {
cb.mu.Lock()
defer cb.mu.Unlock()
if cb.state == StateHalfOpen {
// 半开探测成功,彻底自愈恢复正常
cb.state = StateClosed
cb.failureCount = 0
} else if cb.state == StateClosed {
cb.failureCount = 0
}
}
func (cb *CircuitBreaker) RecordFailure() {
cb.mu.Lock()
defer cb.mu.Unlock()
cb.failureCount++
now := time.Now()
if cb.state == StateHalfOpen || cb.failureCount >= cb.maxFailures {
// 触发熔断
cb.state = StateOpen
cb.lastStateTime = now
}
}
三、故障自动秒级转移(Failover)代理执行器
package gateway
import (
"bytes"
"context"
"io"
"net/http"
"time"
)
type ProxyClient struct {
router *DynamicRouter
httpClient *http.Client
}
func (p *ProxyClient) ExecuteWithFailover(ctx context.Context, payload []byte) ([]byte, error) {
maxAttempts := 3
var lastErr error
for attempt := 0; attempt < maxAttempts; attempt++ {
node, err := p.router.SelectNode()
if err != nil {
return nil, err
}
req, _ := http.NewRequestWithContext(ctx, "POST", node.BaseURL+"/v1/chat/completions", bytes.NewReader(payload))
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", "Bearer "+node.APIKey)
resp, err := p.httpClient.Do(req)
if err != nil || resp.StatusCode >= 500 || resp.StatusCode == 429 {
// 记录失败并触发熔断统计
node.Breaker.RecordFailure()
lastErr = err
if resp != nil {
resp.Body.Close()
}
// 立即轮询下一个可用节点重试
continue
}
defer resp.Body.Close()
body, _ := io.ReadAll(resp.Body)
node.Breaker.RecordSuccess()
return body, nil
}
return nil, lastErr
}
四、生产治理收益
自研网关上线后,团队获得了三个关键掌控力:
- 供应商故障 0 秒级自动切换:某厂商出现机房故障时,网关在连续 3 个请求失败后自动熔断该节点,后续流量 100% 透明转发至备用厂商,业务端完全无感知;
- 流量权重热更新:通过对接 Nacos 或 ETCD,动态修改权重值,瞬间实现多模型的蓝绿发布与 A/B 效果测试;
- 极简架构与零依赖:单二进制文件部署,无任何复杂运维依赖,极度契合小团队的敏捷基因。
更多推荐

所有评论(0)