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

如何解决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):
    }
}

问题原因

  1. write: broken pipe:首次请求结束后客户端进程终止,WebSocket连接未彻底清理,服务端对应的channel未销毁,第二次请求复用旧channel时,尝试向已关闭的连接写入数据。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 20:57:03