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 } } } }
问题分析
从报错信息来看,根源是两点:
- WebSocket连接非并发安全:多个goroutine同时读写连接,导致内部状态混乱
- 未完整处理消息帧:客户端报错
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
相关产品推荐
相关产品推荐

