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

nhooyr WebSocket读写并发异常问题排查求助

WebSocket同时读写时连接异常问题

核心问题:Go-Gin的WebSocket处理器无法同时完成读写操作。写入前设置5秒休眠时,休眠期间读取正常;5秒后尝试写入,读取操作触发报错,错误信息包括:

  • failed to get reader (failed to read frame header: EOF)
  • received close frame: status = StatusPolicyViolation

服务端代码

type DriverWebSocketOutput struct {
    IsBase bool            `json:"is_base"`
    Orders []*models.Order `json:"orders"`
}

func (config *Config) DriverWebSocket(c *gin.Context) {
    ctx := c.Request.Context()
    driverID := auth.GetToken(c).ID()
    driverLoc := &redisclient.DriverLocation{ID: driverID}

    // 将HTTP连接升级为WebSocket
    wsconn, err := websocket.Accept(c.Writer, c.Request, &websocket.AcceptOptions{InsecureSkipVerify: true})
    if err != nil {
        return
    }
    defer wsconn.Close(websocket.StatusInternalError, "")

    // 读取司机位置
    go func() {
        for {
            if err := wsjson.Read(ctx, wsconn, &driverLoc.Loc); err != nil {
                if websocket.CloseStatus(err) == websocket.StatusNormalClosure || websocket.CloseStatus(err) == websocket.StatusGoingAway {
                    return
                }
                fmt.Println("读取司机位置出错", err)//5秒后会出现以下任一错误:
                //failed to get reader: received close frame: status = StatusPolicyViolation and reason = "unexpected data message"
                //failed to read JSON message: failed to get reader: failed to read frame header: EOF
                break
            }
            fmt.Println("司机位置", driverLoc.Loc)//会打印5秒数据!
        }
    }()

    // 向司机发送订单更新
    for {
        out := DriverWebSocketOutput{}//
        time.Sleep(time.Second * 5)
        if err := wsjson.Write(ctx, wsconn, &out); err != nil {
            fmt.Println("向WebSocket写入出错: ", err)
            break
        }

    }
}

客户端代码

go func() {
    received := handlers.DriverWebSocketOutput{}
    for {
        if err = wsjson.Read(c.WSContext, wsconn, &received); err != nil {
            fmt.Println(fmt.Sprintf("无法读取WebSocket消息: %v", err))
            //failed to read JSON message: failed to get reader: previous message not read to completion
            break
        } 
    }
}()

// 向服务器发送司机位置
ticker := time.NewTicker(time.Millisecond * 1000)
defer ticker.Stop()
loop:
for {
    select {
    //发送坐标到服务器
    case <-closeRead.Done():
        fmt.Println("服务器读取完成")
        break loop
    case t := <-ticker.C:
                coords := redisclient.Location{Lat: c.Routing.PolylineCoords[c.Routing.I][0], Lng: c.Routing.PolylineCoords[c.Routing.I][1]}
                if err = wsjson.Write(c.WSContext, wsconn, &coords); err != nil {
                    fmt.Printf("无法通过WebSocket发送位置: %v\n", err)
                    //failed to write JSON message: failed to get writer: WebSocket closed: failed to get reader: previous message not read to completion
                    break loop
                }
          }
    }
}

问题分析

从报错信息来看,根源是两点:

  1. WebSocket连接非并发安全:多个goroutine同时读写连接,导致内部状态混乱
  2. 未完整处理消息帧:客户端报错previous message not read to completion,说明上一次读取操作未完成就发起新请求;服务端的StatusPolicyViolation是因为连接处于未处理完上一条消息的状态,触发协议冲突

解决方案

1. 用互斥锁保护连接读写

WebSocket连接不能被多个goroutine同时操作,必须用sync.Mutex加锁,确保同一时间只有一个操作在执行。

2. 确保消息处理完整性

遇到读取错误时,主动调用NextReader()跳过未处理的消息帧,避免连接状态卡住。

3. 优化循环逻辑

加入上下文取消判断,避免goroutine泄漏;连接关闭时统一处理。

修改后的代码示例

服务端代码

import "sync"

type DriverWebSocketOutput struct {
    IsBase bool            `json:"is_base"`
    Orders []*models.Order `json:"orders"`
}

func (config *Config) DriverWebSocket(c *gin.Context) {
    ctx := c.Request.Context()
    driverID := auth.GetToken(c).ID()
    driverLoc := &redisclient.DriverLocation{ID: driverID}

    // 将HTTP连接升级为WebSocket
    wsconn, err := websocket.Accept(c.Writer, c.Request, &websocket.AcceptOptions{InsecureSkipVerify: true})
    if err != nil {
        return
    }
    defer wsconn.Close(websocket.StatusNormalClosure, "")

    var mu sync.Mutex

    // 读取司机位置
    go func() {
        defer func() {
            mu.Lock()
            wsconn.Close(websocket.StatusNormalClosure, "")
            mu.Unlock()
        }()
        for {
            select {
            case <-ctx.Done():
                return
            default:
                mu.Lock()
                err := wsjson.Read(ctx, wsconn, &driverLoc.Loc)
                mu.Unlock()
                if err != nil {
                    if websocket.CloseStatus(err) == websocket.StatusNormalClosure || websocket.CloseStatus(err) == websocket.StatusGoingAway {
                        return
                    }
                    // 清理未完成的消息帧
                    mu.Lock()
                    if _, rerr := wsconn.NextReader(); rerr != nil {
                        fmt.Println("清理未完成消息出错", rerr)
                    }
                    mu.Unlock()
                    fmt.Println("读取司机位置出错", err)
                    return
                }
                fmt.Println("司机位置", driverLoc.Loc)
            }
        }
    }()

    // 向司机发送订单更新
    for {
        select {
        case <-ctx.Done():
            return
        default:
            time.Sleep(time.Second * 5)
            out := DriverWebSocketOutput{}
            mu.Lock()
            err := wsjson.Write(ctx, wsconn, &out)
            mu.Unlock()
            if err != nil {
                fmt.Println("向WebSocket写入出错: ", err)
                return
            }
        }
    }
}

客户端代码

import "sync"

// ... 其他业务代码

go func() {
    received := handlers.DriverWebSocketOutput{}
    var mu sync.Mutex
    defer func() {
        mu.Lock()
        wsconn.Close(websocket.StatusNormalClosure, "")
        mu.Unlock()
    }()
    for {
        select {
        case <-c.WSContext.Done():
            return
        default:
            mu.Lock()
            err = wsjson.Read(c.WSContext, wsconn, &received)
            mu.Unlock()
            if err != nil {
                if websocket.CloseStatus(err) == websocket.StatusNormalClosure || websocket.CloseStatus(err) == websocket.StatusGoingAway {
                    return
                }
                // 清理未完成的消息帧
                mu.Lock()
                if _, rerr := wsconn.NextReader(); rerr != nil {
                    fmt.Println("清理未完成消息出错", rerr)
                }
                mu.Unlock()
                fmt.Printf("无法读取WebSocket消息: %v\n", err)
                return
            }
        }
    }
}()

// 向服务器发送司机位置
ticker := time.NewTicker(time.Millisecond * 1000)
defer ticker.Stop()
var mu sync.Mutex
loop:
for {
    select {
    case <-closeRead.Done():
        fmt.Println("服务器读取完成")
        break loop
    case <-ticker.C:
        coords := redisclient.Location{Lat: c.Routing.PolylineCoords[c.Routing.I][0], Lng: c.Routing.PolylineCoords[c.Routing.I][1]}
        mu.Lock()
        err = wsjson.Write(c.WSContext, wsconn, &coords)
        mu.Unlock()
        if err != nil {
            fmt.Printf("无法通过WebSocket发送位置: %v\n", err)
            break loop
        }
    }
}

内容的提问来源于stack exchange,提问作者Sami Al-Subhi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 14:48:18