inference-gateway 的 SSE 流式代理:Timeout:0、errgroup 取消与 io.Copy 背压
在 inference-gateway(个人学习项目)里,我用 Go + GoFrame 写了一个前置在 vLLM(OpenAI 兼容接口)之前的推理网关。这里记录并发长连接链路中最核心的一段:SSE 流式反向代理 Streamer 的设计、踩坑和重构。项目未提供仓库地址。
1. 传统 API 网关和 LLM 推理网关的差异
普通微服务网关的请求生命周期通常是短连接、快速响应,几十到数百毫秒。大模型推理(比如 /v1/chat/completions 生成 1000 个 token)会把请求生命周期拉到另一个量级:
| 维度 | 传统微服务网关 | LLM 推理网关 |
|---|---|---|
| 连接时长 | 短平快(50ms ~ 500ms) | 极长连接(5s ~ 60s+) |
| 响应模式 | 一次性完整 JSON 返回 | Server-Sent Events (SSE) 逐字分块回写 |
| 中断代价 | 客户端断开影响小 | 极其昂贵:若不掐断后端,vLLM 会继续算完整段,白白浪费 GPU 算力与 KV Cache |
| 延迟关注点 | RT(总响应时间) | TTFT(Time To First Token 首字延迟)+ TPOT(每个 Token 吐出间隔) |
2. HTTP 客户端与连接池:Timeout:0 和 MaxIdleConnsPerHost
func NewStreamer(r *router.Router) *Streamer {
return &Streamer{
router: r,
client: &http.Client{
Timeout: 0, // 流式不设整体超时;按需在 Transport 层控制
Transport: &http.Transport{
MaxIdleConnsPerHost: 100,
IdleConnTimeout: 90 * time.Second,
},
},
}
}
这里有两个点。
① 为什么 Timeout: 0?
http.Client.Timeout 的语义是:从发起请求到整个响应体(Response Body)完全读完的总耗时。如果设置常见的 Timeout: 10s,用户生成一篇长文章耗时 30 秒,连接会在第 10 秒被 Go 标准库强行截断。所以这里必须设为 0,不限时,真正的请求生命周期交由 context.Context 动态管控。
② 为什么必须调大 MaxIdleConnsPerHost?
Go 标准库默认的 DefaultMaxIdleConnsPerHost 只有 2。推理网关下游有成百上千个并发用户,但上游的 vLLM 实例通常只有固定几个节点,比如 http://localhost:8001。如果保持默认值 2,大量并发请求结束后,多余的 TCP 连接会被直接关闭。下一个请求又必须经历 TCP 三次握手,导致大量连接处于 TIME_WAIT,不仅拖慢首字延迟(TTFT),还会耗尽操作系统的文件描述符。调大至 100 并配合 IdleConnTimeout: 90s,实现与推理后端的高复用长连接池。
3. 客户端断连与 context 级联取消
用户在前端网页点击“停止生成(Stop Generating)”或直接关闭标签页时,网关必须立刻感知并掐断向 vLLM 的请求。
eg, ctx := errgroup.WithContext(r.Context())
outReq, err := http.NewRequestWithContext(
ctx, r.Method, target+r.URL.Path+"?"+r.URL.RawQuery, r.Body,
)
...
resp, err := s.client.Do(outReq)
if err != nil {
// 若客户端已取消,直接静默退出
if ctx.Err() != nil {
return
}
r.Response.WriteStatusExit(502, errJSON("upstream error: "+err.Error()))
return
}
底层网络信号传导链条:
- 客户端(用户点击“停止” / 关网页 / 网络断开)发送 TCP FIN 或 RST 报文
- 操作系统内核协议栈(epoll / kqueue 捕获 Socket EOF)
- Go HTTP Server 运行时监听到连接断开,触发请求的 cancelFunc
r.Context()被取消(<-r.Context().Done()关闭,Err() == context.Canceled),由 errgroup 派生的子 ctx 同步取消s.client.Do(outReq)内部的 RoundTrip 监听到ctx.Done(),立即掐断发送给 vLLM 的 TCP 连接- 网关判断:
err != nil && ctx.Err() != nil,确定是客户端主动离开,不报 502,安全退出
区分上游挂了 vs 客户端断开:
若 ctx.Err() == nil,说明客户端连接完好,确实是上游 vLLM 挂了或网络不通,回写 502 Bad Gateway;
若 ctx.Err() != nil,说明是客户端断开,已无需且无法向客户端回写错误,直接 return 释放资源即可。
4. 流式转发重构:4KB buffer 循环换成 io.Copy + flushWriter
重构前(原生骨架):
buf := make([]byte, 4*1024)
for {
n, err := resp.Body.Read(buf)
if n > 0 {
r.Response.Write(buf[:n])
r.Response.Flush()
}
if err == io.EOF {
return
}
if err != nil {
return
}
}
重构后:
// flushWriter 包装底层 Response,满足 io.Writer 接口
type flushWriter struct {
res *ghttp.Response
}
func (w *flushWriter) Write(p []byte) (int, error) {
w.res.Write(p)
w.res.Flush() // 保证每个 chunk 立即推给客户端,极低 TTFT
return len(p), nil
}
// 转发逻辑中:
writer := &flushWriter{res: r.Response}
eg.Go(func() error {
_, copyErr := io.Copy(writer, resp.Body)
if copyErr != nil && copyErr != io.EOF && copyErr != context.Canceled {
return copyErr
}
return nil
})
_ = eg.Wait()
这里的变化主要在四点:
- 内存与 GC 压力:重构前每个请求都在堆上
make([]byte, 4096),高并发下引发内存抖动;重构后io.Copy内部借用sync.Pool复用临时缓冲区,减轻 GC 负担。 - 面向接口的装饰器模式:将“写入即 Flush”的特性抽象成
flushWriter,契合标准库io.Writer。 - 结构化并发:使用 errgroup 管理协程生命周期,方便后续扩展旁路任务,如流式 Token 计数、延时指标打点等。
- 精细化的错误甄别:将
io.EOF(正常输出完毕)与context.Canceled(用户正常停止生成)明确排除在系统异常之外,避免生产环境日志误报。
5. io.Copy 内部循环:读一块写一块
很多人常有一个误解:io.Copy 是不是把上游所有内容全部读入内存后,才一次性写给客户端的?答案是否定的。io.Copy 本质上是一个“水泵式”的事件循环:
for {
nr, er := src.Read(buf) // 1. 尝试从上游读取
if nr > 0 {
nw, ew := dst.Write(buf[0:nr]) // 2. 只要有数据,立即调用写入
if ew != nil {
break
}
}
if er != nil { // 读到 EOF 或出错才跳出
break
}
}
为什么在 LLM 场景下能做到平滑流式推送?
src.Read随产随还:vLLM 每生成一个 token chunk(数十字节)发上网络,网关的resp.Body.Read就会立即解除阻塞返回;dst.Write + Flush()强行冲刷:普通 Web 框架会在内存中缓存数 KB 数据才发。但flushWriter.Write内部显式调用了Flush(),立刻绕过缓冲区,打包成 HTTP Chunked 数据包推给客户端;- 0% CPU 占用等待:在 vLLM 正在 Decode 下一个 token 的几十毫秒间隙里,网络没有新数据,
src.Read会在操作系统 Socket 上挂起,Goroutine 让出 CPU,完全不消耗 CPU 资源。
6. 背压:下游慢客户端如何不压垮网关内存
公网环境下用户网络状况千差万别,比如弱网、移动端卡顿。如果客户端接收速率极慢(1KB/s),而 vLLM 生成速率极快(100KB/s);如果没有背压机制,网关要么在内存中拼命开 buffer 堆积数据导致 OOM,要么丢包。
代码中的自然背压传递链:
客户端网络慢 / 接收窗口满
│ (TCP 接收窗口耗尽,发送 TCP Zero Window 探测)
网关 TCP 发送缓冲区填满
│
w.res.Write(p) 发生阻塞 (阻断在内核态系统调用)
│
io.Copy 内部暂停,无法进入下一轮循环
│
暂停调用 resp.Body.Read(),上游网关接收缓冲区填满
│ (反向压迫 vLLM 的 TCP 发送窗口)
vLLM 的 Socket 发送变慢 / 推理流程受到反压调节
没有任何显式的队列和复杂锁逻辑,仅靠同步接口组合与 TCP 流量控制,自然实现了全链路背压,保证网关在高并发慢客户端场景下的内存安全。
7. 单测:SSE 流式推送与断连取消
在 internal/proxy/proxy_test.go 中,通过 httptest.Server 与 GoFrame 随机端口测试服务,搭了两个关键测试用例。
① 验证 SSE 流式完整性与实时性:模拟 upstream 以固定时间间隔(如 20ms)分块吐出数据,验证网关能完整流式转发:
func TestStreamer_Forward_SSE(t *testing.T) {
// 模拟 vLLM chunked 发送
upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
flusher := w.(http.Flusher)
for _, chunk := range chunks {
w.Write([]byte(chunk))
flusher.Flush()
time.Sleep(20 * time.Millisecond)
}
}))
defer upstream.Close()
...
}
② 验证客户端断开时的级联取消(GPU 止损单测):模拟下游客户端在读取第 1 个 chunk 后主动执行 cancel(),上游后端必须在限定时间内接收到 <-r.Context().Done():
func TestStreamer_Forward_ClientCancel(t *testing.T) {
upstreamCanceled := make(chan struct{})
upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
...
select {
case <-r.Context().Done():
close(upstreamCanceled) // 上游成功捕获取消!
return
case <-time.After(3 * time.Second):
t.Errorf("upstream did not receive context cancellation")
}
}))
...
// 客户端读完首包后断开
resp.Body.Read(buf)
cancel()
<-upstreamCanceled // 断言上游收到信号
}
执行竞态检测:
go test -v -race ./internal/proxy/...
全部测试通过,且 -race 下没有数据竞态。
8. 总结与架构思考
- 生命周期优先。大模型算力贵,请求取消不只是前端体验问题,也是算力损耗。把客户端
r.Context()严密级联传导到底层连接,是推理网关的基本功。 - 用标准库的小接口组合。不要手写 buffer 循环。自定义
flushWriter实现io.Writer,再交给io.Copy,性能和内存池复用交给标准库,维护面更小。 - 协议细节决定并发表现。分块传输编码(Chunked Transfer Encoding)、TCP 滑动窗口、连接池
MaxIdleConns,高并发性能瓶颈往往藏在网络协议与 I/O 调度的细节里。