如何解决WebSocket的write: broken pipe及write: close sent错误?
WebSocket进度条通信中的管道错误排查与解决
问题场景
通过WebSocket实现进程间通信,向CLI客户端发送下载进度条数据,遇到以下问题:
- 首次请求无报错,第二次请求触发
write: broken pipe错误,但进度条仍能正常完成100%更新 - 处理客户端关闭消息后,
broken pipe错误消失,但出现write: close sent错误
服务端(Go Fiber WebSocket)代码
// middleware app.Use("/ws", func(c *fiber.Ctx) error { if websocket.IsWebSocketUpgrade(c) { c.Locals("allowed", true) return c.Next() } return fiber.ErrUpgradeRequired }) // handler func (s *downloaderService) progressBar(c *websocket.Conn) { channel := api.CreateChannel(c.Params("client")) done := make(chan bool) go func() { for { t, _, err := c.ReadMessage() if err != nil { log.Println("Error reading message:", err) return } if t == websocket.CloseMessage { done <- true } } }() for { select { case <-done: return case data, ok := <-channel.Subscribe(): if !ok { return } progressBar := data.(downloader.Progressbar) payload, err := json.Marshal(progressBar) if err != nil { log.Println("Error marshalling data:", err) break } c.SetWriteDeadline(time.Now().Add(10 * time.Second)) if err := c.WriteMessage(websocket.TextMessage, payload); err != nil { log.Println("Error sending progress data:", err) done <- true } } } }
客户端(Gorilla WebSocket)代码
func main() { interrupt := make(chan os.Signal, 1) signal.Notify(interrupt, []os.Signal{syscall.SIGINT, syscall.SIGKILL, syscall.SIGTERM, syscall.SIGSTOP, os.Interrupt}...) ctx, cancel := context.WithCancel(context.Background()) conn, res, err := websocket.DefaultDialer.DialContext(ctx, ws, nil) if err != nil { log.Fatalf("Error dialing websocket: %v. Status code %d", err, res.StatusCode) return } executeCommand(ctx) progressBar := progressbar() go func() { for { _, message, err := conn.ReadMessage() if err != nil { log.Println("Error reading message:", err) break } var progress progress if err := json.Unmarshal(message, &progress); err != nil { log.Println("Error unmarshalling message:", err) break } if progress.Done { truncateStore() cancel() break } progressBar.update(progress.Index, progress.Downloaded, progress.Size) } }() for { select { case <-ctx.Done(): closeConn(ctx, conn) return case <-interrupt: stopDownload() closeConn(ctx, conn) return } } } func closeConn(ctx context.Context, conn *websocket.Conn) { if err := conn.WriteMessage(websocket.CloseMessage, websocket.FormatCloseMessage(websocket.CloseNormalClosure, "")); err != nil { log.Println("Error sending close signal to server:", err) return } select { case <-ctx.Done(): case <-time.After(time.Second): } }
问题原因
write: broken pipe:首次请求结束后客户端进程终止,WebSocket连接未彻底清理,服务端对应的channel未销毁,第二次请求复用旧channel时,尝试向已关闭的连接写入数据。write: close sent:服务端收到客户端的关闭消息后,未立即终止写循环,仍尝试向已发送关闭帧的连接发送数据。
解决方案
服务端代码修改
func (s *downloaderService) progressBar(c *websocket.Conn) { clientID := c.Params("client") channel := api.CreateChannel(clientID) defer api.DestroyChannel(clientID) // 连接关闭时自动销毁channel done := make(chan struct{}) // 用struct{}替代bool,节省资源 go func() { defer close(done) // 协程退出时关闭done通道,触发主循环退出 for { t, _, err := c.ReadMessage() if err != nil { log.Println("Error reading message:", err) return } if t == websocket.CloseMessage { // 回复客户端关闭确认帧 c.WriteMessage(websocket.CloseMessage, websocket.FormatCloseMessage(websocket.CloseNormalClosure, "")) return } } }() for { select { case <-done: return case data, ok := <-channel.Subscribe(): if !ok { return } progressBar := data.(downloader.Progressbar) payload, err := json.Marshal(progressBar) if err != nil { log.Println("Error marshalling data:", err) continue // 序列化失败时跳过当前数据,不终止循环 } c.SetWriteDeadline(time.Now().Add(10 * time.Second)) if err := c.WriteMessage(websocket.TextMessage, payload); err != nil { log.Println("Error sending progress data:", err) return // 写失败直接退出,不再尝试后续操作 } // 进度完成后主动退出循环 if progressBar.Done { return } } } }
客户端代码修改
func closeConn(ctx context.Context, conn *websocket.Conn) { // 设置写超时,避免阻塞 conn.SetWriteDeadline(time.Now().Add(2 * time.Second)) if err := conn.WriteMessage(websocket.CloseMessage, websocket.FormatCloseMessage(websocket.CloseNormalClosure, "progress completed")); err != nil { log.Println("Error sending close signal:", err) } // 等待服务端响应或超时 select { case <-ctx.Done(): case <-time.After(2 * time.Second): } // 主动关闭连接,确保资源释放 conn.Close() }
关键修改点
- 服务端添加
defer api.DestroyChannel(clientID),确保连接关闭后销毁对应channel,避免后续请求复用无效资源 - 优化
done通道类型为chan struct{},并在read协程退出时关闭通道,确保主循环及时终止 - 收到客户端关闭消息时,服务端回复关闭确认帧,规范WebSocket关闭流程
- 写数据失败时直接退出循环,不再尝试向无效连接写入
- 客户端关闭连接时主动调用
conn.Close(),确保连接资源彻底释放
内容的提问来源于stack exchange,提问作者newtocoding
相关产品推荐
相关产品推荐

