Skip to content

Client端是如何实现高并发的?

一句话概括

miniRPC Client 通过全双工通信 + 异步收发分离 + pending map 请求匹配实现高并发:多个 goroutine 可以同时发起调用(sending 锁串行写入),单个 receive 协程持续读取响应并通过 Seq 序列号匹配回对应的调用,整个过程基于一条 TCP 连接完成多路复用。


客户端并发模型全景

                              一条 TCP 连接(全双工)
                         ┌──────────────────────────────┐
goroutine 1 ── send ──┐  │                              │
goroutine 2 ── send ──┤  │  ======► 请求方向 ======►     │  ──► 服务端
goroutine 3 ── send ──┘  │                              │
    sending 锁串行写入     │  ◄====== 响应方向 ◄======    │  ◄── 服务端
                         └──────────────────────────────┘

                               receive goroutine(唯一)

                          ┌────────────┼────────────┐
                          ▼            ▼            ▼
                    pending[1]   pending[2]   pending[3]
                    call1.Done   call2.Done   call3.Done
                        │            │            │
                        ▼            ▼            ▼
                   goroutine 1  goroutine 2  goroutine 3
                   收到结果      收到结果      收到结果

核心机制一:异步收发分离

Client 的并发能力建立在发送和接收完全解耦的基础上。

传统同步 RPC 的问题

如果一个客户端采用"发一个请求、等一个响应"的同步模型:

goroutine 1:send(req1) → 等待 resp1 → 完成        (200ms)
goroutine 2:              send(req2) → 等待 resp2   (又 200ms)
goroutine 3:                            send(req3)   (又 200ms)
总耗时:600ms

每个调用都独占连接,后续调用必须排队。

miniRPC 的异步模型

miniRPC 将发送和接收分离到不同的执行流中:

goroutine 1:send(req1) → 注册到 pending → 返回(不阻塞)
goroutine 2:send(req2) → 注册到 pending → 返回
goroutine 3:send(req3) → 注册到 pending → 返回

receive goroutine(后台持续运行):
  读到 resp2 → pending[2].Done ← call2   ← 服务端先处理完 req2
  读到 resp1 → pending[1].Done ← call1
  读到 resp3 → pending[3].Done ← call3

总耗时:≈ 200ms(三个请求几乎同时被服务端处理)

发送方只负责"把请求写到连接上然后走人",不等待响应。接收方独立运行,持续从连接读取响应并分发。这就是一条连接上的多路复用


核心机制二:pending map —— 请求与响应的匹配

并发场景下,多个请求同时在途,响应到达的顺序不确定。pending map 解决了"这个响应属于哪个请求"的问题。

工作流程

步骤 1:发送请求时注册

  registerCall(call):
    client.mu.Lock()
    call.Seq = client.seq        // 分配序列号 seq=1
    client.pending[1] = call     // 注册到 pending map
    client.seq++                 // seq 递增为 2
    client.mu.Unlock()

  pending map 状态:{ 1: call1, 2: call2, 3: call3 }

步骤 2:收到响应时匹配

  removeCall(h.Seq):
    client.mu.Lock()
    call := client.pending[2]    // 通过 Seq 找到对应的 Call
    delete(client.pending, 2)    // 从 map 中移除
    client.mu.Unlock()
    return call                  // 返回 call2

  pending map 状态:{ 1: call1, 3: call3 }

步骤 3:通知调用方

  call2.Done <- call2            // 通过 channel 通知 goroutine 2

为什么用 map 而不是队列

因为响应顺序不等于请求顺序。服务端对每个请求独立并行处理,处理快的先返回。如果用队列(FIFO),第一个响应到达时不一定是第一个请求的结果,就无法正确匹配。用 map 以 Seq 为 key 做 O(1) 查找,不依赖顺序。


核心机制三:两把锁的分工

Client 有两把 sync.Mutex,各自保护不同的资源:

go
type Client struct {
    sending sync.Mutex     // 锁 1:保护发送操作
    mu      sync.Mutex     // 锁 2:保护 pending map 和状态字段
}

sending 锁 —— 保护写入连接

go
func (client *Client) send(call *Call) {
    client.sending.Lock()          // 获取发送锁
    defer client.sending.Unlock()

    seq, err := client.registerCall(call)   // 内部用 mu 锁操作 pending
    // ...
    client.header.ServiceMethod = call.ServiceMethod
    client.header.Seq = seq
    client.header.Error = ""

    client.cc.Write(&client.header, call.Args)   // 写入连接
}

sending 锁保证两件事:

  1. Header 复用安全:Client 只有一个 header 实例,多个 goroutine 并发 send 时,必须锁住"填充 header → 写入连接"这个整体操作,否则 goroutine A 填了 header,goroutine B 覆盖了 header,A 写出去的就是 B 的数据。
  2. 写入完整性:和服务端的 sending 锁一样,防止两个 goroutine 的 Header+Body 交叉写入。

mu 锁 —— 保护 pending map

go
func (client *Client) registerCall(call *Call) (uint64, error) {
    client.mu.Lock()
    defer client.mu.Unlock()
    // 操作 pending map、seq、closing、shutdown
}

func (client *Client) removeCall(seq uint64) *Call {
    client.mu.Lock()
    defer client.mu.Unlock()
    // 操作 pending map
}

mu 锁保护 pending map、seq 计数器、closing/shutdown 状态字段。这些字段被 send(多个应用 goroutine)和 receive(后台 goroutine)同时访问,必须加锁。

为什么是两把锁而不是一把

如果只用一把锁,那 send 写连接时(可能需要等待网络 IO),receive 就无法操作 pending map(因为锁被 send 持有)。两把锁分离后:

goroutine 1:sending.Lock() → 写连接(慢)→ sending.Unlock()

receive goroutine:             │  mu.Lock() → removeCall → mu.Unlock()
                                │  ↑ 不受 sending 锁影响,可以并行

发送和接收互不阻塞,最大化并发效率。


核心机制四:Go() 和 Call() —— 异步与同步的两种调用模式

Go() —— 异步调用

go
func (client *Client) Go(serviceMethod string, args, reply interface{}, done chan *Call) *Call {
    if done == nil {
        done = make(chan *Call, 10)
    }
    call := &Call{
        ServiceMethod: serviceMethod,
        Args:          args,
        Reply:         reply,
        Done:          done,
    }
    client.send(call)
    return call    // 立即返回,不等结果
}

Go() 发送请求后立即返回 Call 对象,调用方可以在任意时机通过 <-call.Done 获取结果。这是并发的基础——发起调用的 goroutine 不会被阻塞。

Call() —— 同步调用(基于 Go() 封装)

go
func (client *Client) Call(ctx context.Context, serviceMethod string, args, reply interface{}) error {
    call := client.Go(serviceMethod, args, reply, make(chan *Call, 1))
    select {
    case <-ctx.Done():
        client.removeCall(call.Seq)
        return errors.New("rpc client: call failed: " + ctx.Err().Error())
    case c := <-call.Done:
        return c.Error
    }
}

Call()Go() 的同步包装。它调用 Go() 后立即用 select 阻塞等待,直到:

  • 调用完成(call.Done 收到通知)—— 返回结果
  • context 超时/取消 —— 从 pending 中移除该调用,返回超时错误

关键点:即使是 Call() 同步调用,底层仍然是异步的。多个 goroutine 各自调用 Call() 时,它们的请求会并发发出,不会互相阻塞。每个 goroutine 只阻塞在自己的 select 上等待自己的结果。

goroutine 1:Call("A") → Go("A") → send → select 等待 call1.Done
goroutine 2:Call("B") → Go("B") → send → select 等待 call2.Done
goroutine 3:Call("C") → Go("C") → send → select 等待 call3.Done

                   三个 goroutine 各等各的,互不阻塞

核心机制五:receive 协程 —— 单消费者模型

go
func (client *Client) receive() {
    var err error
    for err == nil {
        var h codec.Header
        if err = client.cc.ReadHeader(&h); err != nil {
            break
        }
        call := client.removeCall(h.Seq)
        switch {
        case call == nil:
            err = client.cc.ReadBody(nil)
        case h.Error != "":
            call.Error = fmt.Errorf(h.Error)
            err = client.cc.ReadBody(nil)
            call.done()
        default:
            err = client.cc.ReadBody(call.Reply)
            if err != nil {
                call.Error = errors.New("reading body " + err.Error())
            }
            call.done()
        }
    }
    client.terminateCalls(err)
}

为什么只有一个 receive goroutine

和服务端 readRequest 必须串行的原因相同:TCP 是字节流,响应的 Header 和 Body 紧密相连,必须由同一个 goroutine 按序读取,否则会导致数据错乱。

但只有一个 receive goroutine 不会成为瓶颈,因为它做的事情非常轻量:

receive 的工作:
  ReadHeader  → 读几十字节,微秒级
  removeCall  → map 查找,纳秒级
  ReadBody    → 读若干字节,微秒级
  call.done() → channel 发送,纳秒级

  然后立即进入下一轮循环,处理下一个响应

真正耗时的操作(服务方法执行)在服务端,receive 只是把结果分发出去,速度极快。

call.done() 如何通知调用方

go
func (call *Call) done() {
    call.Done <- call    // 往 buffered channel 写入
}

Done 是一个带缓冲的 channel,done()call 写入后立即返回,不会阻塞 receive 协程。调用方(在 Call()select 中或在 <-call.Done 处等待)收到通知后继续执行。


异常处理与并发安全

连接断开:terminateCalls

go
func (client *Client) terminateCalls(err error) {
    client.sending.Lock()       // 先拿 sending 锁,阻止新的请求发送
    defer client.sending.Unlock()
    client.mu.Lock()            // 再拿 mu 锁,操作 pending map
    defer client.mu.Unlock()
    client.shutdown = true      // 标记为异常关闭
    for _, call := range client.pending {
        call.Error = err        // 给每个未完成的调用设置错误
        call.done()             // 通知调用方
    }
}

当 receive 协程读取失败时(服务端关闭连接、网络断开等),它会跳出循环并调用 terminateCalls。这个函数做了两件事:

  1. 获取两把锁:先 sendingmu,确保此刻没有正在进行的发送操作,也没有其他 goroutine 在操作 pending map。锁的顺序固定(先 sending 后 mu),避免死锁。
  2. 批量通知:遍历所有未完成的调用,将连接错误填入每个 call,然后 call.done() 通知调用方。调用方从 select 中醒来,拿到错误信息。

context 超时取消

go
case <-ctx.Done():
    client.removeCall(call.Seq)
    return errors.New("rpc client: call failed: " + ctx.Err().Error())

当 context 超时时,Call() 主动从 pending map 中移除该调用。之后如果 receive 协程收到了这个 seq 的响应:

go
call := client.removeCall(h.Seq)
switch {
case call == nil:                      // ← 走这个分支
    err = client.cc.ReadBody(nil)      // 丢弃 Body,保持流位置正确

removeCall 返回 nil(已经被 context 超时移除了),receive 优雅地丢弃这个迟到的响应,不会报错。


与服务端并发模型的对比

维度ServerClient
并发单元每个请求一个 goroutine每个调用方一个 goroutine(由应用代码控制)
读取方式单 goroutine 串行 readRequest单 goroutine 串行 receive
写入保护sending Mutex 保护 sendResponsesending Mutex 保护 send
请求匹配不需要(读到即处理)pending map + Seq 序列号
等待完成sync.WaitGroupDone channel + context
异常处理循环退出 + wg.Wait + cc.CloseterminateCalls 批量通知

两端的核心思想一致:读取串行(TCP 字节流约束)、处理/调用并行(goroutine 廉价)、写入互斥(锁保护完整性)。区别在于 Client 多了 pending map 做请求-响应匹配,因为 Client 需要在一条连接上追踪多个并发调用的状态。


实际并发场景示例

以 example/main.go 中的异步调用为例:

go
call1 := client.Go("Arith.Multiply", MulArgs{A: 3, B: 5}, new(int), nil)
call2 := client.Go("Arith.Multiply", MulArgs{A: 4, B: 9}, new(int), nil)

r1 := <-call1.Done
r2 := <-call2.Done

完整的执行时间线:

主 goroutine:
  Go("3*5") → send(call1) → pending[1]=call1 → 写连接 → 返回 call1
  Go("4*9") → send(call2) → pending[2]=call2 → 写连接 → 返回 call2
  <-call1.Done(阻塞等待)

receive goroutine:                 │
  ReadHeader → Seq=2               │
  removeCall(2) → call2            │
  ReadBody → reply=36              │
  call2.done()  ────────────────────│──► call2.Done 收到通知(但主 goroutine 在等 call1)
  ReadHeader → Seq=1               │
  removeCall(1) → call1            │
  ReadBody → reply=15              │
  call1.done()  ────────────────────┘──► call1.Done 收到通知,主 goroutine 继续
                                        然后 <-call2.Done 立即返回(已有数据)

两个请求在一条连接上并发完成,响应可能乱序到达(这里 Seq=2 先回来),但通过 pending map 正确匹配,调用方拿到正确的结果。


三种 Dial 方式的设计与区别

总体架构

miniRPC 提供了三个拨号函数,形成分层设计

                          XDial("http@localhost:9999")
                          XDial("tcp@localhost:9999")

                         解析 "protocol@address"

                  ┌─────────────────┼─────────────────┐
                  ▼                                   ▼
          DialHTTP("tcp", addr)              Dial("tcp", addr)
                  │                                   │
                  ▼                                   ▼
      dialTimeout(newClientHTTP, ...)     dialTimeout(newClient, ...)
                  │                                   │
          ┌───────┴───────┐                   ┌───────┴───────┐
          ▼               ▼                   ▼               ▼
   net.DialTimeout    newClientHTTP     net.DialTimeout    newClient
   (TCP 连接建立)     (HTTP CONNECT      (TCP 连接建立)    (直接发送
                       协议切换                             Option 握手)
                       → newClient)

三者最终都通过 dialTimeout 建立 TCP 连接,区别在于连接建立后的握手方式不同


Dial —— 直接 TCP 连接

go
func Dial(network, address string, opts ...*codec.Option) (*Client, error) {
    return dialTimeout(newClient, network, address, opts...)
}

最简单直接的方式。连接建立后立即用 newClient 进行 miniRPC 握手:

客户端                                        服务端(server.Accept 监听)

net.DialTimeout("tcp", "localhost:9999")
     │                                            │
     │ ◄───────── TCP 三次握手 ──────────────────► │
     │                                            │
     │  json.Encode(Option)                       │
     │ ──────────────────────────────────────────► │ json.Decode(&opt)
     │                                            │ 校验 MagicNumber
     │                                            │ 创建 Codec
     │                                            │
     │ ◄═══════ 进入 RPC 通信(Codec 编码)═══════► │

使用场景:客户端和服务端之间是直连 TCP,没有 HTTP 代理或反向代理的介入。这是最常用、最高效的方式。


DialHTTP —— 通过 HTTP CONNECT 协议切换

go
func DialHTTP(network, address string, opts ...*codec.Option) (*Client, error) {
    return dialTimeout(newClientHTTP, network, address, opts...)
}

连接建立后,先走一轮 HTTP CONNECT 握手,再进入 miniRPC 握手:

客户端                                        服务端(http.Handle 监听)

net.DialTimeout("tcp", "localhost:9999")
     │                                            │
     │ ◄───────── TCP 三次握手 ──────────────────► │
     │                                            │
     │  "CONNECT /_minirpc_ HTTP/1.0\n\n"         │
     │ ──────────────────────────────────────────► │ ServeHTTP()
     │                                            │ 检查 Method == "CONNECT"
     │  "HTTP/1.0 200 Connected to miniRPC\n\n"   │ Hijack() 夺取连接
     │ ◄────────────────────────────────────────── │
     │                                            │
     │  json.Encode(Option)                       │
     │ ──────────────────────────────────────────► │ json.Decode(&opt)(同 Dial)
     │                                            │
     │ ◄═══════ 进入 RPC 通信(Codec 编码)═══════► │

为什么需要 HTTP CONNECT

核心原因是复用 HTTP 基础设施

  1. 穿透代理:很多企业网络中,客户端不能直接建立 TCP 连接,必须经过 HTTP 代理。HTTP CONNECT 是代理隧道的标准方法——代理看到 CONNECT 请求后会建立一条透明的 TCP 隧道,后续数据原样转发。
  2. 共享端口:服务端可以在同一个端口上同时提供 HTTP 服务(如调试页面 /debug/minirpc)和 RPC 服务。普通 HTTP 请求走正常的 Handler,CONNECT 请求切换到 RPC 模式。
  3. Hijack 机制:服务端的 ServeHTTP 通过 Hijacker.Hijack() 从 HTTP 框架手中"夺取"底层的原始 TCP 连接,之后这个连接就脱离了 HTTP 框架的管理,完全由 miniRPC 的 ServeConn 接管。

newClientHTTP 的实现细节

go
func newClientHTTP(conn net.Conn, opt *codec.Option) (*Client, error) {
    // 第 1 步:发送 HTTP CONNECT 请求
    _, _ = io.WriteString(conn, fmt.Sprintf("CONNECT %s HTTP/1.0\n\n", DefaultRPCPath))

    // 第 2 步:读取 HTTP 响应
    resp, err := http.ReadResponse(bufio.NewReader(conn), &http.Request{Method: "CONNECT"})
    if err == nil && resp.Status == "200 Connected to miniRPC" {
        // 第 3 步:协议切换成功,后续走标准 miniRPC 握手
        return newClient(conn, opt)
    }
    // 协议切换失败
    if err == nil {
        err = errors.New("unexpected HTTP response: " + resp.Status)
    }
    return nil, err
}

注意:CONNECT 成功后调用的是 newClient(conn, opt)——同一个 TCP 连接先走 HTTP 握手,再走 miniRPC 握手,后续的通信和 Dial 完全一样。HTTP CONNECT 只是一个"敲门"环节。


XDial —— 统一入口,自动选择协议

go
func XDial(rpcAddr string, opts ...*codec.Option) (*Client, error) {
    parts := strings.Split(rpcAddr, "@")
    protocol, addr := parts[0], parts[1]
    switch protocol {
    case "http":
        return DialHTTP("tcp", addr, opts...)
    default:
        return Dial(protocol, addr, opts...)
    }
}

XDial 是一个路由层,根据地址格式自动选择连接方式:

调用解析结果实际调用
XDial("tcp@localhost:9999")protocol=tcp, addr=localhost:9999Dial("tcp", "localhost:9999")
XDial("http@localhost:9999")protocol=http, addr=localhost:9999DialHTTP("tcp", "localhost:9999")
XDial("unix@/tmp/rpc.sock")protocol=unix, addr=/tmp/rpc.sockDial("unix", "/tmp/rpc.sock")

为什么需要 XDial

服务发现(Day 6 XClient)场景中,服务发现返回的地址是一个字符串(如 "tcp@10.0.0.1:8080"),调用方不需要关心底层协议细节,只需把地址字符串传给 XDial,它自动做对。

go
// 服务发现场景(Day 6 XClient 中的实际用法)
addr := discovery.GetServer()          // 返回 "tcp@10.0.0.1:8080" 或 "http@10.0.0.2:8080"
client, err := minirpc.XDial(addr)     // 自动选择 Dial 还是 DialHTTP

三者共享的底层:dialTimeout

无论哪种 Dial 方式,最终都经过 dialTimeout,它是超时控制的统一入口

go
func dialTimeout(f newClientFunc, network, address string, opts ...*codec.Option) (*Client, error) {
    opt := parseOptions(opts...)
    conn, err := net.DialTimeout(network, address, opt.ConnectTimeout)   // ① TCP 连接(带超时)
    if err != nil {
        return nil, err
    }

    ch := make(chan clientResult, 1)
    go func() {
        client, err := f(conn, opt)    // ② 在子 goroutine 中执行握手(f 是 newClient 或 newClientHTTP)
        ch <- clientResult{client: client, err: err}
    }()

    if opt.ConnectTimeout == 0 {       // 无超时,直接等结果
        result := <-ch
        return result.client, result.err
    }
    select {                           // ③ 超时控制
    case <-time.After(opt.ConnectTimeout):
        _ = conn.Close()
        return nil, fmt.Errorf("rpc client: connect timeout: expect within %s", opt.ConnectTimeout)
    case result := <-ch:
        return result.client, result.err
    }
}

关键设计:dialTimeout 接收一个 newClientFunc 参数(函数类型),Dial 传入 newClientDialHTTP 传入 newClientHTTP。这样超时控制逻辑只写一次,不同的握手方式通过传入不同的函数来差异化。这是 Go 中用函数作为参数实现策略模式的典型做法。


总结对比

维度DialDialHTTPXDial
握手方式直接发送 Option先 HTTP CONNECT,再发送 Option自动选择
握手步骤1 步2 步(HTTP + Option)取决于协议
服务端要求server.Accept(lis)server.HandleHTTP() + HTTP 服务取决于协议
适用场景内网直连,性能优先需要穿透 HTTP 代理或共享 HTTP 端口服务发现,地址格式统一
性能最优(少一轮 HTTP 握手)略差(多一轮 HTTP 往返)同所选方式
底层实现dialTimeout(newClient, ...)dialTimeout(newClientHTTP, ...)路由到 Dial 或 DialHTTP