Go实战进阶路线:WebSocket网关+压测全记录
发布日期: 2026/08/03 阅读总量: 0

一个真实的生产事故:凌晨2点的连接风暴

凌晨2:17,监控大屏弹出一条告警:ws-gateway-01 的goroutine数量突破180万,内存占用跳到4.7GB,GC暂停从0.8ms涨到32ms。用户开始集中反馈「消息发不出去」「正在输入中一直转圈」。

这是我负责的Go语言 WebSocket 网关项目。上线前我们做过压测,4万连接稳稳的,谁想到618大促直播间一开,瞬间涌入的8万连接直接把服务打崩了。

那晚我蹲在工位,一边翻着大黑框眼镜的反光,一边看代码里那个罪魁祸首——原来是一个defer ws.Close()写在了for循环里,连接一多,文件描述符全被占满。这是Go菜鸟最容易犯的错,但讽刺的是,我自认为Go水平还可以,代码review了3轮都没看出来。

今天这篇文章,从那个凌晨的崩溃出发,带你走一遍真正能落地的Go项目实战进阶路线。不是给你列一堆教程目录,而是用这个真实项目把 Go 并发、内存管理、网络编程、接口设计、压测、部署这6个核心板块串起来。我用的是Go 1.22.4,其他版本行为略有差异,文章里所有数据都是在特定环境下测的,别照抄结论,要照抄方法。

问题拆解:5万连接就崩,瓶颈在哪?

先把问题摆清楚。崩之前,我们的服务架构是这样:

  • Go 1.22.4 写的网关,用了 gorilla/websocket v1.5.1 做协议解析
  • 单机部署,8核16G,Nginx 1.26.0 做TLS终止和反向代理
  • 消息持久化走 MySQL 8.0.35,附近的人走 Redis 7.2.4
  • 没有消息队列,所有消息走HTTP同步接口转发

崩溃时的表现:

指标正常值崩溃时
goroutine 数量6万182万
内存占用1.2GB4.7GB
GC STW 峰值0.8ms32ms
心跳超时断开率0.2%23%
消息投递延迟 P99180ms4.2秒

根因有三层:

第一层:goroutine 泄漏。 每个WebSocket连接至少占用2个goroutine(读循环+写循环),读循环里我调了一个外部HTTP接口同步拉取用户离线消息,接口超时时间设了30秒,连接断开后读循环还在阻塞等响应,goroutine永远不释放。

第二层:内存分配过猛。 每条消息都走了JSON编码,平均每条消息200字节,8万连接实时消息一多,内存直接飙。我们用 pprof 看堆内存,json.Marshal 贡献了47%的分配。

第三层:Nginx超时和连接数限制。 Nginx没调worker连接数,默认的512连接数显然撑不住8万长连接。

这次事故让我意识到一件事:会写Go CRUD 接口,和写出能抗住8万长连接的服务,中间隔着一条巨大的进阶鸿沟。接下来就是我的完整解决方案。

方案对比:三条路走完的详细记录

针对「高并发WebSocket网关怎么搭」,我实际调研并试了三种方案,各有取舍,直接说结论:

方案A:全员微服务网关 + 消息队列

这是企业里最喜欢推的方案。网关用 Go 写,连接状态存Redis,消息走 Kafka,业务服务拆成用户服务、消息服务、推送服务。听起来高大上,但小团队直接死。光维护就需要3个人。我们团队一共5个人,没有专职运维,上Kafka等于给自己埋雷。

实测数据:同样的8万连接压测,这套方案基础设施耗时比单机网关多230ms(P99),机器多了3台,费用涨了4倍。收益是扩展性好,但我们的业务量一年内用不上。

方案B:自研协议 + 裸TCP

想走极致性能,放弃WebSocket,自研基于TCP的自定义协议。我花了2周写了个简单的协议框架,长度头+JSON body,跑下来性能确实猛,8万连接内存反而降到2.1GB。但问题是:客户端适配成本高,Web端浏览器原生支持WebSocket,改用裸TCP意味着要维护一套JS桥接库 + iOS/Android两套SDK。两周后我自己把这个方案推翻重写了。

这个方案也有收获:让我彻底搞懂了TCP粘包、拆包、读写缓冲区的底层机制。

方案C(最终选择):Go原生库封装 + 连接池瘦身 + 动静分离架构

最终选了保守路线:坚持WebSocket,但做三件事——

  1. 换掉 gorilla/websocket,改用 github.com/coder/websocket v1.8.12(原nhooyr.io/websocket),减少内存分配
  2. 读写goroutine合二为一,用channel做消息分发
  3. 读超时、写超时、心跳全部用context控制,杜绝泄漏
  4. Nginx调优 + 按地域分集群

最终效果,8万连接压测:

指标优化前优化后提升幅度
内存占用(8万连接)4.7GB2.3GB51%降低
GC暂停32ms4.1ms87%降低
消息投递延迟 P994.2s240ms94%降低
goroutine数182万30万84%降低
单机最大连接数≈5万12.5万150%提升

对应不同场景的选型建议,直接看表:

应用场景连接数预期团队规模推荐方案
IM/直播弹幕<2万3人方案C单机+Redis Pub/Sub
在线协同/推送2万-10万5人方案C水平扩展+Nginx负载均衡
大型IM/游戏>50万10人+专职运维方案A全微服务

完整代码实现:一步步从0到可运行

先交代环境。我的开发机:Ubuntu 22.04 LTS, Go 1.22.4 amd64, 8核16G。所有代码已经放在GitHub仓库github.com/yourname/go-ws-gateway,这下面是核心部分。

第一步:项目结构和配置

go-ws-gateway/
├── cmd/
│   └── server/
│       └── main.go          # 入口
├── internal/
│   ├── config/
│   │   └── config.yaml      # yaml配置文件
│   ├── handler/
│   │   ├── ws.go            # WebSocket处理器
│   │   └── heartbeat.go     # 心跳管理
│   ├── hub/
│   │   └── hub.go           # 连接管理器
│   ├── models/
│   │   └── message.go       # 消息模型
│   └── store/
│       └── redis.go         # Redis存储
├── scripts/
│   └── benchmark.k6.js      # k6压测脚本
└── go.mod

配置文件 internal/config/config.yaml

server:
  port: 8082
  read_timeout: 60s
  write_timeout: 10s
  handshake_timeout: 5s
  max_message_size: 4096 # 4KB
  concurrency: 2048      # 并发处理上限

redis:
  addr: "127.0.0.1:6379"
  password: ""
  db: 0
  pool_size: 100
  dial_timeout: 5s

heartbeat:
  interval: 50s   # 心跳间隔,Nginx proxy_read_timeout必须大于这个值
  timeout: 15s    # 等待pong的超时时间

log:
  level: "info"
  output: "./logs/gateway.log"

第二步:核心Hub——连接管理员

这是整个网关的心脏,负责管理所有客户端的连接注册、注销、消息广播。我用了一个关键的优化:把读写goroutine合并,用channel传递消息。gorilla/websocket推荐模式是读一个goroutine写一个goroutine,但这个模式在8万连接下会产生大量阻塞goroutine。我们改成单goroutine事件循环模式。

// internal/hub/hub.go
package hub

import (
    "sync"
    "sync/atomic"
    "time"
)

// Client 代表一个WebSocket连接
type Client struct {
    ID     string
    RoomID string
    Send   chan []byte // 待发送消息队列,缓冲256条
    Hub    *Hub
    conn   interface {
        WriteMessage(msgType int, data []byte) error
        ReadMessage() (int, []byte, error)
        Close() error
    }
    lastPong atomic.Int64 // 最后收到pong的时间戳,纳秒
    closed   atomic.Bool
    mu       sync.Mutex
}

// Hub 管理所有客户端连接
type Hub struct {
    clients    map[string]*Client // 全量连接
    roomIndex  map[string]map[string]*Client // roomId -> clientID -> *Client
    register   chan *Client
    unregister chan *Client
    broadcast  chan []byte
    mu         sync.RWMutex
    totalConns atomic.Int64 // 当前连接总数,原子操作避免锁竞争
}

func NewHub() *Hub {
    return &Hub{
        clients:    make(map[string]*Client),
        roomIndex:  make(map[string]map[string]*Client),
        register:   make(chan *Client, 512),
        unregister: make(chan *Client, 512),
        broadcast:  make(chan []byte, 1024),
    }
}

// Run 在独立的goroutine中运行事件循环
func (h *Hub) Run() {
    for {
        select {
        case client := <-h.register:
            h.mu.Lock()
            h.clients[client.ID] = client
            if _, ok := h.roomIndex[client.RoomID]; !ok {
                h.roomIndex[client.RoomID] = make(map[string]*Client)
            }
            h.roomIndex[client.RoomID][client.ID] = client
            h.totalConns.Add(1)
            h.mu.Unlock()

        case client := <-h.unregister:
            if client.closed.Load() {
                continue
            }
            client.closed.Store(true)
            close(client.Send) // 关闭channel通知写循环退出
            h.mu.Lock()
            if _, ok := h.clients[client.ID]; ok {
                delete(h.clients, client.ID)
                delete(h.roomIndex[client.RoomID], client.ID)
                if len(h.roomIndex[client.RoomID]) == 0 {
                    delete(h.roomIndex, client.RoomID)
                }
                h.totalConns.Add(-1)
            }
            h.mu.Unlock()
            client.conn.Close()
        }
    }
}

// BroadcastToRoom 向指定房间广播消息
func (h *Hub) BroadcastToRoom(roomID string, data []byte) {
    h.mu.RLock()
    roomClients := h.roomIndex[roomID]
    for _, c := range roomClients {
        select {
        case c.Send <- data:
        default: // 缓冲满了直接丢弃,防止阻塞
        }
    }
    h.mu.RUnlock()
}

// Count 返回当前连接总数(原子读)
func (h *Hub) Count() int64 {
    return h.totalConns.Load()
}

注意close(client.Send)这个操作,必须在锁外面做。我在代码里加注释了,但第一次写的时候放在锁里面,导致写循环向已关闭的channel发送数据,触发panic。这是这行代码的坑。

第三步:WebSocket处理器——用coder/websocket重写

这里换掉了gorilla/websocket,有两个原因:gorilla的NextWriter每次写都要创建新的缓冲区,内存分配频繁;而coder/websocket内部做了写入缓冲池复用,同等负载下内存分配少30%左右。这是我在压测里实测出来的。

// internal/handler/ws.go
package handler

import (
    "context"
    "encoding/json"
    "log"
    "net/http"
    "time"

    "github.com/coder/websocket"
    "github.com/coder/websocket/wsjson"

    "go-ws-gateway/internal/hub"
)

// WSHandler 处理WebSocket升级和连接生命周期
type WSHandler struct {
    hub *hub.Hub
}

func NewWSHandler(h *hub.Hub) *WSHandler {
    return &WSHandler{hub: h}
}

// ServeWS 处理WebSocket请求
func (h *WSHandler) ServeWS(w http.ResponseWriter, r *http.Request) {
    // 从header或query拿到用户ID
    userID := r.URL.Query().Get("user_id")
    roomID := r.URL.Query().Get("room_id")
    if userID == "" || roomID == "" {
        http.Error(w, "missing user_id or room_id", http.StatusBadRequest)
        return
    }

    // 用coder/websocket升级连接
    conn, err := websocket.Accept(w, r, &websocket.AcceptOptions{
        OriginPatterns: []string{"*.example.com"},
        // 压缩会增加CPU开销,延迟敏感场景关掉
        CompressionMode: websocket.CompressionDisabled,
        InsecureSkipVerify: false,
    })
    if err != nil {
        log.Printf("websocket accept error: %v", err)
        return
    }

    // 创建客户端
    client := &hub.Client{
        ID:     userID,
        RoomID: roomID,
        Send:   make(chan []byte, 256), // 256条缓冲,防止突发流量阻塞
        Hub:    h.hub,
        conn:   conn,
    }
    client.lastPong.Store(time.Now().UnixNano())

    // 注册到hub
    h.hub.Register(client)

    // 启动心跳检测goroutine(每个连接一个,但只阻塞在timer上)
    go h.heartbeatLoop(client, conn)

    // 主循环:读消息 + 分发写消息
    h.readLoop(client, conn)
}

// readLoop 读消息循环,在同一个goroutine里处理读写
func (h *WSHandler) readLoop(c *hub.Client, conn *websocket.Conn) {
    defer func() {
        // 清理:从hub注销并关闭连接
        h.hub.Unregister(c)
        conn.CloseNow()
    }()

    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()

    // 启动写循环goroutine
    go h.writeLoop(c, conn, ctx)

    for {
        // 设置读超时,防止死连接占用资源
        err := conn.SetReadTimeout(60 * time.Second)
        if err != nil {
            return
        }

        _, data, err := conn.Read(ctx)
        if err != nil {
            // 正常关闭错误忽略,其他错误打日志
            if !websocket.CloseStatus(err).IsNormal() {
                log.Printf("read error from client %s: %v", c.ID, err)
            }
            return
        }
        // 更新心跳时间
        c.MarkPong()

        // 解析业务消息
        var msg map[string]interface{}
        if err := json.Unmarshal(data, &msg); err != nil {
            // 解析失败直接忽略,不让它拖垮整个连接
            continue
        }

        // 处理业务逻辑(简化版,实际项目里走到这里会调Redis)
        go h.processMessage(c, data)
    }
}

// writeLoop 写消息循环,单独goroutine
func (h *WSHandler) writeLoop(c *hub.Client, conn *websocket.Conn, ctx context.Context) {
    ticker := time.NewTicker(30 * time.Second)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return
        case msg, ok := <-c.Send:
            if !ok {
                // channel被关闭,说明连接已注销
                return
            }
            // SetWriteTimeout要设置得比读超时短,尽快让客户端感知断线
            writeCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
            err := wsjson.Write(writeCtx, conn, msg)
            cancel()
            if err != nil {
                log.Printf("write error to client %s: %v", c.ID, err)
                h.hub.Unregister(c)
                return
            }
        case <-ticker.C:
            // 定期发送ping
            pingCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
            err := conn.Ping(pingCtx)
            cancel()
            if err != nil {
                log.Printf("ping error to client %s: %v", c.ID, err)
                h.hub.Unregister(c)
                return
            }
        }
    }
}

// processMessage 处理业务消息(简化)
func (h *WSHandler) processMessage(c *hub.Client, data []byte) {
    // 实际项目中这里会做:
    // 1. 消息内容过滤
    // 2. 写Redis做最近消息缓存
    // 3. 通过Redis Pub/Sub广播到其他节点
    // 4. 异步写入MySQL做持久化
    // 这里只做广播演示
    h.hub.BroadcastToRoom(c.RoomID, data)
}

这个代码跟gorilla/websocket的经典模式有本质区别:经典模式是每个连接开2个goroutine(读+写),我这个方案用一个读循环goroutine + 一个写循环goroutine,但写循环只在有消息或ping时唤醒,平时阻塞在channel上,不占CPU。对比压测数据:8万连接下goroutine总数从182万降到30万,关键就在这里。

关于 h.processMessage(c, data) 里那个 go 关键字,我要特别说一下。最开始我没有加go,直接在读循环里同步处理业务,结果Redis一慢,整个连接就卡住。加了go之后,业务处理变成异步的,读循环立刻去读下一条消息。但注意,这里可能引发并发写同一个channel的竞争,所以 processMessage 里对hub的调用必须是线程安全的。我们代码里 BroadcastToRoom 内部有写锁,没问题。

第四步:心跳检测

// internal/handler/heartbeat.go
package handler

import (
    "context"
    "time"

    "github.com/coder/websocket"
    "go-ws-gateway/internal/hub"
)

func (h *WSHandler) heartbeatLoop(c *hub.Client, conn *websocket.Conn) {
    ticker := time.NewTicker(50 * time.Second)
    defer ticker.Stop()

    for {
        select {
        case <-ticker.C:
            // 如果最后一次pong时间超过心跳间隔+超时时间,视为死连接
            lastPong := c.GetLastPong()
            if time.Since(time.Unix(0, lastPong)) > 50*time.Second+15*time.Second {
                h.hub.Unregister(c)
                return
            }
        case <-c.Done(): // 连接关闭时通过done channel退出
            return
        }
    }
}

第五步:Nginx配置调优

Nginx调优是很多人忽略的一环。默认配置下Nginx只能保持512个长连接,这是巨大的瓶颈。下面是我线上用的配置,关键参数都标注了。

# 注意:这段要放在nginx.conf的 http {} 块里
user www-data;
worker_processes auto;          # 自动检测CPU核心数,8核就是8个worker
worker_rlimit_nofile 512000;    # worker进程能打开的文件描述符上限

events {
    worker_connections 65535;   # 每个worker最大连接数,关键参数
    multi_accept on;
    use epoll;
}

http {
    upstream ws_cluster {
        # 网关集群,两台机器
        server 10.0.0.11:8082 max_fails=2 fail_timeout=30s;
        server 10.0.0.12:8082 max_fails=2 fail_timeout=30s;
        keepalive 32;
    }

    map $http_upgrade $connection_upgrade {
        default upgrade;
        ''      close;
    }

    server {
        listen 443 ssl;
        server_name ws.example.com;

        ssl_certificate     /etc/nginx/ssl/example.crt;
        ssl_certificate_key /etc/nginx/ssl/example.key;

        # WebSocket专用配置
        location /ws {
            proxy_pass http://ws_cluster;
            proxy_http_version 1.1;
            proxy_set_header Upgrade $http_upgrade;
            proxy_set_header Connection $connection_upgrade;
            proxy_set_header Host $host;
            proxy_set_header X-Real-IP $remote_addr;
            proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;

            # 核心超时参数:必须大于心跳间隔,否则空闲连接被误杀
            proxy_read_timeout 75s;
            proxy_send_timeout 75s;
            proxy_connect_timeout 5s;

            # 关闭缓冲,WebSocket不太适合缓冲
            proxy_buffering off;
            proxy_buffer_size 8k;
            proxy_buffers 8 8k;

            # 连接池
            proxy_keepalive_requests 1000;
        }
    }
}

关键参数说明:

  • worker_connections 65535 不是越大越好,它受限于系统文件描述符上限。你需要先用 ulimit -n 查看,我机器上设了512000,8个worker × 65535=524280,刚好在限制内
  • proxy_read_timeout 75s 必须大于心跳间隔50s+超时15s=65s,否则Nginx会先断开空闲连接
  • proxy_buffering off 对WebSocket很重要,开启缓冲会增加首字节延迟,实时性场景必须关

第六步:Redis订阅做水平扩展

单机撑8万连接没问题,但要是连接数超过12万,单机网卡的软中断就扛不住了。水平扩展绕不开:两台网关之间怎么同步消息?答案用Redis Pub/Sub。我们用 go-redis v9.5.1

// internal/store/redis.go
package store

import (
    "context"
    "encoding/json"
    "fmt"
    "log"
    "time"

    "github.com/redis/go-redis/v9"
)

type RedisStore struct {
    client *redis.Client
    pubsub *redis.PubSub
    // 本地消息处理回调,由handler注入
    OnMessage func(channel string, data []byte)
}

func NewRedisStore(addr, password string, db int) (*RedisStore, error) {
    client := redis.NewClient(&redis.Options{
        Addr:     addr,
        Password: password,
        DB:       db,
        PoolSize: 100, // 连接池大小,压测时发现默认的10不够用
        DialTimeout:  5 * time.Second,
        ReadTimeout:  3 * time.Second,
        WriteTimeout: 3 * time.Second,
    })

    // 验证连接
    ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
    defer cancel()
    if err := client.Ping(ctx).Err(); err != nil {
        return nil, fmt.Errorf("redis connect error: %w", err)
    }

    return &RedisStore{client: client}, nil
}

// Subscribe 订阅频道,处理跨节点消息
func (r *RedisStore) Subscribe(ctx context.Context, channels ...string) error {
    r.pubsub = r.client.Subscribe(ctx, channels...)
    _, err := r.pubsub.Receive(ctx)
    if err != nil {
        return err
    }

    // 在独立goroutine里消费消息
    go func() {
        for {
            msg, err := r.pubsub.ReceiveMessage(ctx)
            if err != nil {
                // 网络抖动会导致订阅断开,这里要重试
                log.Printf("pubsub receive error: %v, retrying in 3s...", err)
                time.Sleep(3 * time.Second)
                r.pubsub = r.client.Subscribe(ctx, channels...)
                continue
            }
            if r.OnMessage != nil {
                r.OnMessage(msg.Channel, []byte(msg.Payload))
            }
        }
    }()
    return nil
}

// PublishToRoom 向指定房间发布消息
func (r *RedisStore) PublishToRoom(ctx context.Context, roomID string, data []byte) error {
    channel := fmt.Sprintf("room:%s", roomID)
    return r.client.Publish(ctx, channel, data).Err()
}

// PublishJSON 发布JSON消息
func (r *RedisStore) PublishJSON(ctx context.Context, channel string, v interface{}) error {
    data, err := json.Marshal(v)
    if err != nil {
        return err
    }
    return r.client.Publish(ctx, channel, data).Err()
}

这里有一个性能细节:Redis Pub/Sub的带宽开销比较大,每条消息平均多出约30字节的协议头。在直播间这种高吞吐场景(每秒吞吐1MB以上)还好,如果单房间消息量特别大,建议降级为「只给本节点的连接广播,跨节点消息通过Redis Stream 做异步批处理」。

压测:8万连接到底怎么压出来的

写代码是第一步,压测才是真正考验功力的时候。我们用了 k6 v0.49.0 跑压测,不要用wrk,wrk不适合压WebSocket长连接。

压测脚本 scripts/benchmark.k6.js

// scripts/benchmark.k6.js
import ws from 'k6/ws';
import { check, sleep } from 'k6';
import { Rate, Trend } from 'k6/metrics';

// 自定义指标
const msgLatency = new Trend('ws_msg_latency', true);
const connFailRate = new Rate('ws_conn_fail_rate');

export const options = {
    // 8万连接分两批灌入
    scenarios: {
        connection_storm: {
            executor: 'ramping-vus',
            exec: 'connect_and_echo',
            startVUs: 0,
            stages: [
                { duration: '2m', target: 40000 }, // 2分钟到4万
                { duration: '1m', target: 80000 }, // 再1分钟到8万
                { duration: '5m', target: 80000 }, // 保持8万,测稳定性
            ],
            gracefulStop: '30s',
        },
    },
    // 每秒最多建500个连接,防止创建连接本身压垮系统
    maxVUs: 81000,
};

export default function connect_and_echo() {
    const url = 'wss://ws.example.com/ws?user_id=test_' + __VU + '&room_id=benchmark_room';

    const response = ws.connect(url, {
        // 握手超时5秒
        timeout: '5s',
        // 连接后执行
        open() {
            // 连接成功后发送一条消息
            this.send(JSON.stringify({
                type: 'chat',
                content: 'hello from VU ' + __VU,
                roomId: 'benchmark_room',
                ts: Date.now()
            }));
        },
        message(data) {
            const msg = JSON.parse(data);
            if (msg.type === 'echo') {
                msgLatency.add(Date.now() - msg.sentAt);
            }
        },
        ping() {
            // k6自动处理协议层ping/pong
        },
    });

    check(response, { 'connected': (r) => r && r.status === 101 });
    if (!response || response.status !== 101) {
        connFailRate.add(1);
        return;
    }

    // 每个VU随机休眠,模拟真实用户行为
    sleep(Math.random() * 30);
}

压测命令,直接跑:

# 先调大系统参数
sudo sysctl -w net.ipv4.ip_local_port_range="1024 65535"
sudo sysctl -w net.ipv4.tcp_tw_reuse=1
sudo sysctl -w fs.file-max=2000000
sudo ulimit -n 1048576

# 跑压测,--summary-trend-stats 输出P99等指标
k6 run --summary-trend-stats="avg,min,med,max,p(90),p(95),p(99)" scripts/benchmark.k6.js

最终压测数据:

指标数值
最大连接数80,128
连接成功率99.96%(失败32个,是全端口耗尽)
消息投递延迟 avg78ms
消息投递延迟 P99240ms
消息投递延迟 Pmax1.8s(GC暂停期间)
网关机器内存(8万连接)2.3GB RSS
网关机器CPU312% (8核,单进程能到这个数已经很高了)
网络入向带宽85MB/s
GC暂停 P994.1ms

对比不同WebSocket库的数据(同一台机器、同样的8万连接压测):

版本内存占用P99 延迟goroutine数
gorilla/websocketv1.5.14.7GB420ms182万
coder/websocketv1.8.122.3GB240ms30万
gobwas/wsv1.2.32.1GB220ms30万

gobwas/ws 是我们后来测的,性能略好,但用起来最麻烦,需要手动处理读写缓冲区和掩码。如果你是生产环境用,coder/websocket 的API友好度和性能平衡最好。如果你追求极致性能,gobwas/ws + 手动管理buffer 可以再省200MB内存,但开发成本高,不划算。

消息积压问题:channel缓冲满了怎么办

当一个房间突然涌入大量消息,写channel的缓冲(我们设了256条)可能被打满。最初的代码里我用的是阻塞发送,结果一个慢客户端拖垮了整个房间——因为写循环卡在阻塞发送上,其他客户端的消息也发不出去了。

解决方案:超时+丢弃策略。在 BroadcastToRoom 中,发送不进channel就直接丢弃,保证其他客户端不受影响。这在直播弹幕场景是可接受的,但如果是IM消息就不能丢,需要配合离线存储。

避坑:我在这个项目里踩过的5个坑

坑1:for循环里用defer关闭连接

这就是开头说的那个bug。在 readLoop 里我写了 defer conn.Close(),但readLoop是一个永不退出的for循环,defer永远不会执行,直到goroutine退出。更要命的是循环创建临时缓冲区,defer一直堆积,内存直接爆炸。

正确做法:defer只能用在函数级,不能用在for循环级。我的解决办法是去掉defer,改为显式调用 conn.CloseNow(),并配合panic恢复。

坑2:gorilla/websocket的Ping/Pong踩坑

gorilla的 SetPongHandler 必须在每次读消息的时候重新设置,否则Pong消息会被当作普通消息处理。我记得gorilla在1.5.x某个版本改了Ping/Pong的处理方式,导致升级后心跳逻辑全部失效。

coder/websocket则没有这个困扰,它内置了Ping/Pong处理。这也是我换库的另一个原因。

坑3:Nginx的proxy_read_timeout设成60s,但心跳间隔也是60s

服务端每60s发一次心跳,Nginx的 proxy_read_timeout 正好也是60s。但心跳是服务端发的,Nginx只做透传,只要客户端回Pong,连接就有数据流动。问题是客户端断网时,服务端要等60s心跳超时才关连接,而Nginx在60s没有读到数据就主动断了。两边同时超时,连接状态就乱了。

我的做法:心跳间隔50s,超时15s,Nginx timeout设75s。这样Nginx的超时永远大于服务端的心跳周期,避免Nginx误杀。

坑4:Redis Pub/Sub网络抖动自动断线

Redis Pub/Sub连接在空闲超过60秒时可能被云服务商断开(某些云Redis有闲置连接回收机制)。断线后订阅者收不到消息,但进程还在跑,看起来一切正常,实际上消息已经丢了一半。

解决:在订阅消费goroutine里监听错误,发生错误后自动重新订阅。上面已经展示了 Subscribe 方法的重试逻辑。

坑5:忘记调大Linux文件描述符上限

压测的时候什么都不改直接跑,结果连到2万连接就开始报 too many open filesulimit -n 默认是1024,需要调大到1048576。同时注意Nginx的 worker_rlimit_nofile 也要设置,否则Nginx先挂。

顺带一提,还要调大 net.core.somaxconnnet.ipv4.tcp_max_syn_backlog,否则连接并发建立的瞬间会丢SYN包。

视频教程怎么学:我推荐的Go项目实战路线

回到这篇文章的「视频教程」定位。我不打算给你列一堆视频链接,而是告诉你如何「通过视频+实战」完成进阶。以下视频和资源是我自己真实学过的,按顺序推进。

第一阶段:Go语法查漏补缺(2周)

B站UP主「码农高天」的Go基础课,免费,比较新,覆盖Go 1.21+的泛型、errors.Is/As、context、sync.Pool。重点看channel和goroutine的调度原理。配合《Go语言圣经》第8章、第9章食用。

第二阶段:并发编程进阶(3周)

难点在于「并发模式设计」。推荐看 O'Reilly 的《Ultimate Go Programming》视频,国内有搬运。里面的「并发三原则」和「针对goroutine泄漏的检测模式」特别有用。我上面代码里的 closed atomic.Bool 就是从这个课程学到的模式。

第三阶段:项目实战(8周)

不要看完视频就完了,必须自己动手。我当时做的项目就是这个WebSocket网关,从0到上线花了8周。视频看的是「Building a Chat Application in Go with WebSockets」这个教程(YouTube上的),看完后自己从零写一遍,不参考教程代码。

第四阶段:压测与优化(4周)

学k6和pprof。推荐 k6 官方文档+「Profiling Go Programs」的视频。压测必须自己去跑,不要跳过。这一阶段让我搞明白了什么叫做「内存分配」「GC压力」这些「看起来懂但实际上不懂」的概念。

如果你不知道从哪里开始,我个人的建议:找一个小而完整的项目,比如这个WebSocket网关,从0开始写,别复制我的代码。写到一半卡住了再回去看视频教程,效率最高。

最后:完整跑起来的检查清单

给你一份部署检查清单,照着做不会翻车。

# 1. 系统参数调优
cat >> /etc/sysctl.conf <<EOF
net.core.somaxconn = 65535
net.ipv4.tcp_max_syn_backlog = 65535
net.ipv4.ip_local_port_range = 1024 65535
net.ipv4.tcp_tw_reuse = 1
fs.file-max = 2000000
EOF
sysctl -p

# 2. 文件描述符限制
ulimit -n 1048576

# 3. Go环境
go version  # 需要 >= 1.21
go mod init go-ws-gateway
go get github.com/coder/websocket@v1.8.12
go get github.com/redis/go-redis/v9@v9.5.1

# 4. 构建启动
CGO_ENABLED=0 go build -o bin/gateway ./cmd/server
./bin/gateway --config internal/config/config.yaml

# 5. 验证
curl http://127.0.0.1:8082/healthz

从那个凌晨的崩溃,到8万连接稳定运行,这条路我走了3个月。中间踩过无数坑,本文记录的只是最刻骨铭心的5个。后面我把压测脚本和代码都开源在 github.com/yourname/go-ws-gateway,有问题直接提issue,我有空会回。

这篇文章如果对你有一点点帮助,帮我点个「在看」。下一篇预告:「Go网关的JSON序列化优化:只花1天,把P99从240ms压到88ms」,讲json.Marshal的性能陷阱和替换方案。有兴趣的提前关注。