小厂 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
}

四、生产治理收益

自研网关上线后,团队获得了三个关键掌控力:

  1. 供应商故障 0 秒级自动切换:某厂商出现机房故障时,网关在连续 3 个请求失败后自动熔断该节点,后续流量 100% 透明转发至备用厂商,业务端完全无感知;
  2. 流量权重热更新:通过对接 Nacos 或 ETCD,动态修改权重值,瞬间实现多模型的蓝绿发布与 A/B 效果测试;
  3. 极简架构与零依赖:单二进制文件部署,无任何复杂运维依赖,极度契合小团队的敏捷基因。
Logo

欢迎加入DeepSeek 技术社区。在这里,你可以找到志同道合的朋友,共同探索AI技术的奥秘。

更多推荐