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

多Go协程共享net.Conn时部分协程读阻塞问题排查

问题

我正在实现一套客户端与服务端RPC应用,基于TCP传输JSON数据。单客户端请求时,请求响应周期可正常完成,但当10个Go协程共享一个net.Conn对象时,第一个协程退出后,其余协程会在读取连接时阻塞(据观察,服务端仍会返回处于处理中的RPC调用响应),这些协程无法执行waitgroup的done调用,导致整个进程挂起。

main.go中启动协程的代码:

rpcClient := connectAndCreateRpcClient() // has net.Conn from net.Dial

var wg sync.WaitGroup
wg.Add(10)

for adders := 0; adders < 10; adders++ {
    go func() {
        sum := 0
        for i := 0; i < 100; i++ {
            sum += rpcClient.Add(1, 2)
        }
        wg.Done()
    }()
}
wg.Wait()

客户端RPC调用代码片段:

if err != nil {
    return nil, err
}
// .. marshal into jr, append new line as break

jr = append(jr, byte('\n'))
_, err = jrc.conn.Write(jr)

if err != nil {
    return nil, err
}

cr := bufio.NewReader(jrc.conn)

buff, err := cr.ReadBytes(byte('\n'))
if err != nil {
    return nil, err
}

// .. unmarshall and return the result

服务端接收并响应请求的代码片段:

cr := bufio.NewReader(jrst.conn)

for {
    rj, err := cr.ReadBytes(byte('\n'))

    if err != nil {
        if err == io.EOF {
            break
        }
        panic(err)
    }

    // ... un marshall the request ...

    res, err := cbs.Unary(req.Method, req.Params) 
    if err != nil {
        log.Fatalf("Unary return: %s\n", err.Error())
    }
    // .. marshall up the response

    br = append(br, byte('\n'))

    _, err = jrst.conn.Write(br)
    if err != nil {
        return err
    }
}

核心现象是:当main中的10个协程之一完成100次调用后,其余9个协程会阻塞在buff, err := cr.ReadBytes(byte('\n'))处。

请问为何一个协程完成会导致其他协程无法读取?有没有更优的读写实现方式来保持连接正常?


分析与解决方案

问题根源

  • 多bufio.Reader的冲突:每个协程调用Add时都新建bufio.NewReader(jrc.conn),但bufio.Reader会缓存连接中的数据。多个协程的bufio.Reader同时读取同一个net.Conn,会出现一个Reader把其他协程对应的响应数据读走的情况,导致被抢数据的协程阻塞等待。
  • TCP流式传输的乱序问题:TCP是无边界的流式协议,并发发送的请求对应的响应可能乱序返回。加上没有请求ID绑定,客户端无法区分响应属于哪个请求,即使数据没被抢,也可能出现协程拿到不属于自己的响应,后续正确响应到来时没有协程去读,最终导致阻塞。
  • net.Conn的Read非并发安全:net.Conn的Write方法是并发安全的,但Read不是。多个协程同时调用Read会导致数据被随机分配给任意协程,完全无法保证请求和响应的一一对应。

优化实现方案

客户端优化要点

  1. 共享单个bufio.Reader:在RPC客户端初始化时创建唯一的bufio.Reader,所有协程共用这个Reader,避免多个Reader争抢数据。
  2. 添加请求ID绑定响应:给每个请求生成唯一ID,服务端返回响应时携带该ID,客户端读取响应后根据ID匹配到对应的请求,确保响应交付给正确的协程。
  3. 互斥锁保护读写流程:对单个连接的"写请求-读响应"流程加互斥锁,同一时间只允许一个协程执行完整的请求响应周期,彻底避免并发读写的冲突。

优化后的客户端示例代码:

import (
    "bufio"
    "encoding/json"
    "fmt"
    "net"
    "sync"
    "github.com/google/uuid" // 可用其他方式生成唯一ID
)

type JsonRpcRequest struct {
    ID     string      `json:"id"`
    Method string      `json:"method"`
    Params []int       `json:"params"`
}

type JsonRpcResponse struct {
    ID     string      `json:"id"`
    Result interface{} `json:"result"`
    Error  string      `json:"error,omitempty"`
}

type JsonRpcClient struct {
    conn   net.Conn
    reader *bufio.Reader
    mu     sync.Mutex
}

func NewJsonRpcClient(conn net.Conn) *JsonRpcClient {
    return &JsonRpcClient{
        conn:   conn,
        reader: bufio.NewReader(conn),
    }
}

func (jrc *JsonRpcClient) Add(a, b int) (int, error) {
    jrc.mu.Lock()
    defer jrc.mu.Unlock()

    // 构造带唯一ID的请求
    reqID := uuid.New().String()
    req := JsonRpcRequest{
        ID:     reqID,
        Method: "Add",
        Params: []int{a, b},
    }

    // 序列化请求并添加换行符
    reqBytes, err := json.Marshal(req)
    if err != nil {
        return 0, err
    }
    reqBytes = append(reqBytes, '\n')

    // 发送请求
    _, err = jrc.conn.Write(reqBytes)
    if err != nil {
        return 0, err
    }

    // 读取响应
    respBytes, err := jrc.reader.ReadBytes('\n')
    if err != nil {
        return 0, err
    }

    // 反序列化响应
    var resp JsonRpcResponse
    err = json.Unmarshal(respBytes, &resp)
    if err != nil {
        return 0, err
    }

    // 验证响应ID匹配
    if resp.ID != reqID {
        return 0, fmt.Errorf("response ID mismatch: expected %s, got %s", reqID, resp.ID)
    }

    // 解析结果
    result, ok := resp.Result.(float64)
    if !ok {
        return 0, fmt.Errorf("invalid result type")
    }
    return int(result), nil
}

服务端适配优化

服务端需要在响应中回传请求的ID,确保客户端能匹配:

// 服务端处理请求时的逻辑片段
var req JsonRpcRequest
err = json.Unmarshal(rj, &req)
if err != nil {
    // 构造错误响应
    resp := JsonRpcResponse{
        ID:    req.ID,
        Error: err.Error(),
    }
    respBytes, _ := json.Marshal(resp)
    respBytes = append(respBytes, '\n')
    jrst.conn.Write(respBytes)
    continue
}

// 执行方法调用
res, err := cbs.Unary(req.Method, req.Params)
if err != nil {
    // 构造错误响应
    resp := JsonRpcResponse{
        ID:    req.ID,
        Error: err.Error(),
    }
    respBytes, _ := json.Marshal(resp)
    respBytes = append(respBytes, '\n')
    jrst.conn.Write(respBytes)
    continue
}

// 构造成功响应
resp := JsonRpcResponse{
    ID:     req.ID,
    Result: res,
}
respBytes, err := json.Marshal(resp)
if err != nil {
    // 错误处理
}
respBytes = append(respBytes, '\n')
_, err = jrst.conn.Write(respBytes)
if err != nil {
    return err
}

可选优化方向

如果需要更高的并发性能,可以考虑:

  • 连接池:维护一组TCP连接,协程从池中获取连接使用,用完归还,平衡并发能力和资源消耗。
  • 异步请求响应:客户端发送请求后不立即阻塞等待,而是通过通道接收响应,配合请求ID实现异步处理,但复杂度会提升。

内容的提问来源于stack exchange,提问作者user8801733

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 01:59:57