如何修复Gorilla WebSocket中缓冲写入触发SIGTERM的问题
Go多实例高可用WebSocket客户端实现方案
现有代码核心问题排查
- 空指针panic根因:初始化客户端时未完成连接建立就返回实例,且多goroutine操作
wsconn未加有效锁保护,导致wsconn偶发为nil,触发空指针引用。此外关闭逻辑缺失,可能出现往已关闭的sendBuf写入数据的panic。 - SIGTERM信号触发根因:无统一的资源释放时序,连接关闭时存在未处理的网络读写操作、未退出的后台goroutine,触发运行时panic被系统发送SIGTERM信号。使用全局变量时相当于所有实例共享发送队列与锁,规避了单实例chan空/关闭的问题,但全局资源释放时所有实例都会受影响,关闭时依然会崩溃。
- 多实例隔离失效:全局变量的
mutex和sendBuf破坏了实例隔离性,多个客户端实例的消息会串流,无法独立控制生命周期。 - "use of closed network connection"报错根因:连接关闭后没有停止读写goroutine,依然在操作已经关闭的网络套接字。
优化实现方案
1. 结构体设计优化
新增等待组、关闭原子标记保证资源生命周期可控,所有资源均为实例内部持有,天然支持多实例隔离:
type WebSocketClient struct { url string sendBuf chan []byte ctx context.Context ctxCancel context.CancelFunc mu sync.RWMutex wsconn *websocket.Conn wg sync.WaitGroup // 等待所有后台goroutine退出 closed atomic.Bool // 原子标记客户端是否已关闭,避免重复关闭 }
2. 初始化逻辑修正
首次连接建立成功再返回实例,避免返回不可用的空实例:
func NewWebSocketClient(url string, messageHandler WsHandler, errorHandler WsErrHandler) (*WebSocketClient, error) { client := &WebSocketClient{ url: url, sendBuf: make(chan []byte, 10), // 适当加大队列容量避免写入超时 } client.ctx, client.ctxCancel = context.WithCancel(context.Background()) // 首次建立连接,失败直接返回错误 conn, _, err := websocket.DefaultDialer.Dial(url, nil) if err != nil { return nil, err } client.wsconn = conn client.closed.Store(false) // 注册后台goroutine到等待组 client.wg.Add(3) go client.listen(messageHandler, errorHandler) go client.listenWrite() go client.ping() return client, nil }
3. 写入逻辑优化
提前判断实例状态,避免往已关闭的队列写入数据,移除无意义的循环重试:
func (conn *WebSocketClient) Write(data []byte) error { if conn.closed.Load() { return fmt.Errorf("websocket client already closed") } ctx, cancel := context.WithTimeout(conn.ctx, 150*time.Millisecond) defer cancel() select { case conn.sendBuf <- data: return nil case <-ctx.Done(): return ctx.Err() } }
4. 写监听逻辑修正
统一加锁读取连接实例,写入失败触发异步重连:
func (conn *WebSocketClient) listenWrite() { defer conn.wg.Done() for { select { case data, ok := <-conn.sendBuf: if !ok { // 发送队列已关闭,直接退出 return } conn.mu.RLock() ws := conn.wsconn conn.mu.RUnlock() if ws == nil { log.Println("websocket connection not available, drop message") continue } if err := ws.WriteMessage(websocket.TextMessage, data); err != nil { log.Println("write message failed:", err) conn.reconnect() // 写入失败触发重连 } case <-conn.ctx.Done(): return } } }
5. 重连逻辑实现
单独封装重连逻辑,加锁避免并发重连,支持重试机制:
func (conn *WebSocketClient) reconnect() { conn.mu.Lock() defer conn.mu.Unlock() if conn.closed.Load() { return } // 先关闭旧连接 if conn.wsconn != nil { _ = conn.wsconn.Close() } // 最多重试3次,指数退避 var err error var newConn *websocket.Conn for i := 0; i < 3; i++ { select { case <-conn.ctx.Done(): return default: newConn, _, err = websocket.DefaultDialer.Dial(conn.url, nil) if err == nil { conn.wsconn = newConn log.Println("websocket reconnect success") return } time.Sleep(time.Second * time.Duration(i+1)) } } log.Println("websocket reconnect failed after 3 retries") }
6. 优雅关闭逻辑实现
严格按照停止新写入→关闭网络连接→等待goroutine退出的时序释放资源,避免资源泄漏和关闭panic:
func (conn *WebSocketClient) Close() error { if !conn.closed.CompareAndSwap(false, true) { return nil // 避免重复关闭 } // 1. 取消上下文,通知所有后台goroutine退出 conn.ctxCancel() // 2. 关闭发送队列,禁止新的写入请求 close(conn.sendBuf) // 3. 主动关闭WebSocket连接 conn.mu.Lock() if conn.wsconn != nil { // 发送关闭帧,优雅断开服务端连接 _ = conn.wsconn.WriteMessage(websocket.CloseMessage, websocket.FormatCloseMessage(websocket.CloseNormalClosure, "")) _ = conn.wsconn.Close() conn.wsconn = nil } conn.mu.Unlock() // 4. 等待所有后台goroutine退出,最多等待3秒超时 waitCtx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() done := make(chan struct{}) go func() { conn.wg.Wait() close(done) }() select { case <-done: return nil case <-waitCtx.Done(): return fmt.Errorf("close timeout, some goroutine not exit") } }
补充注意事项
- 读监听的
listen方法中,读取到网络关闭错误时也需要触发reconnect逻辑,不要直接退出goroutine - ping方法中需要加锁操作
wsconn,心跳失败时同样触发重连 - 所有操作
wsconn的位置都需要加对应的读写锁,避免竞态问题
内容的提问来源于stack exchange,提问作者Marvin.Hansen
相关产品推荐
相关产品推荐

