Go手写LLM推理网关:SSE流式转发与背压控制实践
2026/9/21 19:45:39 网站建设 项目流程

先说个结论:任何一个真实可用的 LLM 应用,只要用户量上来,迟早得在应用和后端模型之间加一层网关。我这次用 Go 手写了一个专注于流式转发场景的 LLM 推理网关,核心就三件事——SSE 流式转发、级联取消、背压控制。这篇文章会把整个设计和实现思路原原本本拆开讲,包括踩过的坑和排查过程。如果你正在做 LLM Agent、聊天机器人,或者打算把多个模型供应商接入同一个内部平台,这篇很适合你。

1. 先搞清楚推理网关到底在解决什么问题

1.1 没有网关的时候为什么手忙脚乱

大多数团队最开始接 LLM 都是直连。业务代码里直接写 OpenAI SDK,配置一个 api key,然后调用 chat completions 接口。短期看没问题,但用一段时间就会暴露出来几类矛盾。

第一类是多供应商切换问题。今天想换国产模型,明天想接开源私有化部署模型,每个厂商的接口协议、鉴权方式、限流策略都不同,业务代码里开始大批量出现 if else。时间一久,接口调用逻辑和业务逻辑耦合在一起,任何一个供应商接口升级,业务层都得跟着改。第二类是流式连接的管理问题。LLM 接口大多支持 SSE 流式返回,一个对话可能要持续几十秒到几分钟,但业务服务通常没有为长连接做专门优化,超时时间、并发数、断线重试都只能草草设置。第三类是成本和安全问题。如果客户端直接拿到上游 API key,那泄密、盗刷、恶意刷量都很难管控。如果真的直接把一个带 key 的地址发给前端,我见过有人一天被刷掉几千块。

推理网关本质上就是一个中间层,负责把业务服务与模型供应商之间的复杂交互收敛掉。客户端只感知一个统一入口,其他事情全由网关解决。

1.2 这只网关应该承担哪些职责

设计的时候我把网关职责拆成了几层。

  • 协议接入层:对外暴露统一的 HTTP 接口,兼容 OpenAI 风格的请求格式,让业务方接入成本为零。
  • 路由分发层:根据模型名称、租户、优先级把请求路由到不同的上游,支持多供应商配置和故障转移。
  • 流式转发层:负责 SSE 的逐 token 读取、转发、缓冲和格式化输出。
  • 生命周期管理层:把客户端断连、上游断连、超时、任务取消这些事件串成完整的级联流程,保证资源不泄漏。
  • 流量控制层:实现并发限制、排队和熔断,也就是背压控制,防止突发流量打垮上游。

这里最容易被忽略的是生命周期管理层。很多人实现完前面 3 层就上线,结果线上出现大量 goroutine 泄漏和连接堆积。问题通常出在:客户端关掉了页面,但上游还在继续生成 token,网关协程还在傻傻地转发。这正是级联取消要解决的核心场景。

2. SSE 流式转发的关键设计与实现

2.1 SSE 协议在 Go 里到底怎么落地

SSE 全称 Server-Sent Events,本质上就是 HTTP 长连接配合text/event-stream媒体类型。服务端可以在一个连接上多次写入数据,客户端通过 EventSource API 自动接收。协议格式特别简单,每条消息用空行分隔,字段以key: value形式存在,最常用的是data

测试的时候直接curl -N就能看到原始格式。

curl -N https://your-gateway.example.com/v1/chat/completions \ -H "Content-Type: application/json" \ -d '{"model":"gpt-4o-mini","stream":true,"messages":[{"role":"user","content":"hi"}]}'

输出长这样:

data: {"id":"chatcmpl-xxx","choices":[{"delta":{"role":"assistant"},"index":0}]} data: {"id":"chatcmpl-xxx","choices":[{"delta":{"content":"你"},"index":0}]} data: [DONE]

在 Go 标准库中,SSE 其实不需要额外依赖。核心就是http.Flusher接口。普通 HTTP 响应只有在 handler 返回后才会真正把数据刷到客户端,但http.ResponseWriter如果实现了 Flusher,就可以在每次写完数据后手动调用 Flush 强制把缓冲数据推给客户端。这个机制在 Go 的 net/http 里原生支持,只要响应头没有禁用 chunked 编码,底层会自动用分块传输。

2.2 一次性打通客户端到上游的 SSE 管道

网关要做的事很简单:接收客户端请求,再用同样的请求体去请求上游 LLM 接口,然后把上游返回的 SSE 字节流原样转发给客户端。听起来简单,但里面有一个非常关键的坑:ResponseWriter 的缓冲。

我先给出一个最完整的流式转发实现,再逐行解释。

func handleStreamProxy(w http.ResponseWriter, r *http.Request) { // 1. 先为响应设置 SSE 相关请求头 w.Header().Set("Content-Type", "text/event-stream") w.Header().Set("Cache-Control", "no-cache") w.Header().Set("Connection", "keep-alive") w.Header().Set("X-Accel-Buffering", "no") // 关键:关闭 Nginx 缓冲 flusher, ok := w.(http.Flusher) if !ok { http.Error(w, "streaming unsupported", http.StatusInternalServerError) return } // 2. 构造上游请求,重要的是把客户端 context 传递进去 upstreamURL := "https://api.example.com/v1/chat/completions" payload, _ := io.ReadAll(r.Body) req, err := http.NewRequestWithContext(r.Context(), http.MethodPost, upstreamURL, bytes.NewReader(payload)) if err != nil { http.Error(w, err.Error(), http.StatusBadGateway) return } req.Header.Set("Content-Type", "application/json") req.Header.Set("Authorization", "Bearer "+getUpstreamKey(r)) // 3. 发起上游请求 client := &http.Client{Timeout: 0} // 流式请求不能设整体超时 resp, err := client.Do(req) if err != nil { http.Error(w, err.Error(), http.StatusBadGateway) return } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { body, _ := io.ReadAll(resp.Body) http.Error(w, string(body), resp.StatusCode) return } // 4. 关键:逐行读取上游 SSE 响应,再逐行写回客户端 reader := bufio.NewReader(resp.Body) for { line, err := reader.ReadString('\n') if len(line) > 0 { if _, writeErr := w.Write([]byte(line)); writeErr != nil { // 客户端已断开,停止读取上游 return } flusher.Flush() } if err != nil { if err == io.EOF { return } // 上游读取出错,按 SSE 格式发送错误后退出 _ = writeSSEError(w, flusher, err) return } } }

这段代码的核心逻辑是逐行读取上游的 SSE 响应,然后原样写回给客户端,每写一行就Flush一次。这里有几个很容易被忽视的设计决策。

第一个是X-Accel-Buffering: no。如果网关前面还套了 Nginx,默认情况下 Nginx 会缓冲响应,等到整个请求结束才会把数据推给浏览器,这就等于把 SSE 的实时性完全抹掉了。这个响应头是 Nginx 官方支持的关闭缓冲的开关,没有它,流式接口到了 Nginx 这边就变成了"憋大招"。

第二个是http.Client.Timeout必须设置为 0,但这不是让你完全放弃超时控制。流式请求持续时间可能非常长,一个整体超时设置会把长对话直接砍断。正确的超时控制方式是通过context.WithTimeout或者自定义Transport层配置 Dial/ResponseHeader 超时,而不是Client的整体 Timeout。

第三个是bufio.Reader.ReadString按行读取。SSE 协议每条消息以换行符分隔,所以按行读天然适合这个场景。有些实现用io.Copy直接管道复制,虽然代码短,但会丢掉按事件处理数据的机会,也没办法在出错时插入自定义 SSE 事件。我建议网关这种场景用逐行读的方式。

2.3 不能只做拷贝:请求转换与重试的处理

上面是最基础的透传,但实际生产中通常在请求级别就要做一些改造。我这边总结了三类必须处理的情况。

第一类是上游响应格式不统一的问题。不同供应商的流式返回不一定完全兼容 OpenAI 格式,有些会额外加字段,有些字段命名不一样。如果客户端依赖严格格式,网关最好在转发时做一层标准化,比如把choices[0].delta.content统一映射。这个做起来不复杂,将响应按行解析为 JSON,检查是[DONE]还是数据事件,再重新生成 JSON 写回。

第二类是鉴权信息的注入。客户端给网关的 key 不应该直接透传给上游。比较稳妥的做法是网关内部维护一张 key 映射表,入站请求使用网关分发的 key 校验身份,出站请求使用网关配置的上游真实 key。这样上游 key 永远只存在服务端配置中。

第三类是重试策略。流式请求重试比普通请求危险得多,因为模型已经开始生成了,重试会导致重复计费。我建议只在以下情况重试:上游返回 429、5xx、或者连接建立阶段失败;一旦已经收到第一个数据字节,就绝不能重试。毕竟此时计费已经发生或模型已经消费请求,无脑重试只会放大成本。

这里还可以加入一个简单的基于 channel 的背压缓冲机制,防止上游吐数据速度远快于客户端接收速度时内存暴涨。关于这一点,后面背压部分展开讲。

3. 级联取消:一个 context 贯穿全链路

3.1 为什么会泄取消

Go 的 context 是级联取消的地基。客户端断开、上游断开、网关主动取消,这三件事会沿着 context 父子关系传播。但不少人在写网关时只把它当成"超时控制工具",忽略了它更重要的能力:取消信号传播。

想象一个典型场景:用户在浏览器打开页面,点击"开始生成",网关向上游发起请求,上游开始一个 token 一个 token 地返回。这时候用户直接关掉了浏览器。客户端连接断开后,如果网关没有感知到这个事件,它会继续从上游读取数据,往上写的时候才会发现连接已经断了,然后才返回。看起来好像也没啥问题,但如果在每层都加一个超时等待呢?上游会把这一整段响应全部生成完,token 数量可能几千个,计费按 token 算,等于用户的取消操作根本没能止损。更严重的是,如果网关并发处理很多这种请求,上游和网关之间还会堆积大量无效流量,拖慢正常请求的响应速度。

所以级联取消的核心是:取消信号要尽可能快地扩散到所有相关环节,谁也别等谁。

3.2 context 从哪来、到哪去

在 Go 的 net/http 服务中,每个请求自带一个可取消的 context,客户端连接断开时它会自动被取消。这个 context 就是整条取消链路的最上游。

func handleStreamProxy(w http.ResponseWriter, r *http.Request) { ctx := r.Context() go func() { <-ctx.Done() log.Printf("client disconnected: %v", ctx.Err()) }() // 后续所有上游请求,都基于 ctx }

在网关内部向上游发请求时,必须使用http.NewRequestWithContext(r.Context(), ...),而不是裸的http.NewRequest。这样才能实现:客户端一断,上游请求的 context 跟着取消,连接被关闭,模型侧收到断开信号。

不过有同学会问:上游如果已经收到了请求,并且已经开始流式生成了,网关这边 context 取消后,上游真的会停止吗?答案是取决于上游的实现。OpenAI 等主流接口对连接断开是敏感的,gateway 断开连接,上游服务会发现写失败并停止生成。虽然不是百分百保证,但至少绝大多数场景下能见效。比不传 context、傻等上游自己超时要强得多。

3.3 在响应循环里做主动取消检查

除了依赖 Transport 层自动取消,我建议在读取 SSE 的循环里做一次显式的取消检查。为什么需要这层双保险?因为有些网络库或者自定义协议实现可能在 context 取消后并没有立刻终止当前的阻塞读,这会导致 goroutine 无法及时退出。

func pipeSSE(ctx context.Context, reader *bufio.Reader, w http.ResponseWriter, flusher http.Flusher) error { dataCh := make(chan string, 1) errCh := make(chan error, 1) go func() { for { line, err := reader.ReadString('\n') if len(line) > 0 { select { case dataCh <- line: case <-ctx.Done(): return } } if err != nil { errCh <- err return } } }() for { select { case <-ctx.Done(): return ctx.Err() case line := <-dataCh: if _, err := w.Write([]byte(line)); err != nil { return err } flusher.Flush() case err := <-errCh: if err == io.EOF { return nil } return err } } }

这种做法把读上游数据和写客户端拆成了两个独立流程,中间通过 channel 通信。主循环里用 select 同时监听 context 取消和可写数据,这样 context 一旦取消,主循环能立刻退出,无论此时是否正在等待数据。

这里有一个可选的优化:channel 缓冲设置成 1 就够了。因为 SSE 数据本身就是逐行流动的,缓冲区太小会导致频繁切换,太大则会导致上游已经读到很多行、客户端却不消费,内存占用飙升。缓冲为 1 可以兼顾吞吐和内存安全。之所以不设置更大的缓冲,是因为网关不负责聚合和缓存数据,数据越快推到客户端越好。

3.4 超时控制在级联中的位置

级联取消还需要考虑超时控制的层级。我的经验是:给"每一跳"都设置不同的超时,而不是一刀切。

  • 客户端到网关:这个连接的超时由前端处理,网关不主控,主要靠 detect client disconnect。
  • 网关到上游的连接建立阶段:设置较短的超时,比如 10 秒,防止上游不可达时网关长时间挂起。
  • 上游的首字节超时:设置一个中等超时,比如 60 秒。因为模型有时排队很长,首 token 可能要等几秒到几十秒。
  • 整体流式时间:不设上限,但可以在业务层做最大 token 数控制。

在代码里可以这样结合 context。

connectCtx, cancelConnect := context.WithTimeout(ctx, 10*time.Second) defer cancelConnect() req, _ := http.NewRequestWithContext(connectCtx, ...) firstByteCtx, cancelFirstByte := context.WithTimeout(ctx, 60*time.Second) defer cancelFirstByte() resp, err := client.Do(req) // 收到响应头后,立刻把首字节超时取消掉,回到纯粹由客户端断开控制的 ctx

这里最容易被忽略的细节是:一旦进入流式读取阶段,就应该立刻切回原始的r.Context()来控制生命周期,而不是继续使用带超时的 context。否则,一个"最长允许 60 秒生成"的约束会把长对话拦腰截断。真正应当按时间的限制,只放在连接建立和首字节等待这两个阶段。

4. 背压:别让网关替上游挡枪

4.1 背压的本质是保护最脆弱的环节

背压这个词听起来学术,其实特别朴素:当下游消费不过来的时候,上游不该继续全速猛灌,而是要减速或者暂停。想象一个自助餐厅,厨师疯狂炒菜,但取餐的客人只有一两个,最后菜全堆在台面上,浪费又变凉。在流式系统里,客户端网络比较慢,模型侧吐 token 很快,如果网关不控制,数据会积压在内存里,吃光内存之后整个服务被 OS 杀掉。

网关侧背压有两个维度要处理。

第一个维度是连接数背压。当并发进来的请求数超过阈值,网关应该丢弃或排队,而不是无限地创建 goroutine 和上游连接。每个 goroutine 都有栈内存占用,1000 个 goroutine 不算什么,5 万个同时挂着,内存就会非常吓人。

第二个维度是单连接内字节流背压。客户端读取速度慢,上游发送速度快,这时候读上游的循环不能无脑继续读,应该控制读取速率,或者限制缓冲队列深度。

4.2 信号量限流与有界队列

网关内实现背压最简单的方式是信号量。用带缓冲的 channel 实现一个计数信号量。

type ConnLimiter struct { tokens chan struct{} } func NewConnLimiter(limit int) *ConnLimiter { return &ConnLimiter{tokens: make(chan struct{}, limit)} } func (l *ConnLimiter) Acquire(ctx context.Context) error { select { case l.tokens <- struct{}{}: return nil case <-ctx.Done(): return ctx.Err() } } func (l *ConnLimiter) Release() { <-l.tokens }

在进入 handler 时获取信号量,退出时释放。如果当前并发数已达到上限,Acquire 会阻塞等待,直到有连接释放。

但纯粹的信号量只能限制并发数,不能优雅地排队。生产环境里我们更想要的行为是:超过并发上限后,新的请求先排队一会儿,而不是立即拒绝。这里可以用有界队列配合 worker 池来实现。

伪代码思路是这样的。

type RequestQueue struct { queue chan *http.Request workers int }

实际上对于网关这种 I/O 密集场景,我推荐一个更简洁的方案:直接使用带缓冲的请求 channel 配合固定数量的 worker goroutine。worker 数控制并发上游请求数量,队列深度控制等待数量。

  • 当队列未满,新请求进入队列等待 worker 处理。
  • 当队列已满,直接返回 429 或者 503,并在响应头里带上 Retry-After,引导客户端稍后重试。
  • worker 内部的流式转发照常进行,不因为排队就影响已有连接。

这样做的好处是,背压信号(排队等待、限流拒绝)会一直传导到最上游——客户端。客户端看到 429/503,就会自动放慢发送请求的速率,从而保护整个链路。

4.3 配合 Redis 做分布式限流

单机信号量只能说解决了进程内的背压。如果网关部署了多副本,每个副本的限流阈值是独立的,总并发上限会被放大 N 倍。对于 LLM 请求成本较高的场景,还是需要一轮分布式限流。

我用的方案是 Redis + 滑动窗口,或者更简单一点的令牌桶。核心是保证限流判断的原子性,避免并发请求同时通过检查。在 Go 里对 Redis 的操作建议直接用管道或者 Lua 脚本保证原子性。下面是一个基于 Redis + Lua 实现固定窗口限流的例子。

-- KEYS[1]: 限流 key -- ARGV[1]: 窗口大小(秒) -- ARGV[2]: 窗口内最大请求数 local current = tonumber(redis.call('GET', KEYS[1]) or '0') if current >= tonumber(ARGV[2]) then return 0 end redis.call('INCR', KEYS[1]) redis.call('EXPIRE', KEYS[1], tonumber(ARGV[1])) return 1

go-redisEval方法执行这个脚本。脚本会原子地完成"检查 + 计数"两步,避免并发穿透。

func allowRequest(ctx context.Context, rdb *redis.Client, key string, limit int, window int) bool { result, err := rdb.Eval(ctx, luaScript, []string{key}, window, limit).Int() if err != nil { // Redis 故障时,建议降级为放行,但要记录告警 return true } return result == 1 }

我的建议是:Redis 限流用于粗粒度控制,单机信号量用于细粒度保护,两者组合使用。即使 Redis 挂了,也不会完全失去背压保护。这是因为降级为放行策略后,单机限流还在兜底——毕竟协调中心不可用是故障场景,此时优先保证可用性是对的。

4.4 上游慢、客户端更慢:单连接内的速率控制

很多人在背压设计中遗漏了单连接内的场景,我也是在真机上摸爬滚打才发现问题的。

如果客户端是通过很差的网络访问网关,比如办公室 Wi-Fi,它的 TCP 接收窗口可能很小。网关向上游读数据很快,向下游写数据却很慢,这时如果不做处理,网关的写缓冲会持续增长,内存一路膨胀,最后 OOM。

处理思路是限制每个连接内的读缓冲深度。简单一点的做法是:在 SSE 转发循环里,不要再往 w.Write 里写大块数据,而是把每行数据控制在较小的长度,写完立刻 Flush,让 TCP 层的背压自然作用于上游读取。这种方式下如果客户端实在消费不动,会导致我们的写操作阻塞,进而让 main loop 阻塞在读上游的循环里,整个连接就被"卡住"了。这本质上就是 TCP 背压在起作用:客户端处理不过来,窗口变小,我们写不进去,自然也就不再去上游读新数据。

不过这种情况意味着一个慢客户端会占用一个 goroutine 很久。所以还需要配合"最大连接时长"和"空闲超时"来兜底,避免无限期占用。

5. 实践中踩过的坑与排查实录

5.1 idle timeout waiting for SSE:不是代码问题,是代理问题

我之前拿网关对接某个基于 Codex 协议的服务时,频繁出现类似的错误信息:

stream disconnected before completion: idle timeout waiting for sse

翻译过来就是:上游在应该持续推送 SSE 数据时,超过某个时间没推任何数据,然后连接被判定空闲超时。问题排查了很久,最终发现根因往往不在上游,而是网关所在环境或者网关上层的代理把"空闲"误判为"超时"。

常见触发点有两个。第一是模型侧在生成过程中,如果一次思考时间很长,就会出现超过几十秒没有任何数据推送的情况。SSE 连接活着,但没有数据流动。第二是 Nginx 的proxy_read_timeout默认值是 60 秒,如果 60 秒内没有新数据,Nginx 会主动断开后端连接,客户端就会收到连接中断。

解决方案是在 Nginx 配置里把超时时间调大:

location /v1/chat/completions { proxy_http_version 1.1; proxy_set_header Connection ""; proxy_buffering off; proxy_cache off; proxy_read_timeout 3600s; proxy_send_timeout 3600s; }

如果是网关直连上游,那就要检查http.Transport.ResponseHeaderTimeouthttp.Transport.ExpectContinueTimeout是否设了过小的值。另外,有些 LLM 服务端为了避免中间链路把长思考的空闲连接断掉,会主动发: keepalive注释行。这是 SSE 协议里的心跳机制,注释行会被客户端忽略,但能维持连接活性。如果你的上游不支持心跳,网关可以自己加一个心跳定时器:每 15 秒往上游写一个注释行。

if time.Since(lastWrite) > heartbeatInterval { _, _ = w.Write([]byte(": keepalive\n\n")) flusher.Flush() }

加心跳不是万能的,如果代理只是按照字节数来判定超时,那心跳能救;如果代理严格执行"必须业务数据才算活跃",那还是得调代理配置。

5.2 连接挂死与 read buffer 残留

另一个坑出现在bufio.Reader读取上游响应时。如果上游发送的 SSE 消息结尾只有\n而没有\r\n,某些严格实现的解析器会出问题。我遇到过上游把\r\n\n混着发的场景,最稳妥的处理是读取时统一规范化。

方法很简单:读行时用ReadBytes('\n'),然后去掉末尾的\r\n,再按自己的格式重新组织输出。不要直接依赖上游的换行风格。因为不同供应商对 SSE 协议的实现细节差异很大,有的喜欢带\r\n,有的只用\n,还有的在事件之间加多个空行。

line, err := reader.ReadBytes('\n') line = bytes.TrimRight(line, "\r\n") if len(line) == 0 { // 空行是分隔符,直接透传空行 w.Write([]byte("\n")) continue }

这个处理能避免很多肉眼很难发现的兼容性问题。

另外还有一个比较隐蔽的问题:如果网关使用了连接池,上游响应内容很大,连接池复用的时候可能会把上一次未读完的响应体残留带到下一个请求里。解决方式很简单——开启 HTTP 响应的Body必须被完全读取并关闭,否则 Transport 不会把这个连接返回到池中。所以即使读取出错了,也要调用resp.Body.Close(),并且记住主动 drain 掉剩余未读内容。

5.3 可观测性:slog 日志与流式指标

网关的可观测性可以从日志和指标两个层面考虑。我自己使用的是 Go 1.21 开始内置的log/slog,结构化日志非常好用,不需要额外引入第三方日志包。

logger := slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{ Level: slog.LevelDebug, })) slog.SetDefault(logger) slog.Info("stream proxy start", "req_id", reqID, "model", payload.Model, "upstream", upstreamName, "client_ip", clientIP, )

比较重要的埋点时机包括:请求进入、获取上游响应头、第一个字节写回客户端、连接中断、请求完成。通过这些日志可以还原整个请求的生命周期,排查问题时非常有用。

流式网关真正需要关注的指标也和普通 API 不太一样,我总结了几个关键项。

  • 并发流式连接数:判断当前系统压力,和信号量限流的 margin 对比。
  • 首 token 延迟:用户能感受到的第一个字有多快。
  • 平均/最大 token 间延迟:两个 token 之间的平均间隔,评估流畅度。
  • 客户端中断率:客户端断连次数在总请求中的占比。
  • 上游重试率:判断上游稳定性。

这些指标可以通过 Prometheus 客户端库暴露,也可以在内部先用日志统计,不需要一下子全上。

5.4 不要忽略进程内限流的并发安全

如果你在网关里用计数器实现限流,有个很经典的并发 bug:先判断再累加不是原子操作,高并发下会有多个请求同时通过检查,导致实际并发数超过上限。这也是为什么我选择用 buffered channel 作为信号量,因为 channel 的入队操作本身是并发安全的。

// 错误示例:并发不安全 if counter < max { counter++ // 多个协程同时执行时都通过了 if go handle() }

如果确实要用计数器,就用sync/atomicCompareAndSwap或者atomic.AddInt64做原子操作。但 channel 的方式更符合 Go 的惯例,代码也更容易理解。

5.5 小细节:关闭响应体不等于关闭上游连接

还有一个容易被忽略的点:http 客户端resp.Body.Close()的行为。当你Close一个尚未读完的响应体,Transport 不会复用这个连接,而是直接关闭它。这意味着如果我们在 SSE 转发过程中提前返回(比如客户端断开),defer resp.Body.Close()会顺手把上游连接关闭,导致上游发送端收到连接重置信号,从而停止生成。这其实正是我们想要的效果。

但有些场景下我们可能希望能提前"优雅地"关闭,比如让上游少做一点无效计算。这时可以尝试在 context 取消后向上游发送一个自定义的 SSE 事件请求,告知上游停止生成。不过这依赖上游协议是否支持这种控制指令,我没法保证所有平台都适用,所以目前更多是依赖连接断开来实现。

6. 完整网关骨架和后续扩展建议

我实现的核心网关逻辑其实只有 200 行左右,没有用什么重型框架,标准库加 go-redis 就足够了。下面是一个精简但完整的骨架。

package main import ( "bufio" "bytes" "context" "io" "log/slog" "net/http" "os" "time" ) type Gateway struct { limiter *ConnLimiter upstream string } func (g *Gateway) ServeHTTP(w http.ResponseWriter, r *http.Request) { if err := g.limiter.Acquire(r.Context()); err != nil { http.Error(w, "too many concurrent requests", http.StatusTooManyRequests) return } defer g.limiter.Release() ctx := r.Context() req, err := http.NewRequestWithContext(ctx, http.MethodPost, g.upstream, r.Body) if err != nil { http.Error(w, err.Error(), http.StatusBadGateway) return } req.Header = r.Header.Clone() req.Header.Set("Authorization", "Bearer "+os.Getenv("UPSTREAM_KEY")) resp, err := http.DefaultClient.Do(req) if err != nil { http.Error(w, err.Error(), http.StatusBadGateway) return } defer resp.Body.Close() w.Header().Set("Content-Type", "text/event-stream") w.Header().Set("Cache-Control", "no-cache") w.Header().Set("Connection", "keep-alive") w.Header().Set("X-Accel-Buffering", "no") flusher := w.(http.Flusher) reader := bufio.NewReader(resp.Body) for { line, err := reader.ReadBytes('\n') if len(line) > 0 { if _, writeErr := w.Write(line); writeErr != nil { slog.Info("client disconnected", "err", writeErr) return } flusher.Flush() } if err != nil { if err == io.EOF { return } slog.Error("read upstream failed", "err", err) return } } }

整个骨架到这里就能跑起来做基本验证了。如果要在生产环境使用,还有几个扩展方向值得探索。

第一个是支持多租户鉴权和用量统计,每个业务方分配独立 key,网关层按租户统计 token 消耗和请求量,方便做成本分摊。第二个是支持多种上游协议转换,把非 OpenAI 格式的模型接口转换成 OpenAI 风格,这样业务方不需要为不同模型写多套 SDK。第三个是支持模型 fallback,当主模型返回 429 或者 5xx 时,自动切换到备用模型,但要注意流式场景下切换的语义。

我个人在实际操作中的体会是:网关最核心的价值从来不是"转发"这个动作,而是把生命周期、流量控制和可观察性这些横切关注点统一收敛到一层。把 context 的取消信号一路传到底,把背压一路传到顶,剩下的事情就都顺了。希望这份记录能帮到正在踩坑的各位,少烧几个深夜。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询