You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何修复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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.06 08:18:03