Go SSE 代理模式

业务后端作为 SSE 代理透传上游流式响应:正确事件头、禁用缓冲、逐块 Flush、客户端断开检测、超时与错误映射,避免整包缓冲导致「假流式」。

#type / concept #status / growing #tech / dev #tech / dev / backend #resource / go

[!info] 关联笔记

Go SSE 代理模式

这个概念为什么会出现

浏览器或 App 不宜直连上游 AI/模型/内部流式服务,常见原因:

  • 上游密钥不能下发到客户端
  • 需要在边缘做登录态、配额、审计、上下文拼装
  • 需要统一公网域名与 CORS 策略
  • 需要屏蔽上游错误细节或做协议适配

于是业务后端成为 SSE(Server-Sent Events)代理:对客户端表现为 text/event-stream,对上游则是流式 HTTP 客户端。

难点在于:代理普通 JSON 可以整包读完再写;代理 SSE 必须:

  1. 尽快把上游字节送到客户端
  2. 禁止网关/框架/中间件缓冲整响应
  3. 客户端断开时停止拉上游,避免浪费配额
  4. 正确处理超时、非 200、中途断流

[!abstract] 一句话理解 SSE 代理是「带身份与策略的流式透传」:设好事件流响应头,按块读取上游 body,每次写入后 Flush,并用 Request.Context() 感知客户端离开。

最小可运行示例

先把示例放进业务场景,再看代码:

场景:聊天流式接口代理上游 AI SSE

产品页 GET /api/chat/stream 不能把模型 API Key 下发给浏览器。
业务后端做 SSE 代理

  1. 对前端:text/event-stream 真流式
  2. 对上游:带服务端 Token 拉流
  3. 用户关页 → 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

结合场景再看四个关注点

  1. 先断言 http.Flusher
  2. 上游请求绑定 r.Context(),客户端断开则取消上游。
  3. 禁止 io.ReadAll 后再一次性写出。
  4. 每一可交付单元(行/事件)后 Flush

核心概念与准确模型

SSE 在 HTTP 上的形态

  • 长连接响应,Content-Type: text/event-stream
  • 文本帧:event: / data: / 空行分隔等(见 sse
  • 方向以服务器 → 客户端为主(与 WebSocket 全双工不同)

代理通常有两种策略:

策略做法适用
字节/行透传上游已是 SSE,原样转发AI 网关、已统一协议
协议适配上游 chunked JSON/gRPC stream,本地组装为 SSE上游协议不对外

本篇聚焦透传;适配时仍遵守 Flush 与取消。

为何「缓冲」是第一大敌

缓冲可能来自:

  1. 应用代码 ReadAll / 拼完 string 再写
  2. 中间件 gzip、某些日志包装 Writer
  3. 反向代理(nginx proxy_buffering
  4. 云 LB / CDN 默认策略

结果:前端要等整段生成结束才突然刷出——假流式

缓解:

  • 应用层:分块写 + Flush
  • nginx:X-Accel-Buffering: noproxy_buffering off(按部署确认)
  • 中间件:SSE 路径跳过压缩与全量 body 中间件

客户端断开检测

http.Request.Context() 在客户端离开或服务取消请求时 Done。代理循环必须:

select {
case <-r.Context().Done():
    return
default:
}
// 或绑定到 NewRequestWithContext,使上游 Read 可中断(取决于实现)

不检测的后果:用户关页后仍继续消耗上游 token/GPU。

超时怎么设

层级建议
http.ServerReadHeaderTimeout 等防慢连接;流式接口的 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 的包装。机制与标准库相同。

设计动机

  1. 密钥与策略留在服务端
  2. 用户体验要真流式,而不是假进度
  3. 成本控制:断开即停
  4. 可观测:request id 贯穿接入层与上游调用

边界情况与反直觉行为

  1. WriteHeader(200) 再发现上游失败
    无法再改成 JSON 502;只能在流里发 error 事件。可先探测上游或采用「缓冲首包决策」的有限策略(牺牲首字节延迟)。
  2. bufio.Scanner 默认 token 上限
    超长 data: 行会失败,需 Buffer
  3. 包装 ResponseWriter 丢失 Flusher
    运行期才发现「不流式」。
  4. gzip 中间件
    可能缓冲或破坏事件边界。
  5. 多 goroutine 写同一 ResponseWriter
    无同步则数据竞态;透传循环应单写者。
  6. 上游也是代理
    超时与缓冲问题会叠加,需端到端压测。

常见误区

[!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-streamc.JSON(502, ...)
正确:提交响应模式前完成分支决策,或全程走 SSE 事件。

工程实践

  1. 路径级中间件:SSE 跳过 gzip/缓存。
  2. 鉴权在建立上游连接之前完成。
  3. 结构化日志:开始、首字节时间、断开原因、上游状态码、bytes。
  4. 限流与配额:按用户在代理层拒绝,避免打爆上游。
  5. 心跳:防中间设备空闲断连。
  6. 测试:用 httptest 难完整验证 Flush;应用集成测试或人工/脚本读流。
  7. 安全:不要把上游原始错误与密钥回传浏览器。
  8. 可配置上游 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

自测题

概念题

  1. 为什么 X-Accel-Buffering: no 不能单独保证浏览器实时显示?
  2. 为何上游请求要用 NewRequestWithContext(r.Context(), ...)
  3. 已经 WriteHeader(200)Content-Type: text/event-stream 后,上游返回 401,还能改成 HTTP 401 JSON 吗?

代码推理题

中间件按顺序包装:gzip(logging(sseHandler)),可能出现哪些流式故障?

工程思考题

产品要求「用户关闭页面后 1 秒内停止计费上游」。代理层要满足哪些条件才近似可达?

参考答案

展开
  1. 应用层若不 Flush 或中间件缓冲,仍假流式;该头只影响部分 nginx 行为。
  2. 客户端断开时取消上游,节约资源与费用。
  3. 不能可靠改写状态码与 Content-Type;只能流内报错或事先决策。
    代码:gzip 常缓冲;logging 包装 Writer 可能丢 Flusher。
    工程:ctx 绑定上游、上游尊重取消、循环可中断读、计费以实际上游用量并对取消做对账。

延伸阅读与资料来源

资料类型支撑
MDN — Server-sent events文档SSE 协议与事件格式
Package net/http标准库Flusher、Request.Context、Client
Package context标准库取消传播
Go Blog — context博客请求生命周期

笔记元信息

  • 建议文件名:go-sse-proxy-pattern.md
  • 所属阶段:阶段七(Web 后端 / 流式集成)
  • 本篇状态:已深化
  • 建议下一篇:健康检查端点优雅关闭
创建于 2026/6/25 更新于 2026/7/15