GoChat 的 Gateway 同时处理两类流量:注册、登录和会话管理走 HTTP/gRPC,实时消息走 WebSocket/Kafka。这两条路径的生命周期完全不同。HTTP 请求结束后就能释放状态,WebSocket 连接则可能存活几小时,并持续占用 goroutine、文件描述符和发送队列。
原草稿把这套网关写成了已经压测通过的“十万连接”方案,但仓库没有对应的脚本和原始报告。这里不再保留那些数字,而是对照当前代码说清已经做了什么,哪些地方还不能称为生产实现。仓库根模块声明 go 1.23.8,独立的 im-gateway 子模块仍声明 go 1.21;下文以实际实现为准,不把新版本能力当作既有依赖。
网关在消息链路中的位置
GoChat 将接入与业务处理拆开:
HTTP request
-> im-gateway
-> gRPC
-> im-logic
WebSocket message
-> im-gateway
-> Kafka upstream topic
-> im-logic
-> Kafka task/downstream topic
-> target im-gateway
-> WebSocket client
HTTP 路径通过 gRPC 调用 Logic 的 Auth、Conversation 和 Group 服务。WebSocket 上行消息并没有使用 gRPC stream,当前实现序列化 UpstreamMessage 后同步发送到 Kafka。
下行消息按 Gateway ID 区分 topic:
downstreamTopic := fmt.Sprintf(
"%s%s",
kafka.TopicDownstreamPrefix,
wsManager.GetGatewayID(),
)
每个 Gateway 实例只消费自己的下行 topic,再根据 TargetUserID 找本地连接。这套路由的前提是上游已经知道用户连在哪个 Gateway。仓库里有“在 Redis 注册用户在线状态”的 TODO,因此集群级用户路由还没闭环。
一个连接两个 pump
Connection 封装了 WebSocket、用户和一个有界发送队列:
type Connection struct {
Conn *websocket.Conn
UserID string
GatewayID string
Send chan []byte
Close chan struct{}
}
连接建立后,Gateway 启动 readPump 和 writePump:
go wm.readPump(connection)
go wm.writePump(connection)
readPump 是唯一业务读者,负责解析消息、构建上行事件并发往 Kafka。writePump 是唯一业务写者,它从 Send channel 取消息并写入 WebSocket,同时定期发 Ping。
这种设计有两个好处:
- 满足 Gorilla WebSocket 对“一个并发读者、一个并发写者”的约束。业务 goroutine 不会直接抢着写 socket。
- 通过有界 channel 显式定义每个连接最多积压多少待发消息。
当前队列容量是 256。它不是一个通用最佳值,只表示当前的内存与突发流量取舍。每条消息大小如果不受控,256 个 []byte 仍然可能留住大量内存。
反压:队列满后不能继续堆
sendMessage 使用非阻塞发送:
select {
case conn.Send <- data:
return nil
default:
close(conn.Close)
return errors.New("send queue full")
}
这个策略的意图很明确:当客户端长期读不动时,关闭连接,让客户端重连和补拉消息,避免单个慢客户端无界占用 Gateway 内存。
但当前关闭路径不幂等。readPump、writePump、队列满和 Stop 都可能关闭 conn.Close,重复 close 会 panic。更安全的连接结构应该把关闭收口统一起来:
type Connection struct {
Conn *websocket.Conn
UserID string
Send chan []byte
done chan struct{}
closeOnce sync.Once
}
func (c *Connection) Close() {
c.closeOnce.Do(func() {
close(c.done)
_ = c.Conn.Close()
})
}
Send 是否要关闭要另外设计。如果其他 goroutine 仍可能发送,贸然关闭 Send 会把问题从 double close 变成 send on closed channel。
连接表不只需要一把锁
Gateway 当前使用 map[string]*Connection + sync.RWMutex 保存 userID -> connection。对于这种需要“替换旧连接、关闭旧对象、发布新对象”的复合逻辑,普通 map 加锁比 sync.Map 更容易维护不变式。
当前 removeConnection(userID) 只按 UserID 删除。这有一个重连竞态:
- 旧连接 A 断开,但 readPump 还没走到 defer。
- 用户重连,新连接 B 写入 map。
- A 的 defer 调用
removeConnection(userID),把 B 删掉。
删除时应该同时比较连接身份:
func (wm *WebSocketManager) removeConnection(target *Connection) {
wm.mu.Lock()
current, ok := wm.connections[target.UserID]
if ok && current == target {
delete(wm.connections, target.UserID)
}
wm.mu.Unlock()
target.Close()
}
加锁只保证 map 操作不并发冲突,没有自动保证“删除的就是当初那条连接”。这是连接管理里比换成分片锁更优先的正确性问题。
心跳是读截止时间与 Ping/Pong 的配合
readPump 为连接设置 ReadDeadline,每次收到 Pong 都延后截止时间。writePump 定期发 Ping,并为每次写入设置 WriteDeadline。
conn.SetReadDeadline(time.Now().Add(pongTimeout))
conn.SetPongHandler(func(string) error {
return conn.SetReadDeadline(time.Now().Add(pongTimeout))
})
这里的时间参数需要满足 pongTimeout > pingInterval + 可接受抖动。当前默认值是 Ping 间隔 30 秒、Pong 超时 10 秒,初次 deadline 会在第一个 Ping 发出前过期。这两个默认值需要修正,例如 Ping 30 秒、Pong 等待 60 秒,再结合真实弱网数据调整。
传输层 TCP keepalive 可以识别某些断链,但应用层 Ping/Pong 更适合衡量客户端是否仍能处理 WebSocket 协议。两者可以同时存在。
Kafka 上行和下行的边界
上行路径当前调用:
wm.kafkaProducer.Send(context.Background(), topic, data)
这会丢失 WebSocket 请求的取消和超时上下文。更稳妥的方案是为每次发送创建有上限的 Context,并在消息 header 中传播 trace/request ID。
同步发送让单条 WebSocket 的读 pump 在 Kafka 确认前不再读下一条业务消息,好处是流程简单,坏处是 Kafka 延迟会直接传导到客户端读速度。是否换成异步 batch producer,需要同时设计:
- 队列满时怎么办。
- 客户端什么时候得到 ACK。
- 失败重试如何保持
ClientMsgID幂等。 - 服务停机时如何 flush,最长等多久。
下行路径如果找不到本地用户,当前 consumer 直接返错误。消费组会如何重试、是否进入死信队列,以及用户刚好重连到另一节点时如何重新路由,都要由消息投递协议给出答案。
鉴权与连接安全还没完成
当前 HandleWebSocket 只检查 URL 中是否有 token,然后把用户写死为 user123;CheckOrigin 也无条件返回 true。这部分不能直接用于公网。
上线前至少要补齐:
- 验证 JWT 签名、issuer、audience 和过期时间,UserID 只能来自已验证声明。
- 限制 Origin,或在非浏览器客户端中定义清晰的握手鉴权方式。
- 避免将长期凭证放在 URL query,URL 可能进入访问日志、代理和浏览器历史。
- 为握手、单 IP/用户连接数、消息大小和发送频率设置限制。
- 在认证完成前不把连接加入在线表。
优雅停机要等的是真实工作
Gateway 现有 Shutdown 已经列出注销服务、关闭 Kafka、关闭 WebSocket、关闭 gRPC 与协调器的顺序。但它还不是完整的优雅停机:
- Gin 直接调用
Engine.Run,没有保留http.Server调用Shutdown(ctx),无法先停止新 HTTP/WebSocket 握手。 WaitGroup等待的 WebSocketStart函数会立即返回,它没有跟踪每个连接的 read/write pump。- Kafka 生产者关闭前是否 flush,以及消费者是否已停止接收新任务,需要明确保证。
我会按下面的顺序改:
- 将 Gin Engine 放进
http.Server,收到 SIGTERM 后首先关闭 listener,停止新握手。 - 从服务发现和用户路由中摘除当前 Gateway,防止上游继续产生它的下行消息。
- 停止 Kafka consumer 获取新消息,等待正在推送的处理函数。
- 向活跃 WebSocket 发 Close frame,关闭连接并等待所有 pump 退出。
- flush 并关闭 producer,再关闭 gRPC 和协调器。
- 整个过程受一个总 Context deadline 限制,超时后强制退出。
压测前先定义指标
“支撑十万连接”至少需要同时说清:
- 长连接是否都已鉴权,心跳间隔和消息速率是多少。
- 单条消息大小、上行/下行比例和慢客户端比例。
- 网关、Kafka 和客户端压测器的机器配置与网络拓扑。
- 连接成功率、消息成功率、端到端 P50/P95/P99 延迟、重连率和队列满次数。
- RSS、heap、goroutine、FD、GC CPU 占比和停顿。
当前仓库没有这份报告,所以本文不给出虚构数字。等鉴权、路由、关闭竞态和优雅停机修好后,再建立可重复的压测基线,数字才有意义。
本次复盘后的实现顺序
- 修复 Connection 幂等关闭和旧连接误删新连接的竞态。
- 修正 Ping/Pong 时间关系,补弱网和慢消费者测试。
- 完成 JWT、Origin 和连接/消息限流。
- 完成
UserID -> GatewayID在线路由与重连替换协议。 - 用
http.Server和真实连接 WaitGroup 重做停机流程。 - 建立压测脚本、Dashboard 和每次发版的回归基线。
这个顺序先修正确性和安全边界,再谈分片锁、sync.Pool 或 Netpoller 框架。如果连接还会被重复关闭,或者用户路由尚未闭环,更高的 benchmark 数字不能让它变成可用的 IM 网关。