二、RPC 协议设计与实现

2021-02-17T14:21:02+08:00 | 30分钟阅读 | 更新于 2021-02-17T14:21:02+08:00

@

学习目标

学完本章,你应该能够:

  1. 说清楚一个 RPC 协议为什么要分"头"和"体":能用自己的话解释协议头里放路由/元数据、协议体放业务数据的取舍,以及为什么中间件(网关、sidecar)只解析头部就能做路由。
  2. 设计一个最小可用的请求/响应协议:讲得清头部要有哪些字段(长度、版本、序列化标记、压缩标记、消息 ID、服务名、方法名、错误),以及它们分别解决什么问题(粘包、升级、多路复用、版本兼容)。
  3. 实现多序列化协议支持:讲清 Serializer 接口 + 注册表模式如何工作,知道客户端和服务端如何通过协议头里的 1 字节标记协商序列化方式。
  4. 分辨三种调用语义并落地单向调用:能讲清异步/回调/单向调用的区别,能解释"真假单向调用"对资源释放的影响,并看懂 WithOneWay 的实现。
  5. 讲清链路超时控制的原理:面试时能把这个知识点讲成一个连贯的故事——单一超时 vs 链路超时、用 select + channel 监听 context.Done()、为什么用"时间戳"而不是"剩余时间"跨端传递。

前置知识

  • 上一章"最简 RPC 框架"(JSON 传输、4 字节长度前缀、基本的客户端/服务端)。
  • Go 基础:netencoding/json、goroutine 与 channelcontext 包的基本用法。
  • 一点网络常识:TCP 是字节流、有"粘包"问题,HTTP/2、gRPC、Dubbo 名词听说过即可。

本章你会动手做的事

  1. 在纸上画出"请求协议头"的字段布局,标注哪些是定长、哪些是不定长。
  2. EncodeRequest / DecodeRequest 的编码格式 [头长度][头][体长度][体] 在草稿上模拟一遍(给一个假头部,写出最终字节流)。
  3. CallWithTimeout 跑一次带 3 秒超时的调用,故意把服务端 sleep 5 秒,观察客户端是否按时返回 context deadline exceeded

前言:从最简 RPC 到真正的 RPC 协议

在上一章中,我们手写了一个"最简 RPC 框架"——它能用 JSON 格式传递调用信息,实现基本的远程调用。但这个框架有一个明显的局限:消息格式被写死成 JSON,不支持其他序列化协议。

而在实际应用中,一个成熟的 RPC 框架还需要考虑:

  • 支持不同的序列化协议(JSON、Protobuf、Gob 等)
  • 支持数据压缩(gzip、snappy、zstd 等)
  • 支持协议版本升级(未来字段变更时不破坏旧客户端)
  • 支持加密传输
  • 支持链路追踪(trace ID 传递)
  • 支持超时控制单向调用等高级语义

这一切都引出一个核心问题:如何设计一个合适的 RPC 协议?

本教程将带你从协议设计理论出发,逐步实现一个支持多序列化协议、多调用语义和链路超时控制的 RPC 框架。

本教程涵盖以下主题:

  1. 协议设计理论——协议头与协议体,gRPC/Dubbo/TCP 协议对比
  2. 请求与响应协议——字段设计、接口定义、编解码实现
  3. 多序列化协议支持——Serializer 接口与 JSON/Proto 实现
  4. RPC 调用语义——异步调用、回调、单向调用
  5. RPC 超时控制——单一超时与链路超时、跨端传递
  6. 面试要点总结

一、协议设计理论

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:#e1f5fe

2.3 响应协议头设计

响应头部也设计为不定长,包含以下字段:

字段说明用途
长度字段协议头长度 + 协议体长度用于分割消息
版本字段协议版本号用于后续协议升级
序列化协议序列化方式标记用于标记采用的序列化协议
压缩算法压缩方式标记用于标记协议体是如何被压缩的
消息 ID唯一标识用于匹配请求和响应
错误错误信息为了解决第二个返回值的问题

最后的协议体里面就只存放响应数据。

2.4 响应错误为什么放头部

主要是实在没地方放。从理论上来说,有两个选择:

  1. 放头部:做成一个类似于 serviceName 那种,认为是服务调用本身必需的一种数据
  2. 放协议体:也就是认为是响应体的一部分。这种做法会更加符合直觉,但实现起来难度更高,因为你需要区别响应数据里面,哪部分是用户返回的数据,哪部分是用户返回的错误

我们选择放头部,这样协议体可以干净地只包含业务数据,编解码逻辑更简单。


三、接口定义与编解码实现

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[客户端不等待响应
立即释放资源] end

5.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

链路超时控制的核心流程:

  1. 在起始位置(Web 或 BFF),设置好整个链路的超时时间
  2. ctx 会传递到 RPC 客户端,RPC 客户端会监听 ctx 超时
  3. RPC 客户端把剩余的超时时间通过 meta 传递给 RPC 服务端
  4. RPC 服务端收到请求后,用剩余超时时间重建整个 ctx
  5. RPC 服务端把 ctx 传递给业务代码
  6. 如果业务代码脱离了 RPC 框架(例如调用 HTTP 接口),用户需要手动管理链路超时时间
  7. 业务再一次发起 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] end

6.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 协议,并逐步实现了多序列化协议支持、单向调用和链路超时控制。

关键知识点回顾:

  1. 协议设计的基本原则是分成协议头和协议体,头部放路由和元数据,体放业务数据
  2. 服务名/方法名放头部是为了让 sidecar、网关等中间件只解析头部就能做路由
  3. 多序列化协议通过 Serializer 接口 + 注册表模式实现,协议头中用 1 字节标记序列化方式
  4. 单向调用分为真假两种,真实的单向调用通过元数据标记实现,服务端不发响应
  5. 超时控制使用 select + channel 监听 context.Done(),链路超时通过元数据传递时间戳
  6. 跨端传递超时有剩余时间和时间戳两种方案,各有优劣

下一章我们将讲解服务注册与发现,探讨微服务框架中服务之间如何互相找到对方。


自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. RPC 协议为什么要分成协议头和协议体?协议头里放"服务名/方法名"而不是放协议体,核心价值是什么?
  2. 我们给请求加了一个"长度字段",它主要为了解决什么问题?如果客户端连续发了两条消息,没有长度字段会怎样?
  3. 多序列化协议是怎么协商出来的?客户端和服务端各需要做哪些事,才能让"同一个接口、不同序列化方式"跑通?
  4. “真实单向调用"和"虚假单向调用"的区别在哪?为什么说虚假的单向调用"毫无意义”?
  5. 链路超时控制里,跨端传递的是"剩余时间"还是"时间戳"?两种方案各自的坑是什么?超时之后,正在执行的业务逻辑会被中断吗?

动手练习(建议真做一遍)

  1. EncodeRequest 基础上加一个 gzip 压缩分支:当 Compressor == 1 时,对 Bodygzip 再写长度。反序列化端对称解压,打印前后字节数对比。
  2. 把"虚假单向调用"改成"真实单向调用":服务端检测到 meta["one-way"]=="true" 时直接 continue,用 Wireshark 或加日志验证服务端确实没回包。
  3. CallWithTimeoutselect 改成"可中断"版本:用 context.WithCancel 在收到响应后主动 cancel,观察 goroutine 是否立刻退出(用 runtime.NumGoroutine() 打印验证无泄漏)。

本章小结

  • 协议设计的本质是把"信封"和"信纸"分开:协议头放路由与元数据(长度、版本、序列化、压缩、消息 ID、服务名、方法名、错误),协议体放业务数据。
  • 元数据和路由信息放头部是为了让网关、sidecar 只解析头部就能做负载均衡和路由,不用反序列化整个请求体。
  • 多序列化协议Serializer 接口 + 注册表实现,协议头用 1 字节 Serializer 标记协商,客户端按标记编码、服务端按标记查表解码。
  • 单向调用分真假两种,真实单向调用通过 meta["one-way"] 标记实现,服务端不回包,两端尽快释放资源;它适合日志上报、监控打点等"发完就忘"的场景。
  • 超时控制select + channel 监听 context.Done();链路超时通过 meta["deadline"] 时间戳跨端传递,服务端据此重建带超时的 ctx,任何环节超时都不可能超过链路总时间。

下一章我们进入服务注册与发现,看看微服务之间怎么互相"找到"对方。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!