主题
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 锁保证两件事:
- Header 复用安全:Client 只有一个
header实例,多个 goroutine 并发 send 时,必须锁住"填充 header → 写入连接"这个整体操作,否则 goroutine A 填了 header,goroutine B 覆盖了 header,A 写出去的就是 B 的数据。 - 写入完整性:和服务端的
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。这个函数做了两件事:
- 获取两把锁:先
sending后mu,确保此刻没有正在进行的发送操作,也没有其他 goroutine 在操作 pending map。锁的顺序固定(先 sending 后 mu),避免死锁。 - 批量通知:遍历所有未完成的调用,将连接错误填入每个 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 优雅地丢弃这个迟到的响应,不会报错。
与服务端并发模型的对比
| 维度 | Server | Client |
|---|---|---|
| 并发单元 | 每个请求一个 goroutine | 每个调用方一个 goroutine(由应用代码控制) |
| 读取方式 | 单 goroutine 串行 readRequest | 单 goroutine 串行 receive |
| 写入保护 | sending Mutex 保护 sendResponse | sending Mutex 保护 send |
| 请求匹配 | 不需要(读到即处理) | pending map + Seq 序列号 |
| 等待完成 | sync.WaitGroup | Done channel + context |
| 异常处理 | 循环退出 + wg.Wait + cc.Close | terminateCalls 批量通知 |
两端的核心思想一致:读取串行(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 基础设施:
- 穿透代理:很多企业网络中,客户端不能直接建立 TCP 连接,必须经过 HTTP 代理。HTTP CONNECT 是代理隧道的标准方法——代理看到 CONNECT 请求后会建立一条透明的 TCP 隧道,后续数据原样转发。
- 共享端口:服务端可以在同一个端口上同时提供 HTTP 服务(如调试页面
/debug/minirpc)和 RPC 服务。普通 HTTP 请求走正常的 Handler,CONNECT 请求切换到 RPC 模式。 - 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:9999 | Dial("tcp", "localhost:9999") |
XDial("http@localhost:9999") | protocol=http, addr=localhost:9999 | DialHTTP("tcp", "localhost:9999") |
XDial("unix@/tmp/rpc.sock") | protocol=unix, addr=/tmp/rpc.sock | Dial("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 传入 newClient,DialHTTP 传入 newClientHTTP。这样超时控制逻辑只写一次,不同的握手方式通过传入不同的函数来差异化。这是 Go 中用函数作为参数实现策略模式的典型做法。
总结对比
| 维度 | Dial | DialHTTP | XDial |
|---|---|---|---|
| 握手方式 | 直接发送 Option | 先 HTTP CONNECT,再发送 Option | 自动选择 |
| 握手步骤 | 1 步 | 2 步(HTTP + Option) | 取决于协议 |
| 服务端要求 | server.Accept(lis) | server.HandleHTTP() + HTTP 服务 | 取决于协议 |
| 适用场景 | 内网直连,性能优先 | 需要穿透 HTTP 代理或共享 HTTP 端口 | 服务发现,地址格式统一 |
| 性能 | 最优(少一轮 HTTP 握手) | 略差(多一轮 HTTP 往返) | 同所选方式 |
| 底层实现 | dialTimeout(newClient, ...) | dialTimeout(newClientHTTP, ...) | 路由到 Dial 或 DialHTTP |