Go SSE 代理模式
业务后端作为 SSE 代理透传上游流式响应:正确事件头、禁用缓冲、逐块 Flush、客户端断开检测、超时与错误映射,避免整包缓冲导致「假流式」。
[!info] 关联笔记
- 所属 MOC:Go Web 与后端 · 学习路线
- 前置概念:net/http、context、中间件、SSE
- 并列 / 后续:透传模式、多服务 SSE 管道、前端 SSE 消费、优雅关闭
- 容易混淆:WebSocket 代理、整包
ReadAll、反向代理默认缓冲
Go SSE 代理模式
这个概念为什么会出现
浏览器或 App 不宜直连上游 AI/模型/内部流式服务,常见原因:
- 上游密钥不能下发到客户端
- 需要在边缘做登录态、配额、审计、上下文拼装
- 需要统一公网域名与 CORS 策略
- 需要屏蔽上游错误细节或做协议适配
于是业务后端成为 SSE(Server-Sent Events)代理:对客户端表现为 text/event-stream,对上游则是流式 HTTP 客户端。
难点在于:代理普通 JSON 可以整包读完再写;代理 SSE 必须:
- 尽快把上游字节送到客户端
- 禁止网关/框架/中间件缓冲整响应
- 客户端断开时停止拉上游,避免浪费配额
- 正确处理超时、非 200、中途断流
[!abstract] 一句话理解 SSE 代理是「带身份与策略的流式透传」:设好事件流响应头,按块读取上游 body,每次写入后
Flush,并用Request.Context()感知客户端离开。
最小可运行示例
先把示例放进业务场景,再看代码:
场景:聊天流式接口代理上游 AI SSE
产品页 GET /api/chat/stream 不能把模型 API Key 下发给浏览器。
业务后端做 SSE 代理:
- 对前端:
text/event-stream真流式 - 对上游:带服务端 Token 拉流
- 用户关页 →
r.Context()取消 → 停拉上游,省 GPU/配额
以下用标准库突出机制(生产可换 chi/Gin,Flush/ctx 原则不变)。
package main
import (
"bufio"
"fmt"
"io"
"log"
"net/http"
"time"
)
// sseProxy:把上游 AI 的 event-stream 透传给浏览器。
func sseProxy(w http.ResponseWriter, r *http.Request) {
// 必须能 Flush,否则只能整包缓冲 → 假流式
flusher, ok := w.(http.Flusher)
if !ok {
http.Error(w, "streaming unsupported", http.StatusInternalServerError)
return
}
// 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 不要缓冲
w.WriteHeader(http.StatusOK)
flusher.Flush() // 立刻把头打出去,前端 EventSource 可建立
// 2) 上游请求绑定客户端 ctx:关页即取消
req, err := http.NewRequestWithContext(r.Context(), http.MethodGet, "http://upstream.example/stream", nil)
if err != nil {
fmt.Fprintf(w, "event: error\ndata: %v\n\n", err)
flusher.Flush()
return
}
// 密钥只在服务端;此处注入,永不下发浏览器
req.Header.Set("Authorization", "Bearer "+serverSideToken())
client := &http.Client{
// 流式慎用「总 Timeout 盖整段生成」;
// 连接超时放 Transport,空闲靠 ctx 取消
Timeout: 0,
}
resp, err := client.Do(req)
if err != nil {
fmt.Fprintf(w, "event: error\ndata: upstream unavailable\n\n")
flusher.Flush()
return
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
// 错误体限长,避免把上游堆栈整段吐给前端
body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
fmt.Fprintf(w, "event: error\ndata: status %d %s\n\n", resp.StatusCode, string(body))
flusher.Flush()
return
}
// 3) 逐行透传 + 每行 Flush(禁止 ReadAll 后再写)
scanner := bufio.NewScanner(resp.Body)
// 大事件行需调大 buffer
scanner.Buffer(make([]byte, 64*1024), 1024*1024)
for scanner.Scan() {
if err := r.Context().Err(); err != nil {
return // 客户端已断开,停止读上游
}
line := scanner.Text()
_, _ = fmt.Fprintf(w, "%s\n", line)
flusher.Flush() // 每个可交付单元立刻推前端
}
if err := scanner.Err(); err != nil && r.Context().Err() == nil {
log.Printf("upstream read: %v", err)
}
}
// serverSideToken:演示用;生产从密钥管理读取。
func serverSideToken() string { return "secret" }
func main() {
mux := http.NewServeMux()
mux.HandleFunc("GET /api/chat/stream", sseProxy)
srv := &http.Server{
Addr: ":8080",
Handler: mux,
ReadHeaderTimeout: 5 * time.Second,
}
log.Fatal(srv.ListenAndServe())
}
建议运行:
go run .
# 需有真实上游或把 URL 换成本地 mock stream 后再:
# curl -N http://localhost:8080/api/chat/stream
期望:响应头含 text/event-stream,事件按行逐步到达(非等整段结束才刷出)。无真实上游时会立刻 event: error。
结合场景再看四个关注点
- 先断言
http.Flusher。 - 上游请求绑定
r.Context(),客户端断开则取消上游。 - 禁止
io.ReadAll后再一次性写出。 - 每一可交付单元(行/事件)后
Flush。
核心概念与准确模型
SSE 在 HTTP 上的形态
- 长连接响应,
Content-Type: text/event-stream - 文本帧:
event:/data:/ 空行分隔等(见 sse) - 方向以服务器 → 客户端为主(与 WebSocket 全双工不同)
代理通常有两种策略:
| 策略 | 做法 | 适用 |
|---|---|---|
| 字节/行透传 | 上游已是 SSE,原样转发 | AI 网关、已统一协议 |
| 协议适配 | 上游 chunked JSON/gRPC stream,本地组装为 SSE | 上游协议不对外 |
本篇聚焦透传;适配时仍遵守 Flush 与取消。
为何「缓冲」是第一大敌
缓冲可能来自:
- 应用代码
ReadAll/ 拼完 string 再写 - 中间件 gzip、某些日志包装 Writer
- 反向代理(nginx
proxy_buffering) - 云 LB / CDN 默认策略
结果:前端要等整段生成结束才突然刷出——假流式。
缓解:
- 应用层:分块写 +
Flush - nginx:
X-Accel-Buffering: no与proxy_buffering off(按部署确认) - 中间件:SSE 路径跳过压缩与全量 body 中间件
客户端断开检测
http.Request.Context() 在客户端离开或服务取消请求时 Done。代理循环必须:
select {
case <-r.Context().Done():
return
default:
}
// 或绑定到 NewRequestWithContext,使上游 Read 可中断(取决于实现)
不检测的后果:用户关页后仍继续消耗上游 token/GPU。
超时怎么设
| 层级 | 建议 |
|---|---|
http.Server | ReadHeaderTimeout 等防慢连接;流式接口的 WriteTimeout 要谨慎(可能砍长流) |
上游 Client.Timeout | 覆盖整个请求生命周期;长生成常设 0,改用 ctx / 空闲超时 |
context.WithTimeout | 给「最大会话时长」硬上限 |
| Transport | 拨号、TLS、响应头超时与连接池 |
原则:区分「连不上」与「已经在流式生成」。
错误映射
- 上游不可达:对客户端可用 SSE
event: error或在尚未写入流式头之前返回 JSON 502(二选一,勿混用到一半) - 上游 4xx/5xx:读有限错误体,映射为可观测日志 + 对用户安全的信息
- 中途断流:结束响应;前端按重连策略(SSE 自带
Last-Event-ID时需代理是否支持续传——多数 AI 流不支持透明续传)
与反向代理、K8s 的配合
- 空闲超时要 大于 上游心跳间隔;可定期发送 SSE 注释行
: ping - 负载均衡「整段缓冲」类设置必须关
- 优雅关闭时,in-flight SSE 应在宽限期内自然结束或被 ctx 取消,见 go-graceful-shutdown
Gin 场景注意点
Gin 中仍须:
c.Writer.Header().Set("Content-Type", "text/event-stream")
// ...
c.Writer.Flush()
// 使用 c.Request.Context()
且避免在中间件里对 Writer 做破坏 Flusher 的包装。机制与标准库相同。
设计动机
- 密钥与策略留在服务端
- 用户体验要真流式,而不是假进度
- 成本控制:断开即停
- 可观测:request id 贯穿接入层与上游调用
边界情况与反直觉行为
- 先
WriteHeader(200)再发现上游失败
无法再改成 JSON 502;只能在流里发 error 事件。可先探测上游或采用「缓冲首包决策」的有限策略(牺牲首字节延迟)。 bufio.Scanner默认 token 上限
超长data:行会失败,需Buffer。- 包装 ResponseWriter 丢失 Flusher
运行期才发现「不流式」。 - gzip 中间件
可能缓冲或破坏事件边界。 - 多 goroutine 写同一 ResponseWriter
无同步则数据竞态;透传循环应单写者。 - 上游也是代理
超时与缓冲问题会叠加,需端到端压测。
常见误区
[!warning] 常见误区:ReadAll 再当 SSE 返回 错误:
body, _ := io.ReadAll(resp.Body); w.Write(body)
正确:边读边写边 Flush。
[!warning] 常见误区:忽略客户端断开 错误:只
for scanner.Scan()不看 ctx。
正确:NewRequestWithContext(r.Context(), ...)+ 循环检查。
[!warning] 常见误区:用过短的 Client.Timeout 切长对话 错误:60s 超时砍掉正常长生成。
正确:会话级 ctx 上限 + 连接级超时分离。
[!warning] 常见误区:错误时混用 JSON 与已开始的 SSE 错误:已发
text/event-stream又c.JSON(502, ...)。
正确:提交响应模式前完成分支决策,或全程走 SSE 事件。
工程实践
- 路径级中间件:SSE 跳过 gzip/缓存。
- 鉴权在建立上游连接之前完成。
- 结构化日志:开始、首字节时间、断开原因、上游状态码、bytes。
- 限流与配额:按用户在代理层拒绝,避免打爆上游。
- 心跳:防中间设备空闲断连。
- 测试:用
httptest难完整验证 Flush;应用集成测试或人工/脚本读流。 - 安全:不要把上游原始错误与密钥回传浏览器。
- 可配置上游 URL/模型:但凭证仅环境变量/密钥管理。
反例 vs 正例
// ❌ 假流式
body, _ := io.ReadAll(resp.Body)
w.Header().Set("Content-Type", "text/event-stream")
_, _ = w.Write(body)
// ✅ 真流式
// 见上文 scanner + Flush + ctx
本节总结
- SSE 代理 = 服务端持密 + 流式透传 + 策略执行。
- 三件套:正确头、分块 Flush、ctx 取消。
- 缓冲可能来自代码、中间件、nginx,需分层关闭。
- 超时要区分建连与长生成。
- 下一步:前端消费 frontend-sse-consumption;多段管道 multi-service-sse-pipeline。
自测题
概念题
- 为什么
X-Accel-Buffering: no不能单独保证浏览器实时显示? - 为何上游请求要用
NewRequestWithContext(r.Context(), ...)? - 已经
WriteHeader(200)且Content-Type: text/event-stream后,上游返回 401,还能改成 HTTP 401 JSON 吗?
代码推理题
中间件按顺序包装:gzip(logging(sseHandler)),可能出现哪些流式故障?
工程思考题
产品要求「用户关闭页面后 1 秒内停止计费上游」。代理层要满足哪些条件才近似可达?
参考答案
展开
- 应用层若不 Flush 或中间件缓冲,仍假流式;该头只影响部分 nginx 行为。
- 客户端断开时取消上游,节约资源与费用。
- 不能可靠改写状态码与 Content-Type;只能流内报错或事先决策。
代码:gzip 常缓冲;logging 包装 Writer 可能丢 Flusher。
工程:ctx 绑定上游、上游尊重取消、循环可中断读、计费以实际上游用量并对取消做对账。
延伸阅读与资料来源
| 资料 | 类型 | 支撑 |
|---|---|---|
| MDN — Server-sent events | 文档 | SSE 协议与事件格式 |
| Package net/http | 标准库 | Flusher、Request.Context、Client |
| Package context | 标准库 | 取消传播 |
| Go Blog — context | 博客 | 请求生命周期 |