Go 基础体系 · 第 47/113 篇。示例统一基于 Go 1.26.4;核心片段可能省略 package 与 import,完整程序可直接按文中结构运行。
Go WebSocket 与 SSE:实时通信、心跳、背压和断线恢复
本文以 Go 1.26.4 与 coder/websocket v1.8.15 为基准。WebSocket 在 HTTP 握手后切换成全双工消息协议,SSE 则始终是 text/event-stream HTTP 响应。聊天、协同和双向控制适合 WebSocket;通知、进度和文本流只有下行数据时,SSE 的浏览器接口、代理兼容性和恢复语义更简单。
低频确认可用普通 HTTP 写入加 SSE 下行;高频双向小消息适合 WebSocket。
1. 两种协议的线格式与能力边界
WebSocket 握手从 HTTP/1.1 GET 开始,客户端发送 Upgrade: websocket、随机 Sec-WebSocket-Key 和版本,服务端以 101 Switching Protocols 接受。此后 TCP 字节流承载 WebSocket frame;文本、二进制、Ping、Pong、Close 都有明确帧类型,消息可由多个帧组成。浏览器客户端发送的帧必须掩码,服务端帧不掩码。应用操作的是“消息”,不能假设一次底层 Read 对应一个完整业务对象。
GET /ws HTTP/1.1
Host: api.example.com
Connection: Upgrade
Upgrade: websocket
Sec-WebSocket-Version: 13
Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==
Origin: https://app.example.com
SSE 响应使用 UTF-8 文本。一个事件由若干字段行和空行结束;连续 data: 行以换行拼接,id: 更新浏览器保存的最后事件 ID,retry: 给出建议重连毫秒数,以冒号开头的是注释心跳。它没有客户端到服务端的事件通道,也不适合直接传二进制;二进制需要编码,会增加体积。
id: 1842
event: article.updated
data: {"id":"a-7","version":12}
: heartbeat
retry: 3000
data: ready
2. 生命周期:握手成功只是开始
一个实时连接通常经历:HTTP 入口鉴权与配额检查、协议建立、注册到连接表、读写循环运行、主动或异常关闭、注销和资源清理。任何一步失败都必须只清理自己已经取得的资源。认证应在升级或写出 SSE 200 前完成,因为响应开始后已经无法改成规范的 401 JSON。
WebSocket 的请求 context 通常不能被当作连接的永久业务 context:升级返回后 handler 生命周期和连接生命周期的关系取决于库。应用应显式创建连接 context,并在读循环退出、写失败、服务下线或账户撤销时取消它。SSE 仍在 handler 中运行,r.Context() 会在客户端断开或服务器关闭相关请求时取消,循环中的每个阻塞点都要选择它。
关闭 WebSocket 时先停止接收新业务消息,再给写循环一个有界时间发送 Close,最后关闭网络连接。SSE 没有协议 Close 帧;handler 返回就是结束,浏览器默认会重连。若服务端希望客户端停止永久重连,需要用业务事件通知客户端调用 EventSource.close(),HTTP 状态在已经开始流式响应后无能为力。
3. coder/websocket 的关键 API
manifest 指定的 github.com/coder/websocket 以 context 控制 I/O。服务端用 websocket.Accept(w, r, options),客户端用 websocket.Dial(ctx, url, options);Read(ctx) 返回消息类型和完整消息,Write(ctx, type, data) 写一条消息,流式大消息可用 Reader/Writer。Close(status, reason) 做正常关闭,CloseNow() 用于无法继续握手的兜底清理。
c, err := websocket.Accept(w, r, &websocket.AcceptOptions{
OriginPatterns: []string{"app.example.com"},
})
if err != nil {
return
}
defer c.CloseNow()
c.SetReadLimit(64 << 10)
ctx, cancel := context.WithCancel(r.Context())
defer cancel()
for {
typ, payload, err := c.Read(ctx)
if err != nil {
status := websocket.CloseStatus(err)
if status != websocket.StatusNormalClosure && status != websocket.StatusGoingAway {
slog.Warn("websocket read", "status", status, "err", err)
}
return
}
if typ != websocket.MessageText {
_ = c.Close(websocket.StatusUnsupportedData, "text only")
return
}
_ = payload // 校验后交给领域服务
}
读大小限制必须在读取不可信消息前设置;压缩不是大小限制的替代品,解压后的消息仍可能消耗大量内存。websocket.CloseStatus(err) 用于区分正常离开、策略拒绝、消息过大和异常断线。不要把底层错误字符串原样发给客户端,Close reason 受长度限制且可能进入浏览器日志。
4. 一读一写、串行写入与背压
可靠结构通常是每连接一个读循环、一个写循环。多个业务 goroutine 不应无序并发写同一连接,而应把不可变消息投递到有界队列,由唯一写循环串行发送。队列满表示客户端消费速度落后于生产速度,必须选择可解释策略:丢弃可替代的最新状态、合并进度、拒绝生产者,或关闭慢客户端。无界 channel 只是把网络背压变成堆内存增长。
type peer struct {
send chan []byte
cancel context.CancelFunc
}
func (p *peer) offer(message []byte) bool {
copyOfMessage := append([]byte(nil), message...)
select {
case p.send <- copyOfMessage:
return true
default:
p.cancel() // 本系统选择断开慢客户端
return false
}
}
复制消息避免发布者随后复用缓冲区引发数据竞争或内容变化。广播中心也应由单一 goroutine 或锁保护连接集合;取消与注销必须幂等,不能一边遍历 map 一边由其他 goroutine 删除。队列容量按“允许落后的时间 × 峰值消息率 × 平均消息大小”估算,并用真实慢网测试验证,而不是随手填一个大数。
5. 心跳、期限与半开连接
TCP 连接在拔网线、NAT 丢状态或代理静默回收时,不一定立刻报错。WebSocket Ping/Pong 用于检测应用链路活性。coder/websocket 的默认读取路径会处理控制帧;Ping(ctx) 等待对应 Pong,给它设置短于心跳间隔的 timeout。Ping 失败就取消连接,不要无限重试同一条坏链路。
ticker := time.NewTicker(25 * time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
pingCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
err := c.Ping(pingCtx)
cancel()
if err != nil {
return fmt.Errorf("ping: %w", err)
}
}
}
SSE 没有 Ping 帧,通常每 15 至 30 秒写一条 : heartbeat\n\n 并 Flush,间隔小于链路空闲超时。心跳只证明写入路径可用;确认语义仍需客户端 ack。
6. SSE 正确写法、Flush 与恢复
SSE handler 应设置 Content-Type: text/event-stream、Cache-Control: no-cache,并根据代理约定关闭缓冲。Go 1.26.4 可用 http.NewResponseController(w).Flush(),失败时结束流;只做 fmt.Fprintf 而不 Flush,数据可能长期留在服务端或代理缓冲区。
func writeEvent(w io.Writer, id, event string, data []byte) error {
if strings.ContainsAny(id, "\r\n") || strings.ContainsAny(event, "\r\n") {
return errors.New("invalid SSE field")
}
if id != "" {
fmt.Fprintf(w, "id: %s\n", id)
}
if event != "" {
fmt.Fprintf(w, "event: %s\n", event)
}
for _, line := range strings.Split(string(data), "\n") {
fmt.Fprintf(w, "data: %s\n", line)
}
_, err := io.WriteString(w, "\n")
return err
}
浏览器重连时会把最近的 id 放进 Last-Event-ID 请求头。服务端需要按 ID 从有界历史、数据库日志或消息流重放缺口,再切换到实时订阅;仅重复发送当前状态则要明确这是快照语义。ID 必须单调且在重启、多实例间可解释。若请求的 ID 已超出保留窗口,应发送 reset 快照事件,而不是悄悄漏数据。先订阅实时流再读取历史,或用存储游标原子衔接,避免“读完历史到开始订阅”之间丢事件。
7. 可运行的综合服务
下面核心结构同时提供 SSE、普通 HTTP 发布入口和可测试的有界 Hub。生产中可把发布入口换成领域事件,把内存历史换成 Redis Streams、NATS JetStream 或数据库 outbox;连接仍由当前进程拥有。
type Event struct {
ID uint64 `json:"id"`
Text string `json:"text"`
}
type Hub struct {
mu sync.Mutex
nextID uint64
history []Event
clients map[chan Event]struct{}
}
func NewHub() *Hub { return &Hub{clients: make(map[chan Event]struct{})} }
func (h *Hub) Publish(text string) Event {
h.mu.Lock()
defer h.mu.Unlock()
h.nextID++
event := Event{ID: h.nextID, Text: text}
h.history = append(h.history, event)
if len(h.history) > 100 {
h.history = append([]Event(nil), h.history[len(h.history)-100:]...)
}
for ch := range h.clients {
select {
case ch <- event:
default:
close(ch)
delete(h.clients, ch)
}
}
return event
}
func (h *Hub) Subscribe(after uint64) ([]Event, <-chan Event, func()) {
h.mu.Lock()
defer h.mu.Unlock()
backlog := make([]Event, 0, len(h.history))
for _, event := range h.history {
if event.ID > after {
backlog = append(backlog, event)
}
}
ch := make(chan Event, 16)
h.clients[ch] = struct{}{}
var once sync.Once
stop := func() {
once.Do(func() {
h.mu.Lock()
defer h.mu.Unlock()
if _, ok := h.clients[ch]; ok {
delete(h.clients, ch)
close(ch)
}
})
}
return backlog, ch, stop
}
func (h *Hub) Events(w http.ResponseWriter, r *http.Request) {
after, _ := strconv.ParseUint(r.Header.Get("Last-Event-ID"), 10, 64)
backlog, events, stop := h.Subscribe(after)
defer stop()
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
w.Header().Set("X-Accel-Buffering", "no")
controller := http.NewResponseController(w)
send := func(event Event) error {
data, err := json.Marshal(event)
if err != nil { return err }
if err := writeEvent(w, strconv.FormatUint(event.ID, 10), "message", data); err != nil { return err }
return controller.Flush()
}
for _, event := range backlog {
if err := send(event); err != nil { return }
}
heartbeat := time.NewTicker(20 * time.Second)
defer heartbeat.Stop()
for {
select {
case <-r.Context().Done():
return
case event, ok := <-events:
if !ok || send(event) != nil { return }
case <-heartbeat.C:
if _, err := io.WriteString(w, ": heartbeat\n\n"); err != nil { return }
if controller.Flush() != nil { return }
}
}
}
注意示例在锁内对满队列执行关闭和删除,消费者只读 channel;stop 再次执行时会检查存在性,因此不会重复关闭。若发布量很大,不应在全局锁内遍历几十万连接,可按主题分片 Hub,或让发布 goroutine 独占状态。历史切片的复制用于及时释放旧底层数组引用,实际持久恢复要依赖外部日志。
8. 错误、取消和超时如何分类
建立连接前的认证失败、限流和参数错误使用普通 HTTP 状态。连接建立后,WebSocket 用标准 Close status 表达类别:1000 正常完成,1001 服务下线,1008 策略或授权失败,1009 消息过大,1011 服务端意外错误。不要发送保留码,也不要把所有错误都归为 1006;它只用于观察到的异常关闭,不能作为 Close 帧发送。
超时至少分四类:握手超时、读空闲/心跳超时、单次写超时、整体业务超时。它们不应共用一个从建立时就开始倒计时的 context,否则长连接会按固定寿命被误杀。每次 I/O 从连接 context 派生短 timeout,用后立即 cancel。客户端取消、正常关闭通常记低级别日志;协议错误、过载断开、内部错误分别计数,避免正常离线污染告警。
9. 测试:不仅断言能连上
SSE 可用 httptest.NewServer 和真实 http.Client 读取多行事件,验证响应头、Flush、Last-Event-ID 重放与取消后订阅数归零。httptest.ResponseRecorder 不足以证明真实流式时序。WebSocket 测试使用 websocket.Dial 连接测试服务器,覆盖文本/二进制限制、超大消息、正常 Close、心跳超时和并发广播。
Hub 属于并发状态机,应在 go test -race 下反复执行发布、订阅和取消;刻意构造容量为 1 的慢客户端,确认不会阻塞其他连接且只关闭一次。网络级测试还要经过与生产一致的 ingress,验证 Upgrade 头、HTTP 版本、代理缓冲、空闲超时和滚动发布行为。
go test -race ./...
curl -N -H 'Last-Event-ID: 40' http://127.0.0.1:8080/events
websocat -H='Origin: https://app.example.com' ws://127.0.0.1:8080/ws
10. 性能与容量规划
长连接的主要预算不是路由 QPS,而是文件描述符、每连接 goroutine 栈、读写缓冲、应用队列、TLS 状态和内核 socket 缓冲。估算总内存时将这些项乘连接数,再加入广播瞬时分配和 GC 余量。消息编码一次后以只读字节广播,能减少 JSON 重复编码;但为了隔离生命周期可能仍需复制或引用计数,必须以 profile 证明。
压测要模拟快慢客户端、断线风暴、重连抖动和大消息。客户端重连使用指数退避加随机抖动。WebSocket 压缩能省带宽但消耗 CPU;对已压缩或很小消息常得不偿失。
11. 安全边界
浏览器 WebSocket 不受普通 CORS 流程保护,服务端必须校验 Origin,只允许明确域名;非浏览器客户端再通过 token 或 mTLS 认证。查询参数里的 token 容易进入访问日志,优先短期一次性票据、受保护 cookie 或应用层首消息认证,并限制认证前能读取的数据和等待时间。
SSE 使用普通 HTTP,仍需防 CSRF 风格的跨站读取、租户越权和缓存泄漏。设置准确的 CORS、Cache-Control: no-store(敏感流)、认证 cookie 属性与每用户连接上限。两种协议都要限制消息大小、主题订阅数、发布频率和总连接数;业务负载先做 schema 校验,HTML 前端展示时按数据转义,不能因为来自“可信连接”就拼进 DOM。
12. 多实例与生产部署边界
负载均衡器必须支持 WebSocket Upgrade 和长时间流式响应,关闭 SSE 响应缓冲,并把空闲超时设得高于心跳间隔。连接一旦建立就属于某个实例;一致哈希或 sticky session 可以提高局部性,却不能替代共享事件存储。跨实例广播由消息系统分发,每个实例只投递给本地订阅者。消息系统故障时要明确丢弃、断开还是降级到轮询。
滚动下线分为摘流量和排连接:先让 readiness 失败并停止接受新连接,广播 server_draining 或以 1001 关闭,给客户端随机重连提示,在总 shutdown deadline 内等待写循环,最后强制取消。监控至少包括当前连接、建立/关闭原因、认证失败、队列满、消息大小、发送延迟、历史重放量和每实例文件描述符。不要记录完整私密消息或 token。
最终设计原则很直接:协议负责传输,应用负责恢复、授权、背压与一致性。 WebSocket 的双工能力并不自动带来可靠投递,SSE 的自动重连也不自动补齐事件;只有明确消息 ID、队列上限、取消路径、存储窗口和部署链路,两种实时通道才具备可运营性。
系列导航与关联阅读
- 系列入口:Go 完整技术体系学习路线:从语法、并发到框架、中间件与 AI
- 上一篇:Go gRPC 与 Protobuf 完整基础:IDL、Unary、Stream 与拦截器
- 下一篇:Go GraphQL 与 gqlgen:Schema、Resolver、DataLoader 和复杂度控制
- 延伸:Go net/http 基础:Server、Handler、Middleware 与 Client 超时
- 延伸:Go channel 完整基础:发送、接收、缓冲、关闭与所有权
- 延伸:Go Redis 生产模式:缓存一致性、穿透击穿、锁与 Streams
官方资料
本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。

评论
0 条讨论