HTTP 一问一答,服务器不能主动推。聊天消息、行情跳动、任务进度这类场景,靠前端轮询既慢又费流量。WebSocket 在一条 TCP 上双向通信,握手走一次 HTTP,之后服务端随时发。这一篇在 Gin 里接上它,做一个可广播的房间。
1. 升级连接
用 github.com/gorilla/websocket,把进来的 HTTP 请求升级成 WebSocket:
import "github.com/gorilla/websocket" var upgrader = websocket.Upgrader{ CheckOrigin: func(r *http.Request) bool { return true // 生产按 Origin 白名单校验,别直接放行 }, } func (h *WSHandler) Room(c *gin.Context) { conn, err := upgrader.Upgrade(c.Writer, c.Request, nil) if err != nil { return // Upgrade 已写过响应,出错直接返回 } client := &Client{conn: conn, send: make(chan []byte, 256)} hub.register <- client defer func() { hub.unregister <- client }() go client.writePump() client.readPump() }
Upgrade 成功后 conn 接管了这条连接,Gin 的响应机制就退场了,后面读写全走 conn。
2. 用一个 Hub 做广播
多个客户端要互相同步,集中到一个 Hub 管理。Hub 跑在自己的 goroutine 里,靠 channel 收注册、注销、广播消息,避免并发改 map:
type Hub struct { clients map[*Client]bool register chan *Client unregister chan *Client broadcast chan []byte } func (h *Hub) Run() { for { select { case c := <-h.register: h.clients[c] = true case c := <-h.unregister: if _, ok := h.clients[c]; ok { delete(h.clients, c) close(c.send) } case msg := <-h.broadcast: for c := range h.clients { select { case c.send <- msg: default: // 发送缓冲满了,踢掉慢客户端 close(c.send) delete(h.clients, c) } } } } }
所有状态变更只发生在 Run 这一个 goroutine,map 不用额外加锁。
3. 客户端读写分两个协程
一个连接要同时收和发,Go 里各开一个 goroutine。读协程收消息后转交广播,写协程从 send channel 取消息往外发:
func (c *Client) readPump() { defer c.conn.Close() for { _, msg, err := c.conn.ReadMessage() if err != nil { return } hub.broadcast <- msg } } func (c *Client) writePump() { defer c.conn.Close() for msg := range c.send { if err := c.conn.WriteMessage(websocket.TextMessage, msg); err != nil { return } } }
readPump 退出(defer 关连接)会让 writePump 的 range c.send 也结束,两条协程一起收尾。
4. 心跳保活
中间网络会静默断开,双方都不知道。客户端定时发 ping,服务端回 pong,SetReadDeadline 兜底,超时没收到就断:
func (c *Client) readPump() { defer c.conn.Close() c.conn.SetReadDeadline(time.Now().Add(60 * time.Second)) c.conn.SetPongHandler(func(string) error { c.conn.SetReadDeadline(time.Now().Add(60 * time.Second)) return nil }) for { _, msg, err := c.conn.ReadMessage() if err != nil { return } hub.broadcast <- msg } }
写协程里定时发 ping:
func (c *Client) writePump() { ticker := time.NewTicker(30 * time.Second) defer func() { ticker.Stop(); c.conn.Close() }() for { select { case msg, ok := <-c.send: if !ok { c.conn.WriteMessage(websocket.CloseMessage, []byte{}) return } c.conn.WriteMessage(websocket.TextMessage, msg) case <-ticker.C: c.conn.WriteControl(websocket.PingMessage, nil, time.Now().Add(5*time.Second)) } } }
5. 鉴权
升级前先校验登录态,别让匿名连接进房间:
func (h *WSHandler) Room(c *gin.Context) { token := c.Query("token") if !validToken(token) { c.JSON(401, gin.H{"error": "未授权"}) return } conn, err := upgrader.Upgrade(c.Writer, c.Request, nil) ... }
token 放查询参数或 Sec-WebSocket-Protocol 子协议头里带过来,升级请求是普通 HTTP,能走第 6 篇的校验逻辑。
6. nginx 反代要放行升级
WebSocket 握手带 Upgrade 头,nginx 默认不透,会卡在 101 之前:
location /ws/ {
proxy_pass http://127.0.0.1:8080;
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_read_timeout 3600s;
proxy_set_header X-Real-IP $remote_addr;
}
Connection "upgrade" 和 Upgrade 头是握手关键。proxy_read_timeout 调长,否则空闲连接被 nginx 提前断。
7. 上线前清单
CheckOrigin 按域名白名单校验,别放行所有源。Hub 集中管理连接,状态变更只在一个 goroutine。发送缓冲满的慢客户端直接踢,防止一个卡连接拖垮广播。心跳 30 秒 ping、60 秒读超时兜底断连。连接数多的服务,单个 Hub 的 map 要分片或拆房间,否则广播遍历会成瓶颈。nginx 必须配 Upgrade 头和读超时。
下一篇讲可观测性,让线上出问题能靠日志和链路定位,而不是盲猜。