一个真实的生产事故:凌晨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.2GB | 4.7GB |
| GC STW 峰值 | 0.8ms | 32ms |
| 心跳超时断开率 | 0.2% | 23% |
| 消息投递延迟 P99 | 180ms | 4.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,但做三件事——
- 换掉 gorilla/websocket,改用
github.com/coder/websocket v1.8.12(原nhooyr.io/websocket),减少内存分配 - 读写goroutine合二为一,用channel做消息分发
- 读超时、写超时、心跳全部用context控制,杜绝泄漏
- Nginx调优 + 按地域分集群
最终效果,8万连接压测:
| 指标 | 优化前 | 优化后 | 提升幅度 |
|---|---|---|---|
| 内存占用(8万连接) | 4.7GB | 2.3GB | 51%降低 |
| GC暂停 | 32ms | 4.1ms | 87%降低 |
| 消息投递延迟 P99 | 4.2s | 240ms | 94%降低 |
| 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个,是全端口耗尽) |
| 消息投递延迟 avg | 78ms |
| 消息投递延迟 P99 | 240ms |
| 消息投递延迟 Pmax | 1.8s(GC暂停期间) |
| 网关机器内存(8万连接) | 2.3GB RSS |
| 网关机器CPU | 312% (8核,单进程能到这个数已经很高了) |
| 网络入向带宽 | 85MB/s |
| GC暂停 P99 | 4.1ms |
对比不同WebSocket库的数据(同一台机器、同样的8万连接压测):
| 库 | 版本 | 内存占用 | P99 延迟 | goroutine数 |
|---|---|---|---|---|
| gorilla/websocket | v1.5.1 | 4.7GB | 420ms | 182万 |
| coder/websocket | v1.8.12 | 2.3GB | 240ms | 30万 |
| gobwas/ws | v1.2.3 | 2.1GB | 220ms | 30万 |
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 files。ulimit -n 默认是1024,需要调大到1048576。同时注意Nginx的 worker_rlimit_nofile 也要设置,否则Nginx先挂。
顺带一提,还要调大 net.core.somaxconn 和 net.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的性能陷阱和替换方案。有兴趣的提前关注。