一、网络编程与最简 RPC

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

@

学习目标

学完本章你应该能够:

  1. 说清微服务相比单体的核心动机(分而治之)与拆分后带来的四类新问题:服务发现、远程通信、容错、可观测。
  2. 用 Go 的 net 包写出一个最简单的 TCP 服务器与客户端,讲清 Listen / Accept / Dial 的分工,以及"读数据 → 处理 → 回写响应"的三段式。
  3. 解释连接池的四个核心参数(InitialCap / MaxIdle / MaxCap / MaxIdleTime)与 Get / Put 流程,并说出 goroutine 三种使用模式在吞吐与复杂度上的取舍。
  4. 讲清 RPC 的本质(像本地调用一样调远程方法)与必须解决的三个问题:调用信息是什么、客户端怎么捕捉、怎么编码传过去再还原。
  5. 手写最简 RPC 框架时,能讲清代理模式 + 反射 + 长度前缀协议各自解决了什么问题,面试能讲成一段闭环故事。

前置知识

  • Go 基础:goroutinechannelinterface、反射(reflect)入门。
  • 网络常识:TCP 三次握手、流式协议没有消息边界。
  • JSON 序列化基础。

本章你会动手做的事

  1. 跑通 2.3 的 TCP 服务器 + 2.7 的客户端,观察 Accept 后每个连接一个 goroutine 的并发行为。
  2. 用 4.5 的最小连接池把 MaxCap 调到 2,写个并发 10 的压测,观察"获取连接超时"。
  3. 在 6.4 的用户服务上加一个 SayHello 方法,走通代理调用,理解反射是怎么找到方法并执行的。

前言:从单体到微服务

在互联网发展的早期,一个系统通常是一个"单体应用"——所有功能模块打包在一起,部署在同一个进程中。单体应用虽然简单,但随着业务规模增长,它会变得越来越臃肿:改一个小 bug 需要重新部署整个系统,一个模块的内存泄漏会拖垮所有功能。

微服务架构的核心思想就是"分而治之":把一个大系统拆分成多个独立的小服务,每个服务负责一块业务,服务之间通过网络通信。这样每个服务可以独立开发、独立部署、独立扩展。

但拆分之后,新的问题随之而来:

  • 服务 A 怎么找到服务 B 的地址?(服务发现
  • 服务 A 怎么调用服务 B?(远程通信
  • 服务 B 挂了怎么办?(容错机制
  • 怎么监控服务的健康状态?(可观测性

这些问题正是微服务框架要解决的。本教程将从最底层的网络编程讲起,一步步带你理解微服务框架的核心——RPC(远程过程调用),并最终手写一个"最简 RPC 框架"。

本教程涵盖以下主题:

  1. 微服务框架概览——理解微服务框架要解决什么问题
  2. Go 网络编程基础——使用 net 包进行 TCP 通信
  3. goroutine 与连接处理——并发处理连接的三种模式
  4. 连接池——复用连接提升性能
  5. RPC 核心概念——远程过程调用的本质
  6. 最简 RPC 框架实现——从零手写一个 RPC 框架
  7. 面试要点总结

一、微服务框架概览

1.1 什么是微服务架构

微服务架构简单来说,就是指整个系统由多个组件组成,每一个组件都独立管理,组件之间通过网络来通信。

注意:单体应用可以部署多个实例,但是它依旧是单体应用。因为单体应用的不同实例之间不会有交互,它们只是同一份代码的多份拷贝。而微服务的不同服务之间是有交互的——一切问题都可以源于网络间通信。

1.2 微服务框架要解决的核心问题

微服务框架主要解决两个核心问题:通信 + 服务治理

  • 通信:即服务之间如何发起调用,一般是 RPC,或者是 HTTP 直接通信
  • 服务治理:涵盖从服务注册与发现到可观测性的全部内容(包括负载均衡、熔断限流、链路追踪、日志监控等)

1.3 微服务框架的分类

微服务框架可以进一步分成三类:

类型特点代表框架
纯粹的 RPC 框架只负责通信,不涉及服务治理早期 gRPC
服务治理框架不设计自己的通信协议,专注服务治理Kratos、go-zero
大一统的微服务框架既有通信协议,又有服务治理模块Dubbo

1.4 主流框架巡礼

gRPC 与 protobuf

gRPC 比较学院派,它是典型的使用 IDL(Interface Description Language,接口描述语言)来生成代码的 RPC 框架。IDL 是指用一种中间语言来定义接口,而后为其它语言生成对应代码的设计方案。所以 gRPC 是多语言通信的首选。

gRPC 使用的 IDL 是 protobuf。protobuf 是一个独立的 IDL,也就是说你可以用 protobuf 来生成 gRPC 的代码,也可以用 protobuf 来生成其它 RPC 框架的代码。protobuf 也定义了序列化格式,所以我们也常说使用 protobuf 来作为序列化协议。

实践建议:遇事不决用 gRPC。如果是小型系统,可以考虑直接使用 HTTP 接口。

Dubbo

Dubbo 是非常早就出现的微服务框架,一句话总结就是:全家桶。Dubbo 涵盖了从上层服务治理到底层通信协议设计的全方位内容。也因为 Dubbo 历经考验,所以基本上微服务相关的所有的话题,你都可以在 Dubbo 里面找到。它是设计微服务框架非常好的参考对象。

go-micro

go-micro 有自己的协议,它本质上也是利用了 protobuf 作为 IDL。同时它也支持了 gRPC 和 HTTP。go-micro 充分利用了插件机制,用户可以替换掉大部分实现,包括注册中心、负载均衡、底层协议。

Kratos

B 站开源出来的,毛老师作品。主要聚焦在服务治理和快速开发上,也就是说它兼具一个微服务框架的功能,以及一个脚手架的功能。它依托于 gRPC 和 HTTP 来作为底层通信协议。

go-zero

近两年很火的一个框架,跟 Kratos 很像,也同样是聚焦在上层服务治理和快速开发上。go-zero 和 Kratos 都是偏向业务的,里面集成了一些实践中本来应该是微服务框架使用者(而不是设计者)要考虑的内容。

1.5 如何选择微服务框架

  • 懒人方案:选 go-zero 或者 Kratos,集成了一些业务实践,开箱即用
  • 中等方案:选 go-micro、Dubbo,如果你担忧高并发大集群的问题,Dubbo 可能更可靠
  • 最高自由度:直接使用 gRPC 或者 HTTP 协议来通信,服务治理需要的时候自己搞

一句话总结:选 gRPC,剩下就随意了。


二、Go 网络编程基础

微服务之间的通信建立在网络编程之上。Go 语言的 net 包是网络相关的核心包,里面包含了 httprpc 等关键包。理解 net 包是理解微服务框架底层通信的基础。

2.1 通信基本流程

网络通信基本分成两个大阶段:

flowchart LR
    subgraph 创建连接阶段
        A[服务端监听端口] --> B[客户端拨通服务端]
        B --> C[TCP 三次握手协商连接]
    end
    C --> D[连接建立]
    subgraph 通信阶段
        D --> E[客户端发送请求]
        E --> F[服务端读取请求]
        F --> G[服务端处理请求]
        G --> H[服务端写回响应]
        H --> I[客户端读取响应]
    end

net 包里面,最重要的两个调用:

  • Listen(network, addr string):监听某个端口,等待客户端连接
  • Dial(network, addr string):拨号,连上某个服务端

2.2 net.Listen:服务端监听

Listen 是监听一个端口,准备读取数据。它还有几个类似接口:

  • ListenTCP:监听 TCP 连接
  • ListenUDP:监听 UDP 连接
  • ListenIP:监听 IP 连接
  • ListenUnix:监听 Unix 域套接字

这些方法都是返回 Listener 的具体类型,如 TCPListener。一般用 Listen 就可以,除非你需要依赖于具体的网络协议特性。

提示:网络通信用 TCP 还是 UDP 是一个影响巨大的事情,一般确认了就不会改。

2.3 创建 TCP 服务器:完整代码

下面我们从零开始写一个简单的 TCP 服务器。这个服务器接收客户端发来的消息,转换成大写后返回。

package main

import (
	"bufio"
	"fmt"
	"net"
)

// TCPServer 是一个简单的 TCP 服务器结构体
// 它只包含一个地址字段,用于指定服务器监听的地址和端口
type TCPServer struct {
	Addr string // 服务器监听地址,例如 ":8080" 表示监听所有网卡的 8080 端口
}

// Start 启动 TCP 服务器
// 这个方法会一直阻塞,直到服务器关闭
func (s *TCPServer) Start() error {
	// 第一步:使用 net.Listen 创建一个监听器
	// 参数1 "tcp" 表示使用 TCP 协议
	// 参数2 是监听地址,例如 ":8080"
	// 返回的 listener 是一个 net.Listener 接口,用于接受客户端连接
	listener, err := net.Listen("tcp", s.Addr)
	if err != nil {
		return fmt.Errorf("监听失败: %w", err) // %w 是 Go 1.13 引入的错误包装语法
	}
	// defer 确保在函数退出时关闭监听器,防止资源泄漏
	defer listener.Close()

	fmt.Printf("服务器启动成功,监听地址: %s\n", s.Addr)

	// 第二步:在一个 for 循环中不断接受新连接
	// listener.Accept() 会阻塞,直到有客户端连接进来
	// 每次有新连接,都会返回一个 net.Conn 对象,代表这个连接
	for {
		conn, err := listener.Accept()
		if err != nil {
			// 如果接受连接出错,打印错误但不要退出服务器
			// 因为可能只是某个连接出问题,服务器应该继续服务其他客户端
			fmt.Printf("接受连接失败: %v\n", err)
			continue
		}
		// 第三步:为每个连接启动一个独立的 goroutine 来处理
		// 这样服务器就可以同时处理多个客户端的连接(并发)
		// 如果不使用 goroutine,服务器只能一次处理一个客户端
		go s.handleConn(conn)
	}
}

// handleConn 处理单个客户端连接
// 基本流程:读数据 -> 处理数据 -> 回写响应
// 这个方法会在一个独立的 goroutine 中运行
func (s *TCPServer) handleConn(conn net.Conn) {
	// defer 确保连接在函数退出时被关闭
	// 无论函数是正常返回还是 panic,defer 都会执行
	defer conn.Close()

	// 远程地址信息,用于日志打印
	remoteAddr := conn.RemoteAddr().String()
	fmt.Printf("客户端 %s 已连接\n", remoteAddr)

	// 使用 bufio.NewReader 包装 conn,方便按行读取数据
	// bufio.Reader 会缓冲数据,减少系统调用次数,提高读取效率
	reader := bufio.NewReader(conn)

	for {
		// ReadString 读取直到遇到分隔符(这里是换行符 \n)
		// 返回的内容包含分隔符本身
		// 如果客户端关闭了连接,会返回 io.EOF 错误
		line, err := reader.ReadString('\n')
		if err != nil {
			// 如果遇到 EOF,说明客户端主动关闭了连接,这是正常情况
			// ErrUnexpectedEOF 表示读取到一半连接断了
			// 这两种情况都应该直接关闭连接
			fmt.Printf("客户端 %s 断开连接: %v\n", remoteAddr, err)
			return // 退出函数,defer 会关闭 conn
		}

		// 去掉末尾的换行符,得到纯净的消息内容
		msg := line[:len(line)-1]
		fmt.Printf("收到来自 %s 的消息: %s\n", remoteAddr, msg)

		// 处理数据:这里简单地把消息转成大写
		// 在真实场景中,这里可能是查数据库、调用其他服务等
		response := fmt.Sprintf("大写: %s\n", toUpper(msg))

		// 回写响应:把处理结果写回给客户端
		// conn.Write 是把字节切片写入连接,通过网络发送给客户端
		// 即便处理数据出错,也要返回一个错误给客户端
		// 不然客户端不知道服务端处理出错了
		_, err = conn.Write([]byte(response))
		if err != nil {
			fmt.Printf("写回响应失败: %v\n", err)
			return
		}
	}
}

// toUpper 将字符串转换为大写
// 这是一个辅助函数,模拟"处理数据"的逻辑
func toUpper(s string) string {
	result := make([]byte, len(s))
	for i := 0; i < len(s); i++ {
		// 如果是小写字母(a-z),转换成大写(A-Z)
		// 小写字母的 ASCII 码比大写字母大 32
		if s[i] >= 'a' && s[i] <= 'z' {
			result[i] = s[i] - 32
		} else {
			result[i] = s[i]
		}
	}
	return string(result)
}

func main() {
	// 创建一个 TCP 服务器,监听 8080 端口
	// ":8080" 表示监听所有网卡的 8080 端口
	// 如果只想监听本机,可以用 "127.0.0.1:8080"
	server := &TCPServer{Addr: ":8080"}

	// 启动服务器,这个调用会一直阻塞
	if err := server.Start(); err != nil {
		fmt.Printf("服务器启动失败: %v\n", err)
	}
}

2.4 处理连接的核心逻辑

处理连接基本上就是在一个 for 循环内重复三个步骤:

  1. 读数据:读数据要根据上层协议来决定怎么读。例如,简单的 RPC 协议一般分成两段读——先读头部,根据头部得知 Body 有多长,再把剩下的数据读出来。
  2. 处理数据:根据业务逻辑处理请求。
  3. 回写响应:即便处理数据出错,也要返回一个错误给客户端,不然客户端不知道你处理出错了。

2.5 错误处理

在读写的时候,都可能遇到错误。一般来说代表连接已经关掉的是这三个:

  • io.EOF:正常读到末尾,客户端正常关闭了连接
  • io.ErrUnexpectedEOF:读到一半连接断了
  • net.ErrClosed:连接已经被关闭了

实践建议:只要是出错了就直接关闭连接,这样对客户端和服务端代码都简单。不要试图从错误中恢复连接,因为连接可能已经处于不一致的状态。

2.6 net.Dial:客户端连接

net.Dial 是指创建一个连接,连上远端的服务器。它也有几个类似的方法:

  • DialIP
  • DialTCP
  • DialUDP
  • DialUnix
  • DialTimeout:多了一个超时参数

实践建议:直接使用 DialTimeout,因为设置超时可以避免一直阻塞。如果不设置超时,当服务端无响应时,客户端会永远卡在 Dial 调用上。

2.7 创建 TCP 客户端:完整代码

下面写一个与上面 TCP 服务器配套的客户端:

package main

import (
	"bufio"
	"fmt"
	"net"
	"os"
	"time"
)

// TCPClient 是一个简单的 TCP 客户端结构体
type TCPClient struct {
	Addr string // 服务器地址,例如 "127.0.0.1:8080"
}

// Connect 连接服务器并发送消息
func (c *TCPClient) Connect() error {
	// 第一步:使用 net.DialTimeout 连接服务器
	// 参数1 "tcp" 表示使用 TCP 协议
	// 参数2 是服务器地址
	// 参数3 是超时时间——超过这个时间还没连上就返回错误
	// 使用 DialTimeout 而不是 Dial,可以避免一直阻塞
	conn, err := net.DialTimeout("tcp", c.Addr, 3*time.Second)
	if err != nil {
		return fmt.Errorf("连接服务器失败: %w", err)
	}
	// 确保函数退出时关闭连接
	defer conn.Close()

	fmt.Printf("已连接到服务器 %s\n", c.Addr)

	// 使用 bufio 读写,方便处理文本数据
	reader := bufio.NewReader(conn)    // 用于读取服务器返回的数据
	stdinReader := bufio.NewReader(os.Stdin) // 用于读取用户从键盘输入的数据

	for {
		// 提示用户输入
		fmt.Print("请输入消息(输入 quit 退出): ")

		// 从标准输入(键盘)读取一行
		line, err := stdinReader.ReadString('\n')
		if err != nil {
			return fmt.Errorf("读取输入失败: %w", err)
		}

		// 去掉末尾换行符
		msg := line[:len(line)-1]

		// 如果用户输入 quit,退出循环
		if msg == "quit" {
			break
		}

		// 第二步:发送消息到服务器
		// 需要在消息末尾加上换行符,因为服务器用 ReadString('\n') 读取
		_, err = conn.Write([]byte(msg + "\n"))
		if err != nil {
			return fmt.Errorf("发送消息失败: %w", err)
		}

		// 第三步:读取服务器的响应
		// 服务器返回的数据也以换行符结尾
		response, err := reader.ReadString('\n')
		if err != nil {
			return fmt.Errorf("读取响应失败: %w", err)
		}

		fmt.Printf("服务器响应: %s", response)
	}

	return nil
}

func main() {
	// 创建客户端,连接到本机 8080 端口的服务器
	client := &TCPClient{Addr: "127.0.0.1:8080"}

	if err := client.Connect(); err != nil {
		fmt.Printf("客户端出错: %v\n", err)
	}
}

2.8 运行示例

打开两个终端:

# 终端1:启动服务器
go run server.go

# 终端2:启动客户端
go run client.go

客户端输出示例:

已连接到服务器 127.0.0.1:8080
请输入消息(输入 quit 退出): hello
服务器响应: 大写: HELLO
请输入消息(输入 quit 退出): world
服务器响应: 大写: WORLD
请输入消息(输入 quit 退出): quit

三、goroutine 与连接处理

3.1 三种 goroutine 使用模式

在前面的示例代码中,我们在接受连接后就交给另一个 goroutine 去处理。除了这个位置,还有另外两个位置可以使用 goroutine:

flowchart TD
    subgraph 模式一
        A1[Accept] --> B1[新goroutine处理读+处理+写]
    end
    subgraph 模式二
        A2[Accept] --> B2[goroutine读]
        B2 --> C2[新goroutine处理]
        B2 -.继续读下一个请求.-> B2
    end
    subgraph 模式三
        A3[Accept] --> B3[goroutine读+处理]
        B3 --> C3[新goroutine写响应]
        B3 -.继续读下一个请求.-> B3
    end
模式说明TCP 通信效率系统复杂度
模式一每个连接一个 goroutine,负责读+处理+写基础最低
模式二读完请求后交给新 goroutine 处理,当前 goroutine 继续读较高较高
模式三处理完后交给新 goroutine 写响应,当前 goroutine 继续读最高最高

由上至下:TCP 通信效率提高,系统复杂度也提高。

实践建议:因为 goroutine 非常轻量(创建一个 goroutine 只需几 KB 内存),所以即便使用模式一,对于绝大多数应用来说性能也可以满足。准确说,虽然很多人尝试开发新的 net 库来取代 Go 自带的,但实际上这些库普遍存在的问题就是 BUG 多,性能提升有限,但编程模型极其复杂。不到逼不得已不要使用这一类的库。

3.2 模式二实现示例

如果你需要提高单个连接的吞吐量,可以使用模式二——读完请求后立即交给新 goroutine 处理,当前 goroutine 继续读下一个请求:

// handleConnV2 使用模式二处理连接
// 读请求和处理请求分离,提高单个连接的吞吐量
func (s *TCPServer) handleConnV2(conn net.Conn) {
	defer conn.Close()
	remoteAddr := conn.RemoteAddr().String()
	reader := bufio.NewReader(conn)

	for {
		// 读取请求(当前 goroutine 负责)
		line, err := reader.ReadString('\n')
		if err != nil {
			fmt.Printf("客户端 %s 断开: %v\n", remoteAddr, err)
			return
		}
		msg := line[:len(line)-1]

		// 将处理和写响应交给新 goroutine
		// 这样当前 goroutine 可以立即回去读下一个请求
		// 注意:这里不能再使用 conn.Write,因为多个 goroutine
		// 同时写同一个 conn 会导致数据交错
		// 需要用 channel 或锁来同步写操作
		go func(message string) {
			response := fmt.Sprintf("大写: %s\n", toUpper(message))
			// 在真实场景中,这里需要加锁或者用专门的写 goroutine
			conn.Write([]byte(response))
		}(msg)
	}
}

注意:模式二和模式三都涉及到多个 goroutine 同时操作一个 conn。多个 goroutine 同时读或同时写同一个连接会导致数据混乱,需要使用 sync.Mutex 加锁,或者引入专门的"写 goroutine"来串行化写操作。这也是为什么模式越高级,复杂度越高的原因。


四、连接池

4.1 为什么需要连接池

在前面的示例代码中,客户端创建的连接都是一次性使用——用完就关。然而,创建一个连接是非常昂贵的:

  • 要发起系统调用(socket、connect 等)
  • TCP 要完成三次握手
  • 高并发的情况下,可能耗尽文件描述符

连接池就是为了复用这些已经创建好的连接,避免频繁创建和销毁。

4.2 连接池的核心参数

参数说明过小的问题过大的问题
InitialCap初始连接数,初始化时直接创建启动时大部分请求需要创建连接浪费资源
MaxIdle最大空闲连接数无法应付突发流量浪费资源
MaxCap最大连接数限制并发能力耗尽资源
MaxIdleTime最大空闲时间连接可能已失效过期连接被复用

4.3 连接池的 Get/Put 流程

flowchart TD
    subgraph Get 获取连接
        G1[开始获取连接] --> G2{有空闲连接?}
        G2 -- 是 --> G3[从空闲队列取出连接]
        G3 --> G4{连接是否过期?}
        G4 -- 是 --> G5[关闭旧连接,创建新连接]
        G4 -- 否 --> G6[返回连接]
        G5 --> G6
        G2 -- 否 --> G7{未超过最大连接数?}
        G7 -- 是 --> G8[创建新连接]
        G8 --> G6
        G7 -- 否 --> G9[阻塞等待,可设超时]
        G9 --> G2
    end
flowchart TD
    subgraph Put 归还连接
        P1[开始归还连接] --> P2{有阻塞的Get请求?}
        P2 -- 是 --> P3[直接把连接交给阻塞的请求]
        P2 -- 否 --> P4{空闲队列已满?}
        P4 -- 否 --> P5[放入空闲队列]
        P4 -- 是 --> P6[关闭连接]
    end

Get 要考虑:

  • 有空闲连接,直接返回
  • 否则,没超过最大连接数,直接创建新的
  • 否则,阻塞调用方

Put 要考虑:

  • 有 Get 请求被阻塞,把连接丢过去
  • 否则,没超过最大空闲连接数,放到空闲列表
  • 否则,直接关闭

4.4 连接池运作图解

flowchart LR
    subgraph 起步
        S1[空闲队列空] --> S2[创建新连接]
    end
    subgraph 超过上限
        L1[已有10个连接] --> L2[新请求被阻塞]
    end
    subgraph 归还-有阻塞请求
        R1[用完放回] --> R2{有阻塞请求?}
        R2 -- 是 --> R3[唤醒一个请求,转交连接]
    end
    subgraph 归还-放入空闲队列
        R4[用完放回] --> R5{有阻塞请求?}
        R5 -- 否 --> R6{空闲队列未满?}
        R6 -- 是 --> R7[放入空闲队列]
    end
    subgraph 归还-空闲队列满
        R8[用完放回] --> R9{空闲队列满了?}
        R9 -- 是 --> R10[关闭连接]
    end

4.5 简单连接池实现

下面我们手写一个简单的连接池,帮助你理解连接池的核心原理:

package main

import (
	"errors"
	"fmt"
	"net"
	"sync"
	"time"
)

// PoolOption 是连接池的配置参数
type PoolOption struct {
	InitialCap    int           // 初始连接数:启动时预先创建的连接数量
	MaxIdle       int           // 最大空闲连接数:空闲队列最多保存多少个连接
	MaxCap        int           // 最大连接数:同时存在的连接上限
	MaxIdleTime   time.Duration // 最大空闲时间:超过这个时间的空闲连接会被关闭
	Factory       func() (net.Conn, error) // 工厂函数:用于创建新连接
}

// ConnPool 是一个简单的连接池实现
type ConnPool struct {
	mu          sync.Mutex           // 互斥锁,保护并发访问
	conns       chan *idleConn       // 空闲连接队列,用 channel 实现
	factory     func() (net.Conn, error) // 创建新连接的工厂函数
	maxCap      int                  // 最大连接数
	maxIdle     int                  // 最大空闲连接数
	maxIdleTime time.Duration        // 最大空闲时间
	numOpen     int                  // 当前已打开的连接总数(包括正在使用的)
}

// idleConn 包装了一个连接和它的最后使用时间
type idleConn struct {
	conn       net.Conn   // 实际的网络连接
	returnTime time.Time  // 归还到池中的时间
}

// NewConnPool 创建一个新的连接池
func NewConnPool(opt PoolOption) (*ConnPool, error) {
	if opt.MaxIdle <= 0 || opt.MaxCap <= 0 {
		return nil, errors.New("MaxIdle 和 MaxCap 必须大于 0")
	}
	if opt.MaxIdle > opt.MaxCap {
		return nil, errors.New("MaxIdle 不能大于 MaxCap")
	}

	// 创建连接池
	p := &ConnPool{
		conns:       make(chan *idleConn, opt.MaxIdle), // 带缓冲的 channel 作为空闲队列
		factory:     opt.Factory,
		maxCap:      opt.MaxCap,
		maxIdle:     opt.MaxIdle,
		maxIdleTime: opt.MaxIdleTime,
	}

	// 预先创建 InitialCap 个连接
	for i := 0; i < opt.InitialCap; i++ {
		conn, err := opt.Factory()
		if err != nil {
			return nil, fmt.Errorf("创建初始连接失败: %w", err)
		}
		p.numOpen++
		p.conns <- &idleConn{conn: conn, returnTime: time.Now()}
	}

	return p, nil
}

// Get 从连接池获取一个连接
// 如果有空闲连接,直接返回;否则创建新连接;如果已达上限,阻塞等待
func (p *ConnPool) Get() (net.Conn, error) {
	p.mu.Lock()

	// 情况1:空闲队列有连接
	select {
	case ic := <-p.conns:
		p.mu.Unlock()
		// 检查连接是否过期
		if p.maxIdleTime > 0 && time.Since(ic.returnTime) > p.maxIdleTime {
			// 连接已过期,关闭它并创建新的
			ic.conn.Close()
			p.mu.Lock()
			p.numOpen-- // 过期连接被关闭,总数减一
			p.mu.Unlock()
			return p.createNewConn()
		}
		return ic.conn, nil
	default:
		// 空闲队列没有连接
		// 情况2:还没达到最大连接数,创建新连接
		if p.numOpen < p.maxCap {
			p.numOpen++
			p.mu.Unlock()
			return p.createNewConn()
		}
		p.mu.Unlock()

		// 情况3:已达最大连接数,阻塞等待其他连接归还
		// 这里可以加超时控制
		select {
		case ic := <-p.conns:
			return ic.conn, nil
		case <-time.After(3 * time.Second):
			return nil, errors.New("获取连接超时")
		}
	}
}

// createNewConn 使用工厂函数创建新连接
func (p *ConnPool) createNewConn() (net.Conn, error) {
	conn, err := p.factory()
	if err != nil {
		// 创建失败,回退计数
		p.mu.Lock()
		p.numOpen--
		p.mu.Unlock()
		return nil, fmt.Errorf("创建连接失败: %w", err)
	}
	return conn, nil
}

// Put 将连接归还到连接池
func (p *ConnPool) Put(conn net.Conn) error {
	p.mu.Lock()

	// 尝试将连接放入空闲队列
	select {
	case p.conns <- &idleConn{conn: conn, returnTime: time.Now()}:
		// 成功放入空闲队列
		p.mu.Unlock()
		return nil
	default:
		// 空闲队列已满,关闭连接
		p.numOpen--
		p.mu.Unlock()
		conn.Close()
		return nil
	}
}

// Close 关闭连接池,释放所有连接
func (p *ConnPool) Close() {
	p.mu.Lock()
	defer p.mu.Unlock()

	close(p.conns)
	for ic := range p.conns {
		ic.conn.Close()
	}
}

func main() {
	// 创建连接池的工厂函数
	// 这里以 TCP 连接为例
	factory := func() (net.Conn, error) {
		return net.Dial("tcp", "127.0.0.1:8080")
	}

	// 创建连接池
	pool, err := NewConnPool(PoolOption{
		InitialCap:  2,           // 初始创建 2 个连接
		MaxIdle:     5,           // 最多空闲 5 个连接
		MaxCap:      10,          // 最多 10 个连接
		MaxIdleTime: 30 * time.Second, // 空闲超过 30 秒的连接会被关闭
		Factory:     factory,
	})
	if err != nil {
		fmt.Printf("创建连接池失败: %v\n", err)
		return
	}
	defer pool.Close()

	// 从连接池获取连接
	conn, err := pool.Get()
	if err != nil {
		fmt.Printf("获取连接失败: %v\n", err)
		return
	}

	// 使用连接发送数据
	conn.Write([]byte("hello\n"))

	// 读取响应
	buf := make([]byte, 1024)
	n, _ := conn.Read(buf)
	fmt.Printf("收到响应: %s", buf[:n])

	// 用完归还连接(而不是关闭)
	pool.Put(conn)
}

4.6 sql.DB 中的连接池管理

Go 标准库 database/sql 中的 sql.DB 就内置了连接池。它也基本遵循前面总结的原理:

  • 利用 channel 来管理空闲连接
  • 利用一个队列来阻塞请求

sql.DB 有很多细节,这里我们只看它怎么管理连接的:

  • 获取连接conn 方法):基本过程和前面讲的差不多,但它是从队尾开始拿空闲连接的。为什么?因为队首的空闲连接更可能已经超过了最大空闲时间(先放进去的更容易过期)。
  • 归还连接putConn 方法):因为 DB 比较复杂,所以在 putConn 的时候要做很多校验,维持好整体状态:处理 ErrBadConn 的情况、确保 dc(driverConn)没有任何人在使用、处理超时。

sql.DB 解决过期连接的懒惰策略可以类比其它如本地缓存的策略——Lazy Evaluation(惰性求值),即只有在真正使用连接时才检查它是否过期,而不是用定时器主动清理。

sql.DB 常用的连接池配置方法:

// 设置连接池参数
db.SetMaxOpenConns(100)            // 最大连接数
db.SetMaxIdleConns(10)             // 最大空闲连接数
db.SetConnMaxLifetime(time.Hour)   // 连接最大存活时间
db.SetConnMaxIdleTime(10 * time.Minute) // 连接最大空闲时间

五、RPC 核心概念

5.1 什么是 RPC

RPC 的全称是 Remote Procedure Call,即远程过程调用。核心就是:如同本地调用一般调用服务器上的方法

想象一下你在本地调用一个函数:

// 本地调用:直接调用同进程内的函数
result := userService.GetById(123)

RPC 要做的事情就是让你可以像调用本地函数一样,调用远程服务器上的函数:

// 远程调用:看起来像本地调用,但实际上请求被发送到了远程服务器
result := userServiceProxy.GetById(123)

因此要解决的问题就是:怎么把左边的本地调用映射过去右边的远程服务。

5.2 RPC 要解决的核心问题

flowchart LR
    subgraph 客户端
        C1[调用 userService.GetById 123] --> C2[代理捕捉调用信息]
        C2 --> C3[编码为字节流]
        C3 --> C4[通过网络发送]
    end
    C4 -->|网络| S1
    subgraph 服务端
        S1[接收数据] --> S2[解码还原调用信息]
        S2 --> S3[查找 userService 服务]
        S3 --> S4[反射执行 GetById 方法]
        S4 --> S5[编码响应]
        S5 --> S6[写回响应]
    end
    S6 -->|网络| C5
    subgraph 客户端
        C5[接收响应] --> C6[解码响应] --> C7[返回结果]
    end

5.3 调用信息

要完成这种映射,首先要解决第一个问题:映射什么

举个例子,假如我们在客户端调用的是 userService.GetById,传入的参数是 int 类型的值 123。那么服务端怎么知道客户端调用的是 userService.GetById,参数是 int 类型的 123

答案很简单:我们把这些信息传过去给服务端,这些信息统称为调用信息

调用信息需要包含:

  • 服务名userService
  • 方法名GetById
  • 参数值123

要不要参数类型? 如果你在支持重载的语言上设计微服务框架,并且决定支持重载,那么你就需要传递参数类型,否则就不需要。Go 语言不支持方法重载,所以不需要传参数类型。

5.4 客户端捕捉本地调用

既然要传递调用信息,那么问题就在于:RPC 客户端怎么获得这些调用信息?用户调用的是 userService.GetById(123),底层框架怎么知道 userServiceGetById123 这些信息?

主要有两种策略:

策略说明代表框架
代码生成通过 IDL 生成客户端代码,生成的代码中已经包含了调用信息的封装gRPC、go-micro
代理机制在运行时动态生成代理对象,拦截方法调用Dubbo

5.5 代理模式

Go 语言中实现 RPC 客户端的关键技术是代理模式:定义一个结构体,为结构体里面的方法类型字段注入调用逻辑。

注意:Go 是没有办法修改方法实现的,所以我们只能迂回救国——不是修改原方法,而是创建一个代理结构体,让用户调用代理的方法。

为了简化微服务框架的代码,我们约定一个方法签名规范:

  • 每一个方法第一个参数必须是 context.Context,第二个就是请求结构体指针,并且只有这两个参数
  • 返回值的第一个是响应,并且必须是指针,第二个是 error,并且只有这两个返回值
// 约定的方法签名
// 第一个参数:context.Context(用于控制超时和传值)
// 第二个参数:请求结构体指针
// 返回值1:响应结构体指针
// 返回值2:error
func (s *UserService) GetById(ctx context.Context, req *GetByIdRequest) (*GetByIdResponse, error)

这种限制主要就是为了简化微服务框架的代码。在真实生产中,你可以保持这个限制,也可以考虑去掉。


六、最简 RPC 框架实现

现在我们将前面学到的所有知识整合起来,从零手写一个最简 RPC 框架。这个框架包含:

  • 客户端:利用反射生成代理,捕捉调用信息,编码后发送到服务端
  • 服务端:接收数据,还原调用信息,利用反射执行方法,写回响应

6.1 整体架构

flowchart TB
    subgraph 客户端 Client
        CL[调用代理方法] --> RF[反射获取调用信息]
        RF --> EN[JSON编码 + 添加长度前缀]
        EN --> SD[发送到服务端]
        SD --> WR[等待并解析响应]
    end
    subgraph 服务端 Server
        AC[Accept 接受连接] --> RD[读取长度前缀 + 读取消息体]
        RD --> DE[JSON解码还原调用信息]
        DE --> FS[根据服务名查找服务]
        FS --> RM[反射执行方法]
        RM --> EN2[JSON编码响应 + 添加长度前缀]
        EN2 --> SD2[写回响应]
    end
    SD -->|TCP 网络| AC
    SD2 -->|TCP 网络| WR

6.2 定义数据结构

首先定义 RPC 请求和响应的数据结构,以及服务注册表:

package mrpc

import (
	"context"
	"encoding/binary"
	"encoding/json"
	"errors"
	"fmt"
	"io"
	"net"
	"reflect"
	"sync"
	"time"
)

// ============================================================
// 第一部分:数据结构定义
// ============================================================

// Request 是 RPC 请求的传输结构
// 客户端把调用信息编码成 Request,再序列化成 JSON 发送
type Request struct {
	Service string        `json:"service"` // 服务名,例如 "UserService"
	Method  string        `json:"method"`  // 方法名,例如 "GetById"
	Args    []interface{} `json:"args"`    // 参数列表,例如 [123]
}

// Response 是 RPC 响应的传输结构
// 服务端处理完请求后,把结果编码成 Response,再序列化成 JSON 返回
type Response struct {
	Code int         `json:"code"` // 状态码:0 表示成功,非 0 表示失败
	Msg  string      `json:"msg"`  // 错误信息
	Data interface{} `json:"data"` // 返回数据
}

// ============================================================
// 第二部分:服务端实现
// ============================================================

// Server 是 RPC 服务端
// 它负责监听端口、接收连接、解析请求、执行方法、返回响应
type Server struct {
	addr string                    // 监听地址
	mu   sync.RWMutex              // 保护 serviceMap 的并发访问
	// serviceMap 存储已注册的服务
	// key: 服务名(如 "UserService")
	// value: 服务实例的反射值
	serviceMap map[string]reflect.Value
}

// NewServer 创建一个新的 RPC 服务器
func NewServer(addr string) *Server {
	return &Server{
		addr:       addr,
		serviceMap: make(map[string]reflect.Value),
	}
}

// Register 注册一个服务到服务器
// 参数 svc 是服务实例,它的方法会被远程调用
func (s *Server) Register(svc interface{}) error {
	s.mu.Lock()
	defer s.mu.Unlock()

	// 使用反射获取服务的类型信息
	// reflect.TypeOf 返回接口值的动态类型
	svcType := reflect.TypeOf(svc)
	svcName := svcType.Elem().Name() // 获取结构体名称作为服务名

	// 检查是否已经注册过同名服务
	if _, exists := s.serviceMap[svcName]; exists {
		return fmt.Errorf("服务 %s 已存在", svcName)
	}

	// 将服务实例的反射值存入 map
	// reflect.ValueOf 返回接口值的反射值
	s.serviceMap[svcName] = reflect.ValueOf(svc)
	return nil
}

// Start 启动 RPC 服务器
func (s *Server) Start() error {
	listener, err := net.Listen("tcp", s.addr)
	if err != nil {
		return fmt.Errorf("监听失败: %w", err)
	}
	defer listener.Close()

	fmt.Printf("RPC 服务器启动,监听地址: %s\n", s.addr)

	for {
		conn, err := listener.Accept()
		if err != nil {
			continue
		}
		// 每个连接交给独立的 goroutine 处理
		go s.handleConn(conn)
	}
}

// handleConn 处理单个客户端连接
func (s *Server) handleConn(conn net.Conn) {
	defer conn.Close()

	for {
		// 第一步:读取请求
		// 先读取 4 字节的长度前缀(表示消息体有多少字节)
		// 再根据长度读取完整的消息体
		req, err := s.readRequest(conn)
		if err != nil {
			// 如果是 EOF,说明客户端关闭了连接
			if errors.Is(err, io.EOF) {
				return
			}
			fmt.Printf("读取请求失败: %v\n", err)
			return
		}

		// 第二步:处理请求
		resp := s.handleRequest(req)

		// 第三步:写回响应
		if err := s.writeResponse(conn, resp); err != nil {
			fmt.Printf("写回响应失败: %v\n", err)
			return
		}
	}
}

// readRequest 从连接中读取一个完整的 RPC 请求
// 通信协议:[4字节长度][消息体JSON]
func (s *Server) readRequest(conn net.Conn) (*Request, error) {
	// 先读 4 字节的长度前缀
	// 使用 binary.BigEndian 将 4 个字节解读为一个 uint32 整数
	// 这 4 个字节表示后面消息体的字节长度
	lengthBuf := make([]byte, 4)
	if _, err := io.ReadFull(conn, lengthBuf); err != nil {
		return nil, err
	}
	msgLen := binary.BigEndian.Uint32(lengthBuf)

	// 根据长度读取完整的消息体
	msgBuf := make([]byte, msgLen)
	if _, err := io.ReadFull(conn, msgBuf); err != nil {
		return nil, err
	}

	// 将 JSON 反序列化成 Request 结构体
	var req Request
	if err := json.Unmarshal(msgBuf, &req); err != nil {
		return nil, fmt.Errorf("JSON 反序列化失败: %w", err)
	}
	return &req, nil
}

// writeResponse 将响应写回客户端
// 通信协议:[4字节长度][消息体JSON]
func (s *Server) writeResponse(conn net.Conn, resp *Response) error {
	// 将 Response 序列化成 JSON
	data, err := json.Marshal(resp)
	if err != nil {
		return fmt.Errorf("JSON 序列化失败: %w", err)
	}

	// 先写 4 字节的长度前缀
	lengthBuf := make([]byte, 4)
	binary.BigEndian.PutUint32(lengthBuf, uint32(len(data)))
	if _, err := conn.Write(lengthBuf); err != nil {
		return err
	}

	// 再写消息体
	if _, err := conn.Write(data); err != nil {
		return err
	}
	return nil
}

// handleRequest 处理单个请求:查找服务 -> 反射执行方法
func (s *Server) handleRequest(req *Request) *Response {
	s.mu.RLock()
	svcValue, ok := s.serviceMap[req.Service]
	s.mu.RUnlock()

	// 检查服务是否存在
	if !ok {
		return &Response{Code: 404, Msg: fmt.Sprintf("服务 %s 不存在", req.Service)}
	}

	// 获取方法参数的类型信息,用于构造反射调用的参数
	svcType := svcValue.Type()
	// 根据方法名查找方法
	// 这里简化处理:约定方法有两个参数 (context.Context, *Request)
	// 所以我们需要构造这两个参数的类型
	method, ok := svcType.MethodByName(req.Method)
	if !ok {
		return &Response{Code: 404, Msg: fmt.Sprintf("方法 %s 不存在", req.Method)}
	}

	// 构造方法参数
	// 约定:第一个参数是 context.Context,第二个参数是请求结构体
	// 反射调用时,第一个参数是接收者(服务实例本身)
	in := make([]reflect.Value, len(method.Type.In()))
	in[0] = svcValue // 接收者

	// 构造 context.Context 参数(使用 context.Background)
	if len(in) > 1 {
		in[1] = reflect.ValueOf(context.Background())
	}

	// 构造请求参数
	// 因为 JSON 反序列化后,参数是 []interface{}
	// 我们需要将每个参数转换为方法期望的类型
	for i := 2; i < len(in); i++ {
		// 获取方法第 i 个参数的类型
		argType := method.Type.In(i)
		// 将 JSON 解析出的参数转换为目标类型
		// 这里通过 JSON 序列化再反序列化来实现类型转换
		argBytes, _ := json.Marshal(req.Args[i-2])
		argValue := reflect.New(argType)
		json.Unmarshal(argBytes, argValue.Interface())
		in[i] = argValue.Elem()
	}

	// 反射调用方法
	// method.Func.Call 返回 []reflect.Value,即方法的返回值列表
	out := method.Func.Call(in)

	// 约定:返回值第一个是响应,第二个是 error
	var resp Response
	if len(out) >= 2 {
		// 检查是否有 error
		if errInterface := out[1].Interface(); errInterface != nil {
			resp = Response{Code: 500, Msg: errInterface.(error).Error()}
		} else {
			// 成功,取第一个返回值作为数据
			resp = Response{Code: 0, Msg: "success", Data: out[0].Interface()}
		}
	}
	return &resp
}

6.3 客户端代理实现

客户端的核心是利用反射生成代理。代理对象在用户调用方法时,自动拦截调用,将调用信息编码后发送到服务端:

// ============================================================
// 第三部分:客户端实现
// ============================================================

// Client 是 RPC 客户端
// 它负责连接服务器、发送请求、接收响应
type Client struct {
	conn net.Conn // 与服务端的 TCP 连接
}

// NewClient 创建并连接一个 RPC 客户端
func NewClient(addr string) (*Client, error) {
	// 使用 DialTimeout 避免一直阻塞
	// 设置 5 秒超时,如果服务器无响应则返回错误
	conn, err := net.DialTimeout("tcp", addr, 5*time.Second)
	if err != nil {
		return nil, fmt.Errorf("连接服务器失败: %w", err)
	}
	return &Client{conn: conn}, nil
}

// Close 关闭客户端连接
func (c *Client) Close() {
	c.conn.Close()
}

// Call 是底层的 RPC 调用方法
// 参数:
//   - service: 服务名
//   - method: 方法名
//   - args: 参数列表
// 返回:响应数据或错误
func (c *Client) Call(service, method string, args ...interface{}) (*Response, error) {
	// 第一步:构造请求
	req := &Request{
		Service: service,
		Method:  method,
		Args:    args,
	}

	// 第二步:将请求序列化成 JSON
	data, err := json.Marshal(req)
	if err != nil {
		return nil, fmt.Errorf("JSON 序列化失败: %w", err)
	}

	// 第三步:发送请求
	// 通信协议:[4字节长度][消息体JSON]
	lengthBuf := make([]byte, 4)
	binary.BigEndian.PutUint32(lengthBuf, uint32(len(data)))
	if _, err := c.conn.Write(lengthBuf); err != nil {
		return nil, fmt.Errorf("发送长度前缀失败: %w", err)
	}
	if _, err := c.conn.Write(data); err != nil {
		return nil, fmt.Errorf("发送请求失败: %w", err)
	}

	// 第四步:读取响应
	// 先读 4 字节长度前缀
	respLenBuf := make([]byte, 4)
	if _, err := io.ReadFull(c.conn, respLenBuf); err != nil {
		return nil, fmt.Errorf("读取响应长度失败: %w", err)
	}
	respLen := binary.BigEndian.Uint32(respLenBuf)

	// 再读消息体
	respBuf := make([]byte, respLen)
	if _, err := io.ReadFull(c.conn, respBuf); err != nil {
		return nil, fmt.Errorf("读取响应体失败: %w", err)
	}

	// 反序列化响应
	var resp Response
	if err := json.Unmarshal(respBuf, &resp); err != nil {
		return nil, fmt.Errorf("JSON 反序列化响应失败: %w", err)
	}
	return &resp, nil
}

// ============================================================
// 第四部分:反射生成代理
// ============================================================

// NewProxy 使用反射为指定的服务接口生成代理
// 参数:
//   - client: RPC 客户端
//   - service: 服务名
// 返回一个代理对象,调用它的方法会自动发起 RPC 调用
func NewProxy(client *Client, serviceName string) *Proxy {
	return &Proxy{
		client:      client,
		serviceName: serviceName,
	}
}

// Proxy 是通用代理结构体
// 它拦截对服务方法的调用,将调用转发到远程服务器
type Proxy struct {
	client      *Client // RPC 客户端
	serviceName string  // 服务名
}

// CallMethod 是通用的方法调用入口
// 用户通过这个方法来调用远程服务
// 参数:
//   - method: 方法名
//   - args: 参数列表
func (p *Proxy) CallMethod(method string, args ...interface{}) (*Response, error) {
	return p.client.Call(p.serviceName, method, args...)
}

6.4 完整示例:用户服务

下面我们用这个最简 RPC 框架来实现一个"用户服务"的完整示例:

package main

import (
	"context"
	"encoding/json"
	"fmt"
	"mrpc"
	"time"
)

// ============================================================
// 定义服务接口和实现
// ============================================================

// GetByIdRequest 是 GetById 方法的请求参数
type GetByIdRequest struct {
	Id int `json:"id"` // 用户 ID
}

// GetByIdResponse 是 GetById 方法的响应
type GetByIdResponse struct {
	Id   int    `json:"id"`   // 用户 ID
	Name string `json:"name"` // 用户名
	Age  int    `json:"age"`  // 年龄
}

// UserService 是用户服务
// 注意:方法的签名必须符合约定
// (ctx context.Context, req *GetByIdRequest) (*GetByIdResponse, error)
type UserService struct{}

// GetById 根据用户 ID 查询用户信息
// 这是一个模拟实现,真实场景中会查数据库
func (s *UserService) GetById(ctx context.Context, req *GetByIdRequest) (*GetByIdResponse, error) {
	// 模拟数据库查询
	if req.Id == 123 {
		return &GetByIdResponse{
			Id:   123,
			Name: "张三",
			Age:  25,
		}, nil
	}
	// 用户不存在
	return nil, fmt.Errorf("用户 %d 不存在", req.Id)
}

// ============================================================
// 服务端
// ============================================================

func main() {
	// --- 启动服务端 ---
	server := mrpc.NewServer(":9090")

	// 注册 UserService 服务
	// 传入的是指针,因为方法定义在 *UserService 上
	if err := server.Register(&UserService{}); err != nil {
		fmt.Printf("注册服务失败: %v\n", err)
		return
	}

	// 在另一个 goroutine 中启动服务器
	go func() {
		if err := server.Start(); err != nil {
			fmt.Printf("服务器启动失败: %v\n", err)
		}
	}()

	// --- 客户端调用 ---
	// 等待服务器启动
	time.Sleep(100 * time.Millisecond)

	// 创建客户端
	client, err := mrpc.NewClient("127.0.0.1:9090")
	if err != nil {
		fmt.Printf("连接服务器失败: %v\n", err)
		return
	}
	defer client.Close()

	// 创建代理
	proxy := mrpc.NewProxy(client, "UserService")

	// 通过代理调用远程方法
	// 就像调用本地方法一样简单!
	resp, err := proxy.CallMethod("GetById", 123)
	if err != nil {
		fmt.Printf("RPC 调用失败: %v\n", err)
		return
	}

	if resp.Code != 0 {
		fmt.Printf("服务端返回错误: %s\n", resp.Msg)
		return
	}

	// 将响应数据转换回 GetByIdResponse 结构体
	// 因为 JSON 反序列化后 Data 是 interface{} 类型
	dataBytes, _ := json.Marshal(resp.Data)
	var user GetByIdResponse
	json.Unmarshal(dataBytes, &user)

	fmt.Printf("查询成功: %+v\n", user)
	// 输出: 查询成功: {Id:123 Name:张三 Age:25}
}

6.5 通信协议详解

我们的最简 RPC 使用了长度前缀 + JSON 的通信协议:

flowchart LR
    subgraph 消息格式
        A[4字节: 消息体长度] --> B[N字节: JSON消息体]
    end

为什么需要长度前缀?

TCP 是流式协议,没有消息边界。如果客户端连续发送两条消息,服务端可能一次读到一条半消息,或者半条消息。长度前缀告诉服务端"接下来有多少字节是一条完整的消息",从而正确拆分消息。

这就是 PDF 中提到的"先读头部,根据头部得知 Body 有多长,再把剩下的数据读出来"。

// 发送端:先发4字节长度,再发消息体
lengthBuf := make([]byte, 4)
binary.BigEndian.PutUint32(lengthBuf, uint32(len(data)))
conn.Write(lengthBuf)    // 4字节长度前缀
conn.Write(data)         // 消息体

// 接收端:先读4字节长度,再读对应长度的消息体
lengthBuf := make([]byte, 4)
io.ReadFull(conn, lengthBuf)           // 先读4字节
msgLen := binary.BigEndian.Uint32(lengthBuf)
msgBuf := make([]byte, msgLen)
io.ReadFull(conn, msgBuf)              // 再读 msgLen 字节

6.6 最简 RPC 总结

flowchart TB
    subgraph 客户端
        C1[初始化代理] --> C2[代理利用反射获得调用信息]
        C2 --> C3[将调用信息编码成字节流]
        C3 --> C4[加上长度字段发送到服务端]
        C4 --> C5[等待并且解析响应]
    end
    subgraph 服务端
        S1[启动服务器监听端口] --> S2[接收连接并读取数据]
        S2 --> S3[将数据还原回调用信息]
        S3 --> S4[根据服务名查找注册的服务]
        S4 --> S5[利用反射执行方法调用]
        S5 --> S6[写回响应]
    end
    C4 -->|TCP| S2
    S6 -->|TCP| C5

客户端步骤:

  1. 初始化代理
  2. 代理会利用反射获得调用信息
  3. 将调用信息编码成字节流,加上长度字段
  4. 将数据发送到服务端
  5. 等待并且解析响应

服务端步骤:

  1. 启动服务器监听端口
  2. 接收连接,并且读取数据
  3. 将数据还原回调用信息
  4. 根据服务名查找该实例上注册的服务
  5. 利用反射执行方法调用
  6. 写回响应

七、面试要点总结

7.1 网络编程

  • 网络基础知识:包含 TCP 和 UDP 的基础知识,三次握手和四次挥手
  • Go TCP 服务器:如何利用 Go 写一个简单的 TCP 服务器。直接面 net 里面的 API 是很少见的,但如果有编程题环节,可能会让你直接写一个简单的 TCP 服务器
  • goroutine 和连接的关系:可以在不同的环节使用不同的 goroutine,以充分利用 TCP 的全双工通信
  • 连接池参数:初始连接、最大空闲连接、最大连接数、最大空闲时间
  • 连接池运作原理:拿连接会发生什么,放回去又会发生什么
  • sql.DB 过期连接:懒惰策略,可以类比其它如本地缓存的策略
  • 手写连接池:注重考察代码能力的公司可能会让你手写代码

7.2 微服务框架

  • 微服务框架是什么:主要就是解决两个问题——通信和服务治理
  • 为什么使用微服务架构:本质上是为了分而治之,将业务拆分之后独立治理、部署
  • RPC 框架和 RESTful 的区别:两者基本没关联,全是区别。唯一的关联就是 RPC 框架可以利用 RESTful 来实现。RESTful 是指符合 REST 风格的 HTTP 接口,而 RPC 指的是远程过程调用,从本质上就是两回事
  • RPC 框架和 Web 框架的区别:基本也没什么关联,都是区别。唯一的共同点是可以通过对 Web 框架进行封装来实现 RPC 通信

7.3 RPC 核心

  • 什么是 RPC:远程过程调用,类似的还有 RMI(远程方法调用)
  • RPC 相比 HTTP 的优势:不必关心 HTTP 调用的细节,对于使用者来说就如同本地调用一般
  • RPC 框架的要点:客户端捕捉调用信息,编码成二进制,发送到服务端。服务端查找本地服务,执行调用,写回响应。任何一个 RPC 框架都类似
  • RPC 框架怎么捕捉本地调用信息:主要依赖于代理模式和代码生成技术
  • 什么是代理模式/动态代理模式:动态代理可以看做是动态生成的代理,一般是指运行时生成的代理
  • 动态代理技术能用来做什么:四个字,为所欲为。在这里就是用来发起 RPC 调用,然后再返回响应

总结

本教程从微服务框架概览出发,讲解了 Go 网络编程的基础知识(net 包、TCP 服务器与客户端、错误处理、goroutine 使用模式),深入分析了连接池的原理与实现,最后从零手写了一个最简 RPC 框架。

关键知识点回顾:

  1. 网络通信的核心是 net.Listen(服务端)和 net.Dial(客户端)
  2. 处理连接的基本流程是:读数据 -> 处理数据 -> 回写响应
  3. 连接池通过复用连接来避免频繁创建/销毁的开销
  4. RPC 的本质是"如同本地调用一般调用远程方法"
  5. 代理模式反射是实现 RPC 客户端的核心技术
  6. 长度前缀 + JSON 是最简单的 RPC 通信协议

下一章我们将深入讲解 RPC 协议的设计与实现,包括更完善的协议设计、序列化方案选择等内容。


自测题与动手练习

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

  1. 单体应用也可以部署多个实例,那它和微服务多服务部署的本质区别是什么?(提示:看不同实例之间有没有网络间交互
  2. goroutine 模式二、模式三都把"读"和"处理/写"拆到不同 goroutine。为什么多个 goroutine 同时 Write 同一个 conn 会出乱子?应该怎么规避?
  3. 连接池 Get 时,空闲队列里有连接但已经超过 MaxIdleTime 过期了,该怎么处理?为什么 sql.DB 从队尾(而不是队首)取空闲连接?
  4. RPC 客户端要告诉服务端"调用了哪个服务、哪个方法、什么参数"。Go 为什么不需要把参数类型也传过去?什么语言场景下才必须传?
  5. 最简 RPC 用"4 字节长度前缀 + JSON"的格式。如果去掉长度前缀、直接发 JSON,连续发两条消息服务端会怎样?这跟 TCP 的什么特性有关?

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

  1. 在 2.3 的 TCP 服务器基础上,把 handleConn 改成模式二(读请求后交给新 goroutine 处理),体会单连接吞吐提升,同时给它加一把 sync.Mutex 保护 conn.Write,观察是否还出现数据交错。
  2. 把 4.5 的连接池 MaxCap 改成 2,写个并发 10 的循环抢连接,观察"获取连接超时"什么时候触发、空闲队列满时多余连接如何被关掉。
  3. 在 6.4 的用户服务里新增一个 SayHello(ctx, *HelloRequest) (*HelloResponse, error) 方法,注册进 UserService 并走通代理调用,真正理解 method.Func.Call 是怎么找到并执行方法的。

本章小结

  • 微服务本质是分而治之:拆分后问题收敛为「通信 + 服务治理」两类;框架选型上"遇事不决用 gRPC",业务向可优先 Kratos / go-zero。
  • 网络通信建立在 net.Listen / Accept(服务端)与 net.Dial(客户端)之上;处理连接的三段式是「读数据 → 处理 → 回写响应」;出错直接关连接最简单。
  • goroutine 模式越高级吞吐越高但复杂度越高;连接池靠复用连接省去频繁建连开销,核心是 Get / Put 对空闲队列与上限的管理;sql.DB 用惰性策略检查过期连接。
  • RPC 的本质是"像本地调用一样调远程方法":靠代理模式 + 反射捕捉调用信息,靠「长度前缀 + JSON」解决 TCP 流式拆包。
  • 最简 RPC 已完整串起「client 编码发送 → server 解码反射执行 → 回写响应」的闭环,是后续理解 gRPC / Kitex 等工业级框架的基石。
About Me

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

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

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

目标

学AI,加油!加油!