Go 基础体系 · 第 47/113 篇。示例统一基于 Go 1.26.4;核心片段可能省略 package 与 import,完整程序可直接按文中结构运行。

Go WebSocket 与 SSE:实时通信、心跳、背压和断线恢复

本文以 Go 1.26.4coder/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/WriterClose(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-streamCache-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 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。