学习目标
学完本章,你应该能够:
- 说清楚一个 RPC 协议为什么要分"头"和"体":能用自己的话解释协议头里放路由/元数据、协议体放业务数据的取舍,以及为什么中间件(网关、sidecar)只解析头部就能做路由。
- 设计一个最小可用的请求/响应协议:讲得清头部要有哪些字段(长度、版本、序列化标记、压缩标记、消息 ID、服务名、方法名、错误),以及它们分别解决什么问题(粘包、升级、多路复用、版本兼容)。
- 实现多序列化协议支持:讲清
Serializer接口 + 注册表模式如何工作,知道客户端和服务端如何通过协议头里的 1 字节标记协商序列化方式。 - 分辨三种调用语义并落地单向调用:能讲清异步/回调/单向调用的区别,能解释"真假单向调用"对资源释放的影响,并看懂
WithOneWay的实现。 - 讲清链路超时控制的原理:面试时能把这个知识点讲成一个连贯的故事——单一超时 vs 链路超时、用
select + channel监听context.Done()、为什么用"时间戳"而不是"剩余时间"跨端传递。
前置知识:
- 上一章"最简 RPC 框架"(JSON 传输、4 字节长度前缀、基本的客户端/服务端)。
- Go 基础:
net、encoding/json、goroutine 与channel、context包的基本用法。 - 一点网络常识:TCP 是字节流、有"粘包"问题,HTTP/2、gRPC、Dubbo 名词听说过即可。
本章你会动手做的事:
- 在纸上画出"请求协议头"的字段布局,标注哪些是定长、哪些是不定长。
- 把
EncodeRequest/DecodeRequest的编码格式[头长度][头][体长度][体]在草稿上模拟一遍(给一个假头部,写出最终字节流)。 - 用
CallWithTimeout跑一次带 3 秒超时的调用,故意把服务端 sleep 5 秒,观察客户端是否按时返回context deadline exceeded。
前言:从最简 RPC 到真正的 RPC 协议
在上一章中,我们手写了一个"最简 RPC 框架"——它能用 JSON 格式传递调用信息,实现基本的远程调用。但这个框架有一个明显的局限:消息格式被写死成 JSON,不支持其他序列化协议。
而在实际应用中,一个成熟的 RPC 框架还需要考虑:
- 支持不同的序列化协议(JSON、Protobuf、Gob 等)
- 支持数据压缩(gzip、snappy、zstd 等)
- 支持协议版本升级(未来字段变更时不破坏旧客户端)
- 支持加密传输
- 支持链路追踪(trace ID 传递)
- 支持超时控制和单向调用等高级语义
这一切都引出一个核心问题:如何设计一个合适的 RPC 协议?
本教程将带你从协议设计理论出发,逐步实现一个支持多序列化协议、多调用语义和链路超时控制的 RPC 框架。
本教程涵盖以下主题:
- 协议设计理论——协议头与协议体,gRPC/Dubbo/TCP 协议对比
- 请求与响应协议——字段设计、接口定义、编解码实现
- 多序列化协议支持——Serializer 接口与 JSON/Proto 实现
- RPC 调用语义——异步调用、回调、单向调用
- RPC 超时控制——单一超时与链路超时、跨端传递
- 面试要点总结
一、协议设计理论
1.1 为什么需要协议设计
类比:最简 RPC 就像一个把"信纸内容 + 信封信息"全写在同一张纸上、塞进信封就寄走的邮寄方式。对方收到后,必须从头到尾读完整张纸,才知道该把这封信交给谁、要不要回信、用什么语言读。想换成"挂号信"(带追踪、带压缩)?对不起,纸张格式是写死的,改不了。
真正的 RPC 协议相当于把"信封"和"信纸"分开:**信封(协议头)**写清收件人、邮路标记、是否加急;**信纸(协议体)**只写正事。邮局(网关/sidecar)只看信封就能分拣转发,不用拆信。
在最简 RPC 中,我们直接用 JSON 序列化整个 Request 结构体,然后加上一个 4 字节的长度前缀发送出去。这种方式简单粗暴,但存在很多问题:
- 无法切换序列化协议:JSON 写死在代码里,想换成 Protobuf 需要改大量代码
- 无法支持压缩:大数据传输时没有压缩机制,浪费带宽
- 无法支持版本升级:协议字段变更后,旧客户端无法兼容
- 无法传递元数据:链路追踪、A/B 测试标记等无处安放
解决这些问题的关键就是设计一个结构化的 RPC 协议。
1.2 业界协议参考
gRPC 协议
gRPC 协议分成头部和 body 两个部分。这是因为 gRPC 是直接基于 HTTP/2 来实现的,所以 gRPC 的头部就放在 HTTP 协议头,body 就放在 HTTP 协议体里面。
HTTP/2 HEADERS (gRPC 头部)
:method = POST
:path = /package.Service/Method
content-type = application/grpc
grpc-encoding = gzip
grpc-timeout = 1S
...
HTTP/2 DATA (gRPC body)
[Compressed-Flag][Message-Length][Message-Data]
Dubbo 协议
Dubbo 协议按照它自己的说法,是分成定长部分和非定长部分。其中定长部分一般叫协议头,变长部分一般叫做协议体。
| Magic | Flag | Status | RequestID (8 bytes) | Data Length |
| Header (16 bytes fixed) |
| Body (variable) |
Dubbo 在 3.0 里面推出了 Triple 协议,就很接近 gRPC 的设计了。原生的 Dubbo 协议反而和 gRPC 很不一样。
TCP 协议
gRPC 和 Dubbo 协议都可以看做是应用层协议,其它协议设计也是类似的,比如 TCP 协议本身也是分成头部和数据部分。
1.3 协议设计总结
通过对比不同协议,我们可以总结出协议设计的通用原则:
- 协议一般都是分成两个部分:协议头和协议体
- 协议头包含接收方"如何处理这个消息"的必要信息,具体来说包含描述协议本身的数据,和描述这次请求的数据
- 协议体则大多数情况下存放请求数据
从形态上来说:
- 协议体必然是变长的
- 协议头可以是定长的,也可以是变长的
- 在变长协议头的情况下,需要有两个字段来描述长度:一个描述协议头有多长,一个描述协议体有多长
flowchart TB
subgraph 消息格式
H[协议头 Header] --> B[协议体 Body]
end
subgraph 协议头内容
H1[描述协议本身的数据
版本/序列化协议/压缩算法] --> H2[描述这次请求的数据
服务名/方法名/消息ID/元数据]
end
subgraph 协议体内容
B1[请求参数 / 响应数据]
end
H --> H1
B --> B1二、请求与响应协议设计
2.1 请求协议头设计
现在设计我们自己的 RPC 协议。请求头部设计为不定长,包含以下字段:
固定字段:
| 字段 | 说明 | 用途 |
|---|---|---|
| 长度字段 | 协议头长度 + 协议体长度 | 用于分割消息(解决粘包问题) |
| 版本字段 | 协议版本号 | 用于后续协议升级 |
| 序列化协议 | 序列化方式标记 | 用于标记采用的序列化协议 |
| 压缩算法 | 压缩方式标记 | 用于标记协议体是如何被压缩的 |
| 消息 ID | 唯一标识 | 用于后续支持多路复用 |
| 服务名 | 目标服务名 | 服务端据此查找服务 |
| 方法名 | 目标方法名 | 服务端据此查找方法 |
不固定字段:
主要是链路元数据,例如 trace ID、A/B 测试标记、全链路压测的标记位等。这些数据以 key-value 形式存储。
最后的协议体里面就只存放请求参数。
2.2 为什么服务名/方法名放头部
服务名、方法名和元数据都可以考虑放到请求参数(协议体)里,之所以放在头部是因为:
如果我们的微服务请求要经过网关、sidecar(service mesh),那么放在头部,这些中间件就可以考虑只解析头部字段,而不必解析整个请求。
例如在 sidecar 上做负载均衡,那么只需要解析到服务名和方法名就可以,根据服务名和方法名找到可用节点,然后做负载均衡。如果这些信息放在协议体里,sidecar 就需要反序列化整个请求体才能获取,性能开销很大。
flowchart LR
C[客户端] -->|请求| S1[Sidecar
只解析头部]
S1 -->|转发| S2[Sidecar
只解析头部]
S2 -->|转发| SRV[服务端
解析完整请求]
style S1 fill:#e1f5fe
style S2 fill:#e1f5fe2.3 响应协议头设计
响应头部也设计为不定长,包含以下字段:
| 字段 | 说明 | 用途 |
|---|---|---|
| 长度字段 | 协议头长度 + 协议体长度 | 用于分割消息 |
| 版本字段 | 协议版本号 | 用于后续协议升级 |
| 序列化协议 | 序列化方式标记 | 用于标记采用的序列化协议 |
| 压缩算法 | 压缩方式标记 | 用于标记协议体是如何被压缩的 |
| 消息 ID | 唯一标识 | 用于匹配请求和响应 |
| 错误 | 错误信息 | 为了解决第二个返回值的问题 |
最后的协议体里面就只存放响应数据。
2.4 响应错误为什么放头部
主要是实在没地方放。从理论上来说,有两个选择:
- 放头部:做成一个类似于 serviceName 那种,认为是服务调用本身必需的一种数据
- 放协议体:也就是认为是响应体的一部分。这种做法会更加符合直觉,但实现起来难度更高,因为你需要区别响应数据里面,哪部分是用户返回的数据,哪部分是用户返回的错误
我们选择放头部,这样协议体可以干净地只包含业务数据,编解码逻辑更简单。
三、接口定义与编解码实现
3.1 改造 Request 和 Response
首先改造我们的 Request 和 Response,把刚才决定的协议字段加进去:
package mrpc
import (
"encoding/binary"
"encoding/json"
"fmt"
"io"
)
// ============================================================
// 第一部分:协议数据结构定义
// ============================================================
// RequestHeader 是请求的协议头
// 包含描述协议本身的数据和描述这次请求的数据
type RequestHeader struct {
// --- 描述协议本身的数据 ---
Version uint8 // 协议版本号:用于后续协议升级,当前版本为 1
Serializer uint8 // 序列化协议标记:1=JSON, 2=Protobuf, 3=Gob
Compressor uint8 // 压缩算法标记:0=不压缩, 1=gzip, 2=snappy
MessageID uint64 // 消息ID:唯一标识一个请求,用于多路复用
// --- 描述这次请求的数据 ---
Service string // 服务名:服务端据此查找目标服务
Method string // 方法名:服务端据此查找目标方法
Meta map[string]string // 元数据:链路追踪ID、A/B测试标记、超时时间等
}
// Request 是完整的 RPC 请求
// 由协议头和协议体组成
type Request struct {
Header RequestHeader // 协议头:包含路由和元数据信息
Body []byte // 协议体:序列化后的请求参数(已经是字节流)
}
// ResponseHeader 是响应的协议头
type ResponseHeader struct {
// --- 描述协议本身的数据 ---
Version uint8 // 协议版本号
Serializer uint8 // 序列化协议标记
Compressor uint8 // 压缩算法标记
MessageID uint64 // 消息ID:用于匹配请求和响应
// --- 描述这次响应的数据 ---
Error string // 错误信息:空字符串表示成功
}
// Response 是完整的 RPC 响应
type Response struct {
Header ResponseHeader // 协议头
Body []byte // 协议体:序列化后的响应数据
}
// ============================================================
// 第二部分:Request 编解码
// ============================================================
// EncodeRequest 将 Request 编码成字节流
// 编码格式:[头部长度(4字节)][头部JSON][体长度(4字节)][体数据]
// 这里用长度前缀来解决 TCP 粘包问题
func EncodeRequest(req *Request) ([]byte, error) {
// 第一步:将 Header 序列化为 JSON
// Header 是结构化数据,用 JSON 编码方便扩展
headerData, err := json.Marshal(req.Header)
if err != nil {
return nil, fmt.Errorf("编码请求头失败: %w", err)
}
// 第二步:组装最终的字节流
// 格式:[头部长度(4字节)] + [头部JSON] + [体长度(4字节)] + [体数据]
// 使用 4 字节的 uint32 来表示长度,最大支持约 4GB 的消息
result := make([]byte, 0, 4+len(headerData)+4+len(req.Body))
// 写入头部长度(4字节,大端序)
headerLenBuf := make([]byte, 4)
binary.BigEndian.PutUint32(headerLenBuf, uint32(len(headerData)))
result = append(result, headerLenBuf...)
// 写入头部数据
result = append(result, headerData...)
// 写入体长度(4字节,大端序)
bodyLenBuf := make([]byte, 4)
binary.BigEndian.PutUint32(bodyLenBuf, uint32(len(req.Body)))
result = append(result, bodyLenBuf...)
// 写入体数据
result = append(result, req.Body...)
return result, nil
}
// DecodeRequest 从字节流中解码出 Request
// 解码过程与编码过程正好相反:先读头部长度 -> 读头部 -> 读体长度 -> 读体
func DecodeRequest(data []byte) (*Request, error) {
if len(data) < 8 {
return nil, fmt.Errorf("数据太短,无法解码")
}
// 第一步:读取头部长度
headerLen := binary.BigEndian.Uint32(data[:4])
if len(data) < int(4+headerLen+4) {
return nil, fmt.Errorf("数据不完整,头部长度不匹配")
}
// 第二步:读取并反序列化头部
headerData := data[4 : 4+headerLen]
var header RequestHeader
if err := json.Unmarshal(headerData, &header); err != nil {
return nil, fmt.Errorf("解码请求头失败: %w", err)
}
// 第三步:读取体长度
bodyLen := binary.BigEndian.Uint32(data[4+headerLen : 4+headerLen+4])
if len(data) < int(4+headerLen+4+bodyLen) {
return nil, fmt.Errorf("数据不完整,体长度不匹配")
}
// 第四步:读取体数据
bodyData := data[4+headerLen+4 : 4+headerLen+4+bodyLen]
// 组装 Request
req := &Request{
Header: header,
Body: make([]byte, bodyLen),
}
copy(req.Body, bodyData)
return req, nil
}
// ============================================================
// 第三部分:Response 编解码
// ============================================================
// EncodeResponse 将 Response 编码成字节流
// 编码格式与 Request 一致:[头部长度(4字节)][头部JSON][体长度(4字节)][体数据]
func EncodeResponse(resp *Response) ([]byte, error) {
// 将 Header 序列化为 JSON
headerData, err := json.Marshal(resp.Header)
if err != nil {
return nil, fmt.Errorf("编码响应头失败: %w", err)
}
// 组装字节流
result := make([]byte, 0, 4+len(headerData)+4+len(resp.Body))
// 写入头部长度
headerLenBuf := make([]byte, 4)
binary.BigEndian.PutUint32(headerLenBuf, uint32(len(headerData)))
result = append(result, headerLenBuf...)
// 写入头部数据
result = append(result, headerData...)
// 写入体长度
bodyLenBuf := make([]byte, 4)
binary.BigEndian.PutUint32(bodyLenBuf, uint32(len(resp.Body)))
result = append(result, bodyLenBuf...)
// 写入体数据
result = append(result, resp.Body...)
return result, nil
}
// DecodeResponse 从字节流中解码出 Response
// 解码过程与编码过程相反
func DecodeResponse(data []byte) (*Response, error) {
if len(data) < 8 {
return nil, fmt.Errorf("数据太短,无法解码")
}
// 读取头部长度
headerLen := binary.BigEndian.Uint32(data[:4])
if len(data) < int(4+headerLen+4) {
return nil, fmt.Errorf("数据不完整")
}
// 反序列化头部
headerData := data[4 : 4+headerLen]
var header ResponseHeader
if err := json.Unmarshal(headerData, &header); err != nil {
return nil, fmt.Errorf("解码响应头失败: %w", err)
}
// 读取体长度和体数据
bodyLen := binary.BigEndian.Uint32(data[4+headerLen : 4+headerLen+4])
bodyData := data[4+headerLen+4 : 4+headerLen+4+bodyLen]
resp := &Response{
Header: header,
Body: make([]byte, bodyLen),
}
copy(resp.Body, bodyData)
return resp, nil
}
// ============================================================
// 第四部分:网络读写辅助函数
// ============================================================
// WriteRequest 将 Request 写入网络连接
// 先编码成字节流,再写入 conn
func WriteRequest(w io.Writer, req *Request) error {
data, err := EncodeRequest(req)
if err != nil {
return err
}
_, err = w.Write(data)
return err
}
// ReadRequest 从网络连接读取 Request
// 先读取全部数据,再解码
func ReadRequest(r io.Reader) (*Request, error) {
// 先读 4 字节的头部长度
headerLenBuf := make([]byte, 4)
if _, err := io.ReadFull(r, headerLenBuf); err != nil {
return nil, err
}
headerLen := binary.BigEndian.Uint32(headerLenBuf)
// 读取头部数据
headerData := make([]byte, headerLen)
if _, err := io.ReadFull(r, headerData); err != nil {
return nil, err
}
// 读 4 字节的体长度
bodyLenBuf := make([]byte, 4)
if _, err := io.ReadFull(r, bodyLenBuf); err != nil {
return nil, err
}
bodyLen := binary.BigEndian.Uint32(bodyLenBuf)
// 读取体数据
bodyData := make([]byte, bodyLen)
if _, err := io.ReadFull(r, bodyData); err != nil {
return nil, err
}
// 组装完整数据并解码
fullData := make([]byte, 0, 4+len(headerData)+4+len(bodyData))
fullData = append(fullData, headerLenBuf...)
fullData = append(fullData, headerData...)
fullData = append(fullData, bodyLenBuf...)
fullData = append(fullData, bodyData...)
return DecodeRequest(fullData)
}
// WriteResponse 将 Response 写入网络连接
func WriteResponse(w io.Writer, resp *Response) error {
data, err := EncodeResponse(resp)
if err != nil {
return err
}
_, err = w.Write(data)
return err
}
// ReadResponse 从网络连接读取 Response
func ReadResponse(r io.Reader) (*Response, error) {
// 读头部长度
headerLenBuf := make([]byte, 4)
if _, err := io.ReadFull(r, headerLenBuf); err != nil {
return nil, err
}
headerLen := binary.BigEndian.Uint32(headerLenBuf)
// 读头部数据
headerData := make([]byte, headerLen)
if _, err := io.ReadFull(r, headerData); err != nil {
return nil, err
}
// 读体长度
bodyLenBuf := make([]byte, 4)
if _, err := io.ReadFull(r, bodyLenBuf); err != nil {
return nil, err
}
bodyLen := binary.BigEndian.Uint32(bodyLenBuf)
// 读体数据
bodyData := make([]byte, bodyLen)
if _, err := io.ReadFull(r, bodyData); err != nil {
return nil, err
}
// 组装并解码
fullData := make([]byte, 0, 4+len(headerData)+4+len(bodyData))
fullData = append(fullData, headerLenBuf...)
fullData = append(fullData, headerData...)
fullData = append(fullData, bodyLenBuf...)
fullData = append(fullData, bodyData...)
return DecodeResponse(fullData)
}
3.2 编解码流程图解
flowchart LR
subgraph 编码 EncodeRequest
E1[RequestHeader] -->|json.Marshal| E2[头部JSON字节]
E2 -->|拼接长度前缀| E3[完整字节流]
E4[Body字节] -->|拼接| E3
end
subgraph 解码 DecodeRequest
D1[完整字节流] -->|读4字节| D2[头部长度]
D2 -->|读取头部| D3[头部JSON字节]
D3 -->|json.Unmarshal| D4[RequestHeader]
D1 -->|读4字节| D5[体长度]
D5 -->|读取体| D6[Body字节]
end编码差不多就是一个个字段拼起来,解码过程也是类似的,一个个字段处理过去。
四、多序列化协议支持
4.1 Serializer 接口设计
在实际应用中,大多数 RPC 框架都可以支持不同的序列化协议。我们通过定义一个 Serializer 接口来实现这一点:
package mrpc
// ============================================================
// 序列化协议接口与实现
// ============================================================
// Serializer 序列化协议接口
// 不同的序列化协议(JSON、Protobuf、Gob)都实现这个接口
// 这样客户端和服务端可以通过同一个接口使用不同的序列化方式
type Serializer interface {
// Encode 将数据序列化为字节流
// 参数 data 是任意类型的结构体指针
// 返回序列化后的字节切片
Encode(data interface{}) ([]byte, error)
// Decode 将字节流反序列化为数据
// 参数 data 是目标结构体的指针(用于接收反序列化结果)
// 参数 bytes 是序列化后的字节切片
Decode(data interface{}, bytes []byte) error
// Code 返回序列化协议的唯一标识码
// 用于在协议头中标记使用的序列化方式
// 客户端和服务端之间需要协商一致
// 例如:1=JSON, 2=Protobuf, 3=Gob
Code() uint8
}
4.2 JSON 序列化协议实现
import "encoding/json"
// JSONSerializer 使用 JSON 格式进行序列化
// JSON 是最通用的序列化格式,可读性好,但体积较大
// 适合开发调试阶段使用
type JSONSerializer struct{}
// Encode 将数据序列化为 JSON 字节流
func (s *JSONSerializer) Encode(data interface{}) ([]byte, error) {
return json.Marshal(data)
}
// Decode 将 JSON 字节流反序列化为数据
func (s *JSONSerializer) Decode(data interface{}, bytes []byte) error {
return json.Unmarshal(bytes, data)
}
// Code 返回 JSON 序列化协议的标识码
// 约定:1 代表 JSON
func (s *JSONSerializer) Code() uint8 {
return 1
}
4.3 Protobuf 序列化协议实现
import "google.golang.org/protobuf/proto"
// ProtoSerializer 使用 Protobuf 格式进行序列化
// Protobuf 体积小、速度快,但需要预先定义 .proto 文件并生成代码
// 适合生产环境使用
type ProtoSerializer struct{}
// Encode 将数据序列化为 Protobuf 字节流
// 注意:data 必须是 proto.Message 接口的实现
func (s *ProtoSerializer) Encode(data interface{}) ([]byte, error) {
// 类型断言:将 interface{} 转换为 proto.Message
msg, ok := data.(proto.Message)
if !ok {
return nil, fmt.Errorf("数据未实现 proto.Message 接口")
}
return proto.Marshal(msg)
}
// Decode 将 Protobuf 字节流反序列化为数据
func (s *ProtoSerializer) Decode(data interface{}, bytes []byte) error {
msg, ok := data.(proto.Message)
if !ok {
return fmt.Errorf("数据未实现 proto.Message 接口")
}
return proto.Unmarshal(bytes, msg)
}
// Code 返回 Protobuf 序列化协议的标识码
// 约定:2 代表 Protobuf
func (s *ProtoSerializer) Code() uint8 {
return 2
}
4.4 序列化协议注册表
为了让客户端和服务端能够根据协议头中的 Serializer 字段找到对应的序列化器,我们需要一个注册表:
import "sync"
// ============================================================
// 序列化协议注册表
// ============================================================
// SerializerRegistry 序列化协议注册表
// 根据 Code 查找对应的 Serializer 实例
// 客户端和服务端都需要注册自己支持的序列化协议
type SerializerRegistry struct {
mu sync.RWMutex // 读写锁,保护并发访问
serializers map[uint8]Serializer // 存储已注册的序列化器
}
// NewSerializerRegistry 创建一个新的序列化协议注册表
func NewSerializerRegistry() *SerializerRegistry {
return &SerializerRegistry{
serializers: make(map[uint8]Serializer),
}
}
// Register 注册一个序列化协议
// 参数 serializer 是实现了 Serializer 接口的实例
func (r *SerializerRegistry) Register(serializer Serializer) {
r.mu.Lock()
defer r.mu.Unlock()
// 以 Code 为 key 存储序列化器
r.serializers[serializer.Code()] = serializer
}
// Get 根据 Code 获取对应的序列化器
func (r *SerializerRegistry) Get(code uint8) (Serializer, error) {
r.mu.RLock()
defer r.mu.RUnlock()
s, ok := r.serializers[code]
if !ok {
return nil, fmt.Errorf("未找到序列化协议: code=%d", code)
}
return s, nil
}
4.5 改造客户端和服务端
改造后的客户端和服务端都需要持有序列化注册表:
// ============================================================
// 改造后的客户端
// ============================================================
// ClientV2 是支持多序列化协议的 RPC 客户端
type ClientV2 struct {
conn net.Conn // TCP 连接
serializer Serializer // 默认使用的序列化器
registry *SerializerRegistry // 序列化协议注册表
messageID uint64 // 消息ID计数器(用于生成唯一ID)
mu sync.Mutex // 保护 messageID 的并发访问
}
// NewClientV2 创建新的客户端
// 参数 serializer 指定默认使用的序列化协议
func NewClientV2(addr string, serializer Serializer) (*ClientV2, error) {
conn, err := net.DialTimeout("tcp", addr, 5*time.Second)
if err != nil {
return nil, fmt.Errorf("连接服务器失败: %w", err)
}
// 创建序列化注册表并注册默认序列化器
registry := NewSerializerRegistry()
registry.Register(serializer)
return &ClientV2{
conn: conn,
serializer: serializer,
registry: registry,
}, nil
}
// CallV2 发起 RPC 调用(支持多序列化协议)
// 参数:
// - ctx: 上下文,用于超时控制
// - service: 服务名
// - method: 方法名
// - args: 请求参数
func (c *ClientV2) CallV2(ctx context.Context, service, method string, args interface{}) (*Response, error) {
// 第一步:使用序列化器将参数编码为字节流
body, err := c.serializer.Encode(args)
if err != nil {
return nil, fmt.Errorf("序列化请求参数失败: %w", err)
}
// 第二步:构造请求
// 生成唯一的消息ID
c.mu.Lock()
c.messageID++
msgID := c.messageID
c.mu.Unlock()
// 从 context 中提取需要传递的元数据
meta := extractMetaFromContext(ctx)
req := &Request{
Header: RequestHeader{
Version: 1, // 协议版本
Serializer: c.serializer.Code(), // 序列化协议标识
Compressor: 0, // 暂不压缩
MessageID: msgID, // 消息ID
Service: service, // 服务名
Method: method, // 方法名
Meta: meta, // 元数据
},
Body: body, // 序列化后的请求参数
}
// 第三步:发送请求
if err := WriteRequest(c.conn, req); err != nil {
return nil, fmt.Errorf("发送请求失败: %w", err)
}
// 第四步:读取响应
// 使用 select + channel 实现超时控制
respCh := make(chan *Response, 1)
errCh := make(chan error, 1)
go func() {
resp, err := ReadResponse(c.conn)
if err != nil {
errCh <- err
return
}
respCh <- resp
}()
// 等待响应或超时
select {
case resp := <-respCh:
// 检查响应头中的错误信息
if resp.Header.Error != "" {
return resp, fmt.Errorf("服务端错误: %s", resp.Header.Error)
}
return resp, nil
case err := <-errCh:
return nil, fmt.Errorf("读取响应失败: %w", err)
case <-ctx.Done():
// 超时或被取消
return nil, ctx.Err()
}
}
// extractMetaFromContext 从 context 中提取需要传递的元数据
// 例如:trace ID、超时时间戳、单向调用标记等
func extractMetaFromContext(ctx context.Context) map[string]string {
meta := make(map[string]string)
// 提取超时时间戳(如果有)
if deadline, ok := ctx.Deadline(); ok {
// 将截止时间转为 Unix 毫秒时间戳传递
meta["deadline"] = strconv.FormatInt(deadline.UnixMilli(), 10)
}
// 提取单向调用标记(如果有)
if v := ctx.Value("one-way"); v != nil {
meta["one-way"] = "true"
}
return meta
}
// ============================================================
// 改造后的服务端
// ============================================================
// ServerV2 是支持多序列化协议的 RPC 服务端
type ServerV2 struct {
addr string // 监听地址
mu sync.RWMutex // 保护 serviceMap
serviceMap map[string]reflect.Value // 服务注册表
registry *SerializerRegistry // 序列化协议注册表
}
// NewServerV2 创建新的服务端
func NewServerV2(addr string) *ServerV2 {
registry := NewSerializerRegistry()
// 注册默认的序列化协议
registry.Register(&JSONSerializer{})
registry.Register(&ProtoSerializer{})
return &ServerV2{
addr: addr,
serviceMap: make(map[string]reflect.Value),
registry: registry,
}
}
// handleConnV2 处理连接(支持多序列化协议)
func (s *ServerV2) handleConnV2(conn net.Conn) {
defer conn.Close()
for {
// 读取请求
req, err := ReadRequest(conn)
if err != nil {
if errors.Is(err, io.EOF) {
return
}
return
}
// 根据请求头中的序列化协议标识,找到对应的序列化器
serializer, err := s.registry.Get(req.Header.Serializer)
if err != nil {
// 找不到序列化器,返回错误响应
s.sendError(conn, req.Header.MessageID, err.Error())
continue
}
// 根据服务名查找服务
s.mu.RLock()
svcValue, ok := s.serviceMap[req.Header.Service]
s.mu.RUnlock()
if !ok {
s.sendError(conn, req.Header.MessageID, fmt.Sprintf("服务 %s 不存在", req.Header.Service))
continue
}
// 根据方法名查找方法
svcType := svcValue.Type()
method, ok := svcType.MethodByName(req.Header.Method)
if !ok {
s.sendError(conn, req.Header.MessageID, fmt.Sprintf("方法 %s 不存在", req.Header.Method))
continue
}
// 使用序列化器将 Body 反序列化为方法参数
// 约定:方法第二个参数是请求结构体指针
argType := method.Type.In(2) // In(0)=接收者, In(1)=context, In(2)=请求参数
argValue := reflect.New(argType.Elem()) // 创建请求结构体的指针
if err := serializer.Decode(argValue.Interface(), req.Body); err != nil {
s.sendError(conn, req.Header.MessageID, fmt.Sprintf("反序列化参数失败: %v", err))
continue
}
// 根据元数据重建 context
// 如果元数据中包含 deadline,则设置超时
ctx := context.Background()
ctx = rebuildContext(ctx, req.Header.Meta)
// 反射调用方法
in := []reflect.Value{
svcValue, // 接收者
reflect.ValueOf(ctx), // context.Context
argValue, // 请求参数
}
out := method.Func.Call(in)
// 处理返回值
var resp Response
if len(out) >= 2 && !out[1].IsNil() {
// 有错误返回
errMsg := out[1].Interface().(error).Error()
resp.Header = ResponseHeader{
Version: 1,
Serializer: req.Header.Serializer,
MessageID: req.Header.MessageID,
Error: errMsg,
}
} else {
// 成功,序列化响应数据
respData, err := serializer.Encode(out[0].Interface())
if err != nil {
s.sendError(conn, req.Header.MessageID, fmt.Sprintf("序列化响应失败: %v", err))
continue
}
resp.Header = ResponseHeader{
Version: 1,
Serializer: req.Header.Serializer,
MessageID: req.Header.MessageID,
}
resp.Body = respData
}
// 写回响应
if err := WriteResponse(conn, &resp); err != nil {
return
}
}
}
// sendError 发送错误响应
func (s *ServerV2) sendError(conn net.Conn, msgID uint64, errMsg string) {
resp := &Response{
Header: ResponseHeader{
Version: 1,
MessageID: msgID,
Error: errMsg,
},
}
WriteResponse(conn, resp)
}
// rebuildContext 根据元数据重建 context
// 如果元数据中包含 deadline,则创建带超时的 context
func rebuildContext(ctx context.Context, meta map[string]string) context.Context {
if deadlineStr, ok := meta["deadline"]; ok {
// 解析时间戳
deadlineMs, err := strconv.ParseInt(deadlineStr, 10, 64)
if err == nil {
deadline := time.UnixMilli(deadlineMs)
// 用剩余时间创建带超时的 context
// 如果 deadline 已经过去,timeout 会是负数,context 会立即超时
timeout := time.Until(deadline)
if timeout > 0 {
ctx, _ = context.WithDeadline(ctx, deadline)
}
}
}
return ctx
}
4.6 初始化与使用
func main() {
// --- 服务端初始化 ---
server := NewServerV2(":9091")
// 注册默认序列化协议(JSON 和 Proto 已在 NewServerV2 中注册)
// 注册服务
server.Register(&UserService{})
go server.Start()
// --- 客户端初始化 ---
// 创建 JSON 序列化器
jsonSerializer := &JSONSerializer{}
// 创建客户端,指定使用 JSON 序列化
client, err := NewClientV2("127.0.0.1:9091", jsonSerializer)
if err != nil {
fmt.Printf("连接失败: %v\n", err)
return
}
// 发起调用
ctx := context.Background()
resp, err := client.CallV2(ctx, "UserService", "GetById", &GetByIdRequest{Id: 123})
if err != nil {
fmt.Printf("调用失败: %v\n", err)
return
}
// 反序列化响应
var user GetByIdResponse
json.Unmarshal(resp.Body, &user)
fmt.Printf("查询结果: %+v\n", user)
}
flowchart LR
subgraph 客户端
C1[选择序列化器
JSON或Proto] --> C2[序列化参数为Body]
C2 --> C3[组装Request
Header+Body]
C3 --> C4[发送到服务端]
end
C4 -->|网络| S1
subgraph 服务端
S1[读取Request] --> S2[根据Header.Serializer
查找序列化器]
S2 --> S3[反序列化Body为参数]
S3 --> S4[反射执行方法]
S4 --> S5[序列化响应为Body]
S5 --> S6[写回Response]
end
S6 -->|网络| C5[读取Response]
C5 --> C6[反序列化Body为结果]五、RPC 调用语义
5.1 三种调用语义
RPC 调用语义是指客户端发起调用后的行为模式,主要有三种:
| 调用语义 | 说明 | Go 中的实现 |
|---|---|---|
| 异步调用 | 发起调用后可以做别的事,一段时间后再处理结果 | 用户自己开 goroutine 处理 |
| 单向调用(one-way) | 发起调用后不需要结果 | 通过元数据标记,服务端不回写响应 |
| 回调 | 发起调用时注册回调函数,结果返回后执行回调 | 用户自己开 goroutine 处理 |
Go 语言的特殊性:Go 里面的框架一般不会提供异步调用,因为用户可以自己开协程处理。类似的,回调在 Go 里面也不需要支持,因为用户同样可以开协程处理。
在 Java 之类的语言里面,异步调用是通过 Future 来达成的。即最开始调用的时候,用户只能拿到一个 Future 对象,等用户调用 Future 对象上的 Get 方法时就会被阻塞,直到拿到响应。
5.2 单向调用(One-way)
单向调用分成真假两种:
- 虚假的单向调用:是指用户的请求发过去之后,服务端会把响应发回来,但是客户端收到响应之后会直接丢弃
- 真实的单向调用:是指用户的请求发过去之后,服务端发现这是一个单向调用,直接就不会发回响应,读完请求之后,连接等资源就被释放了
核心在于:单向调用一定是为了尽快释放两端资源。因此虚假的单向调用,毫无意义。
在 Go 里面,虚假的单向调用,用户完全可以自己开一个 goroutine,只不过在 goroutine 里面用户不会处理响应而已。
flowchart TB
subgraph 虚假单向调用
F1[客户端发送请求] --> F2[服务端处理]
F2 --> F3[服务端发回响应]
F3 --> F4[客户端丢弃响应]
F4 --> F5[浪费了网络和服务端资源]
end
subgraph 真实单向调用
T1[客户端发送请求] --> T2[服务端处理]
T2 --> T3[服务端不发响应
立即释放连接资源]
T3 --> T4[客户端不等待响应
立即释放资源]
end5.3 单向调用的协议支持
为了支持真实的单向调用,我们需要显式地告诉服务端,这是一个单向调用。根据我们的协议设计,这个信息通过元数据传递:
// ============================================================
// 单向调用支持
// ============================================================
// WithOneWay 在 context 中设置单向调用标记
// 用户通过这个函数创建带标记的 context,传给 RPC 客户端
// 客户端检测到标记后,不会等待响应
// 服务端检测到标记后,不会发回响应
func WithOneWay(ctx context.Context) context.Context {
// 使用 context.WithValue 在 context 中携带单向调用标记
// key 使用自定义类型,避免与其他包的 key 冲突
return context.WithValue(ctx, oneWayKey{}, "true")
}
// oneWayKey 是单向调用标记的 context key 类型
// 使用自定义类型而不是字符串,避免 key 冲突
type oneWayKey struct{}
// isOneWay 检查 context 是否标记为单向调用
func isOneWay(ctx context.Context) bool {
v := ctx.Value(oneWayKey{})
return v != nil
}
5.4 客户端单向调用实现
// CallOneWay 发起单向调用
// 与普通调用的区别:
// 1. 在元数据中带上 one-way 标记
// 2. 发送请求后不等待响应
func (c *ClientV2) CallOneWay(ctx context.Context, service, method string, args interface{}) error {
// 在 context 中设置单向调用标记
ctx = WithOneWay(ctx)
// 序列化请求参数
body, err := c.serializer.Encode(args)
if err != nil {
return fmt.Errorf("序列化请求参数失败: %w", err)
}
// 生成消息ID
c.mu.Lock()
c.messageID++
msgID := c.messageID
c.mu.Unlock()
// 提取元数据(包含 one-way 标记)
meta := extractMetaFromContext(ctx)
// 构造请求
req := &Request{
Header: RequestHeader{
Version: 1,
Serializer: c.serializer.Code(),
MessageID: msgID,
Service: service,
Method: method,
Meta: meta,
},
Body: body,
}
// 发送请求后直接返回,不等待响应
// 这样客户端资源(goroutine、内存等)可以立即释放
return WriteRequest(c.conn, req)
}
5.5 服务端单向调用处理
// handleConnV3 处理连接(支持单向调用)
func (s *ServerV2) handleConnV3(conn net.Conn) {
defer conn.Close()
for {
req, err := ReadRequest(conn)
if err != nil {
if errors.Is(err, io.EOF) {
return
}
return
}
// 检查是否为单向调用
isOneWayCall := req.Header.Meta["one-way"] == "true"
// 处理请求(查找服务、反序列化、反射调用等逻辑与之前相同)
resp := s.processRequest(req)
// 如果是单向调用,不发回响应,直接处理下一个请求
// 这样服务端的连接和 goroutine 资源可以更快释放
if isOneWayCall {
continue
}
// 非单向调用,写回响应
if err := WriteResponse(conn, resp); err != nil {
return
}
}
}
5.6 单向调用使用场景
单向调用一般用在你不需要返回值的地方,即便失败了也无所谓的地方。例如:
- 日志上报:服务端只需要记录日志,不需要返回结果
- 消息推送:发送通知后不需要确认
- 监控指标上报:上报数据后不需要响应
- 缓存更新:异步更新缓存,不需要等待结果
六、RPC 超时控制
6.1 两种超时控制
在微服务里面,超时控制分为两种:
- 单一服务超时控制:每个服务独立设置自己的超时时间
- 链路超时控制:整个调用链共享一个超时时间
例如 A -> B -> C 这种调用关系,假如说 A -> B 设置的超时时间是 1s:
- 如果是链路超时控制,那么
B -> C这个调用,超时时间是不允许超过 1s 的 - 如果是单一服务超时控制,那么
B -> C这个调用可以设置超过 1s 的超时控制
绝大多数微服务框架的超时控制,都只支持单一服务超时控制。但链路超时控制对于防止雪崩效应非常重要。
6.2 链路超时控制原理
flowchart TB
subgraph "A -> B 调用 (总超时 3s)"
A1[Web/BFF 设置链路超时
ctx deadline=3s] --> A2[RPC客户端监听ctx超时]
A2 --> A3[通过meta传递剩余超时
假设还剩2.5s]
A3 -->|网络| B1
end
subgraph "B -> C 调用"
B1[RPC服务端收到请求
用剩余超时重建ctx] --> B2[业务代码使用ctx]
B2 --> B3[RPC客户端继续监听ctx
传递剩余超时给C]
B3 -->|网络| C1[C服务端重建ctx]
end链路超时控制的核心流程:
- 在起始位置(Web 或 BFF),设置好整个链路的超时时间
ctx会传递到 RPC 客户端,RPC 客户端会监听ctx超时- RPC 客户端把剩余的超时时间通过
meta传递给 RPC 服务端 - RPC 服务端收到请求后,用剩余超时时间重建整个
ctx - RPC 服务端把
ctx传递给业务代码 - 如果业务代码脱离了 RPC 框架(例如调用 HTTP 接口),用户需要手动管理链路超时时间
- 业务再一次发起 RPC 调用时,
ctx会再被传到 RPC 客户端,继续传递剩余超时时间
6.3 链路超时控制总结
- RPC 客户端需要监听
ctx,并且要把剩余超时时间传递给 RPC 服务端 - RPC 服务端收到请求之后,要检测元数据里面有没有携带剩余超时时间,然后重建
context.Context - 如果用户的业务里面用到了其它中间件,那么用户可能需要自己手动管理超时元数据,继续传递给其它中间件
- 任何一个环节的超时时间,都不可能超过链路超时时间
6.4 超时监听的位置
理论上来说,超时监听可以:
- 只在客户端监听
- 客户端和服务端同时监听
- 但不能只在服务端监听(如果服务端发现超时并返回超时响应时网络故障,客户端收不到响应就会一直等)
flowchart LR
P1["①获取连接前检测"] --> P2["②发送请求前检测"]
P2 --> P3["③等待响应时检测
(核心位置)"]
P3 --> P4["④读取响应后检测"]理论上可以在多个位置检测超时:
- 位置1:从连接池拿连接之前就先做一个检测
- 位置2:拿到连接之后,在发送请求之前做检测,如果超时就不需要再发送请求了
- 位置3:等待响应的过程中超时,那么就不需要读取响应了
- 位置4:读取响应之后超时了,那么就没必要解析结果了
实际中并不会做得那么复杂,而是只在位置1的时候检测一下,然后就发起调用,同时监听超时。
6.5 客户端超时监听实现
使用 select - case + channel 来实现超时监听,这是 Go 中最经典的超时控制模式:
// CallWithTimeout 带超时控制的 RPC 调用
// 使用 select + channel 同时监听响应和超时
func (c *ClientV2) CallWithTimeout(ctx context.Context, service, method string, args interface{}) (*Response, error) {
// 序列化请求参数
body, err := c.serializer.Encode(args)
if err != nil {
return nil, fmt.Errorf("序列化失败: %w", err)
}
// 生成消息ID
c.mu.Lock()
c.messageID++
msgID := c.messageID
c.mu.Unlock()
// 提取元数据(包含超时时间戳)
meta := extractMetaFromContext(ctx)
// 构造请求
req := &Request{
Header: RequestHeader{
Version: 1,
Serializer: c.serializer.Code(),
MessageID: msgID,
Service: service,
Method: method,
Meta: meta,
},
Body: body,
}
// 发送请求
if err := WriteRequest(c.conn, req); err != nil {
return nil, fmt.Errorf("发送请求失败: %w", err)
}
// 使用 select + channel 监听响应和超时
// 这是 Go 中最经典的超时控制写法
respCh := make(chan *Response, 1) // 带缓冲的 channel,防止 goroutine 泄漏
errCh := make(chan error, 1)
// 启动 goroutine 读取响应
go func() {
resp, err := ReadResponse(c.conn)
if err != nil {
errCh <- err
return
}
respCh <- resp
}()
// select 会同时监听多个 channel
// 哪个 channel 先有数据,就执行哪个 case
select {
case resp := <-respCh:
// 正常收到响应
if resp.Header.Error != "" {
return resp, fmt.Errorf("服务端错误: %s", resp.Header.Error)
}
return resp, nil
case err := <-errCh:
// 读取响应出错
return nil, fmt.Errorf("读取响应失败: %w", err)
case <-ctx.Done():
// context 超时或被取消
// 注意:超时之后并没有中断正在执行的调用
// 只是 RPC 客户端丢掉了后面的响应而已
// 每次都要新创建一个 channel,性能损耗比较大
// 但这种实现很简单也很清晰
return nil, ctx.Err()
}
}
这种实现的缺点:
- 超时之后并没有中断正在执行的调用,只是 RPC 客户端丢掉了后面的响应而已
- 每次都要新创建一个 channel,性能损耗比较大
但对于绝大多数应用来说,这种简单清晰的实现已经足够了。
6.6 跨端传递超时时间
现在完成了本地的超时控制,我们需要将链路超时时间传递给服务端。类似于前面传递的单向调用,也是在元数据部分,带一个 key-value 过去。
那么问题来了,超时时间传什么内容呢?有两种方案:
方案一:传递剩余超时时间
传递类似于 1s、2s、1500ms 这种数值,一般传递一个单位为毫秒的数值。
缺点:难以估计网络传输时间。
例如在 RPC 客户端计算剩余超时时间的时候还有 2000ms,发送到 RPC 服务端的时候,RPC 服务端需要考虑扣减网络传输时间,比如说 10ms,那么 RPC 服务端实际剩下的超时时间就只有 1990ms 了。难就难在,怎么知道花了 10ms?
不过,网络传输时间极短,相比超时时间设置,基本可以忽略不计,比如说 10ms 相比 2000ms 实在是微不足道。可以通过测试来判断在平均大小请求的条件下,两个节点之间的传输时间,测试要完全模拟线上环境。
方案二:传递超时时间戳
传递的是过期的时间戳,例如在 1640970000000(2022-01-01 01:00:00.000)过期。
缺点:
- 时钟同步问题:即某一个特定的时刻,时钟读数在两台机器上是不同的
- 时间戳一般传递 Unix 时间戳,至少需要 64 位。在我们的设计里面传递的是字符串,需要更多的字节
这里我们采用时间戳的方案。两个方案之间,没有特别大的优劣之分,看喜好选择就可以。
6.7 超时时间传递实现
// ============================================================
// 超时时间跨端传递实现
// ============================================================
// 设置超时时间戳到元数据
// 在前面的 extractMetaFromContext 函数中已经实现:
// 如果 context 有 deadline,会将 deadline 转为 Unix 毫秒时间戳
// 存入 meta["deadline"]
//
// 服务端在 rebuildContext 函数中会读取这个时间戳,
// 用它创建新的带超时的 context
// 客户端设置 meta 的过程(已在 extractMetaFromContext 中实现):
func extractMetaFromContext(ctx context.Context) map[string]string {
meta := make(map[string]string)
// 如果 context 设置了 deadline,将其转为时间戳传递
if deadline, ok := ctx.Deadline(); ok {
// 将截止时间转为 Unix 毫秒时间戳
// 使用时间戳方案,服务端可以直接用这个时间戳重建 context
meta["deadline"] = strconv.FormatInt(deadline.UnixMilli(), 10)
}
// 提取单向调用标记
if v := ctx.Value(oneWayKey{}); v != nil {
meta["one-way"] = "true"
}
return meta
}
// 服务端重建 context 的过程(已在 rebuildContext 中实现):
func rebuildContext(ctx context.Context, meta map[string]string) context.Context {
if deadlineStr, ok := meta["deadline"]; ok {
// 解析时间戳
deadlineMs, err := strconv.ParseInt(deadlineStr, 10, 64)
if err == nil {
deadline := time.UnixMilli(deadlineMs)
// 用这个 deadline 创建新的 context
// 如果 deadline 已经过去,context 会立即超时
ctx, _ = context.WithDeadline(ctx, deadline)
}
}
return ctx
}
flowchart TB
subgraph 客户端
C1["创建带超时的ctx
context.WithDeadline"] --> C2["提取deadline
转为时间戳"]
C2 --> C3["放入meta('deadline')"]
C3 --> C4[发送请求]
end
C4 -->|网络| S1
subgraph 服务端
S1[收到请求] --> S2["读取meta('deadline')"]
S2 --> S3[解析时间戳]
S3 --> S4["用deadline重建ctx
context.WithDeadline"]
S4 --> S5[业务代码使用新ctx]
end6.8 完整超时控制使用示例
func main() {
// 创建客户端
client, _ := NewClientV2("127.0.0.1:9091", &JSONSerializer{})
defer client.Close()
// 创建带 3 秒超时的 context
// 这 3 秒是整个调用链的超时时间
// 如果服务端再调用其他服务,剩余超时时间会继续传递
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel() // 确保资源释放
// 发起调用,传入带超时的 ctx
resp, err := client.CallWithTimeout(ctx, "UserService", "GetById", &GetByIdRequest{Id: 123})
if err != nil {
if errors.Is(err, context.DeadlineExceeded) {
fmt.Println("调用超时")
return
}
fmt.Printf("调用失败: %v\n", err)
return
}
// 处理响应
var user GetByIdResponse
json.Unmarshal(resp.Body, &user)
fmt.Printf("查询结果: %+v\n", user)
}
七、面试要点总结
7.1 RPC 协议设计
- RPC 协议主要包含什么:消息头和消息体
- 头部包含什么?有什么用:回答几个主要的头部字段——长度字段、版本字段、序列化协议、压缩算法、消息 ID、服务名、方法名、错误(响应头)
- RPC 里面为什么包含长度字段:主要是为了切割消息,也就所谓的粘包问题
- 头部可以是变长的吗:当然可以
- 大概描述一下 gRPC 协议:关键点是 gRPC 利用了 HTTP 协议,然后再描述一下 gRPC 的几个常见头部。问到 Dubbo 协议也是类似
- 为什么尽量把元数据之类的东西放在消息头:主要是考虑到 sidecar、网关之类的东西,这样可以做到解析部分数据,提高性能
- 如何在 RPC 协议里支持不同序列化协议/压缩算法:但凡问到如何在 RPC 协议上支持 XXX 功能,核心都是要在协议本身加上对应的字段,要考虑是放到协议体还是协议头,然后客户端和服务端都要做相应的修改
7.2 调用语义
- 什么是异步调用?Go 里面怎么实现:核心还是那句话,在有 goroutine 的情况下,框架设计者没有必要支持
- 什么是回调?Go 里面怎么实现:同上
- 什么是 one-way(单向)调用:要注意详细解释真伪两种 one-way 的形态,最关键的地方在于,服务端究竟会不会回写响应,而后进一步讨论两种形态对性能的影响
- one-way 用在什么场景:一般就是用在你不需要返回值的地方,即便失败了也无所谓的地方
- 使用 one-way 有什么优点:性能,能尽早释放资源
7.3 超时控制
- RPC 怎么控制超时:客户端控制、服务端控制的特点
- 为什么要使用超时控制:及时释放资源
- 超时时间没有设置好会有什么问题:过短——大部分请求超时;过长——浪费资源,甚至引起 goroutine 泄漏
- 超时时间在链路中传递,传递的是什么:剩余超时时间或者超时时间戳。进一步可以讨论超时时间戳与时钟同步的问题
- 超时之后可以中断业务执行吗:不可以
- 链路超时怎么实现:核心就是在链路中传递超时时间,这个主要依赖于在 RPC 协议里面传递超时时间。注意,如果你是一个中间件设计者,还要考虑用户可能希望重置整个链路的超时时间,那么你要设计类似的接口
总结
本教程从最简 RPC 的局限性出发,讲解了如何设计一个完整的 RPC 协议,并逐步实现了多序列化协议支持、单向调用和链路超时控制。
关键知识点回顾:
- 协议设计的基本原则是分成协议头和协议体,头部放路由和元数据,体放业务数据
- 服务名/方法名放头部是为了让 sidecar、网关等中间件只解析头部就能做路由
- 多序列化协议通过 Serializer 接口 + 注册表模式实现,协议头中用 1 字节标记序列化方式
- 单向调用分为真假两种,真实的单向调用通过元数据标记实现,服务端不发响应
- 超时控制使用
select + channel监听context.Done(),链路超时通过元数据传递时间戳 - 跨端传递超时有剩余时间和时间戳两种方案,各有优劣
下一章我们将讲解服务注册与发现,探讨微服务框架中服务之间如何互相找到对方。
自测题与动手练习
自测题(合上书能答出来,才算懂):
- RPC 协议为什么要分成协议头和协议体?协议头里放"服务名/方法名"而不是放协议体,核心价值是什么?
- 我们给请求加了一个"长度字段",它主要为了解决什么问题?如果客户端连续发了两条消息,没有长度字段会怎样?
- 多序列化协议是怎么协商出来的?客户端和服务端各需要做哪些事,才能让"同一个接口、不同序列化方式"跑通?
- “真实单向调用"和"虚假单向调用"的区别在哪?为什么说虚假的单向调用"毫无意义”?
- 链路超时控制里,跨端传递的是"剩余时间"还是"时间戳"?两种方案各自的坑是什么?超时之后,正在执行的业务逻辑会被中断吗?
动手练习(建议真做一遍):
- 在
EncodeRequest基础上加一个 gzip 压缩分支:当Compressor == 1时,对Body先gzip再写长度。反序列化端对称解压,打印前后字节数对比。 - 把"虚假单向调用"改成"真实单向调用":服务端检测到
meta["one-way"]=="true"时直接continue,用 Wireshark 或加日志验证服务端确实没回包。 - 把
CallWithTimeout的select改成"可中断"版本:用context.WithCancel在收到响应后主动cancel,观察 goroutine 是否立刻退出(用runtime.NumGoroutine()打印验证无泄漏)。
本章小结
- 协议设计的本质是把"信封"和"信纸"分开:协议头放路由与元数据(长度、版本、序列化、压缩、消息 ID、服务名、方法名、错误),协议体放业务数据。
- 元数据和路由信息放头部是为了让网关、sidecar 只解析头部就能做负载均衡和路由,不用反序列化整个请求体。
- 多序列化协议靠
Serializer接口 + 注册表实现,协议头用 1 字节Serializer标记协商,客户端按标记编码、服务端按标记查表解码。 - 单向调用分真假两种,真实单向调用通过
meta["one-way"]标记实现,服务端不回包,两端尽快释放资源;它适合日志上报、监控打点等"发完就忘"的场景。 - 超时控制用
select + channel监听context.Done();链路超时通过meta["deadline"]时间戳跨端传递,服务端据此重建带超时的ctx,任何环节超时都不可能超过链路总时间。
下一章我们进入服务注册与发现,看看微服务之间怎么互相"找到"对方。