学习目标
学完本章你应该能够:
- 说清微服务相比单体的核心动机(分而治之)与拆分后带来的四类新问题:服务发现、远程通信、容错、可观测。
- 用 Go 的
net包写出一个最简单的 TCP 服务器与客户端,讲清Listen/Accept/Dial的分工,以及"读数据 → 处理 → 回写响应"的三段式。 - 解释连接池的四个核心参数(InitialCap / MaxIdle / MaxCap / MaxIdleTime)与 Get / Put 流程,并说出 goroutine 三种使用模式在吞吐与复杂度上的取舍。
- 讲清 RPC 的本质(像本地调用一样调远程方法)与必须解决的三个问题:调用信息是什么、客户端怎么捕捉、怎么编码传过去再还原。
- 手写最简 RPC 框架时,能讲清代理模式 + 反射 + 长度前缀协议各自解决了什么问题,面试能讲成一段闭环故事。
前置知识:
- Go 基础:
goroutine、channel、interface、反射(reflect)入门。 - 网络常识:TCP 三次握手、流式协议没有消息边界。
- JSON 序列化基础。
本章你会动手做的事:
- 跑通 2.3 的 TCP 服务器 + 2.7 的客户端,观察
Accept后每个连接一个 goroutine 的并发行为。 - 用 4.5 的最小连接池把
MaxCap调到 2,写个并发 10 的压测,观察"获取连接超时"。 - 在 6.4 的用户服务上加一个
SayHello方法,走通代理调用,理解反射是怎么找到方法并执行的。
前言:从单体到微服务
在互联网发展的早期,一个系统通常是一个"单体应用"——所有功能模块打包在一起,部署在同一个进程中。单体应用虽然简单,但随着业务规模增长,它会变得越来越臃肿:改一个小 bug 需要重新部署整个系统,一个模块的内存泄漏会拖垮所有功能。
微服务架构的核心思想就是"分而治之":把一个大系统拆分成多个独立的小服务,每个服务负责一块业务,服务之间通过网络通信。这样每个服务可以独立开发、独立部署、独立扩展。
但拆分之后,新的问题随之而来:
- 服务 A 怎么找到服务 B 的地址?(服务发现)
- 服务 A 怎么调用服务 B?(远程通信)
- 服务 B 挂了怎么办?(容错机制)
- 怎么监控服务的健康状态?(可观测性)
这些问题正是微服务框架要解决的。本教程将从最底层的网络编程讲起,一步步带你理解微服务框架的核心——RPC(远程过程调用),并最终手写一个"最简 RPC 框架"。
本教程涵盖以下主题:
- 微服务框架概览——理解微服务框架要解决什么问题
- Go 网络编程基础——使用
net包进行 TCP 通信 - goroutine 与连接处理——并发处理连接的三种模式
- 连接池——复用连接提升性能
- RPC 核心概念——远程过程调用的本质
- 最简 RPC 框架实现——从零手写一个 RPC 框架
- 面试要点总结
一、微服务框架概览
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 包是网络相关的核心包,里面包含了 http、rpc 等关键包。理解 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 循环内重复三个步骤:
- 读数据:读数据要根据上层协议来决定怎么读。例如,简单的 RPC 协议一般分成两段读——先读头部,根据头部得知 Body 有多长,再把剩下的数据读出来。
- 处理数据:根据业务逻辑处理请求。
- 回写响应:即便处理数据出错,也要返回一个错误给客户端,不然客户端不知道你处理出错了。
2.5 错误处理
在读写的时候,都可能遇到错误。一般来说代表连接已经关掉的是这三个:
io.EOF:正常读到末尾,客户端正常关闭了连接io.ErrUnexpectedEOF:读到一半连接断了net.ErrClosed:连接已经被关闭了
实践建议:只要是出错了就直接关闭连接,这样对客户端和服务端代码都简单。不要试图从错误中恢复连接,因为连接可能已经处于不一致的状态。
2.6 net.Dial:客户端连接
net.Dial 是指创建一个连接,连上远端的服务器。它也有几个类似的方法:
DialIPDialTCPDialUDPDialUnixDialTimeout:多了一个超时参数
实践建议:直接使用
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
endflowchart TD
subgraph Put 归还连接
P1[开始归还连接] --> P2{有阻塞的Get请求?}
P2 -- 是 --> P3[直接把连接交给阻塞的请求]
P2 -- 否 --> P4{空闲队列已满?}
P4 -- 否 --> P5[放入空闲队列]
P4 -- 是 --> P6[关闭连接]
endGet 要考虑:
- 有空闲连接,直接返回
- 否则,没超过最大连接数,直接创建新的
- 否则,阻塞调用方
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[关闭连接]
end4.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[返回结果]
end5.3 调用信息
要完成这种映射,首先要解决第一个问题:映射什么?
举个例子,假如我们在客户端调用的是 userService.GetById,传入的参数是 int 类型的值 123。那么服务端怎么知道客户端调用的是 userService.GetById,参数是 int 类型的 123?
答案很简单:我们把这些信息传过去给服务端,这些信息统称为调用信息。
调用信息需要包含:
- 服务名:
userService - 方法名:
GetById - 参数值:
123
要不要参数类型? 如果你在支持重载的语言上设计微服务框架,并且决定支持重载,那么你就需要传递参数类型,否则就不需要。Go 语言不支持方法重载,所以不需要传参数类型。
5.4 客户端捕捉本地调用
既然要传递调用信息,那么问题就在于:RPC 客户端怎么获得这些调用信息?用户调用的是 userService.GetById(123),底层框架怎么知道 userService、GetById、123 这些信息?
主要有两种策略:
| 策略 | 说明 | 代表框架 |
|---|---|---|
| 代码生成 | 通过 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 网络| WR6.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客户端步骤:
- 初始化代理
- 代理会利用反射获得调用信息
- 将调用信息编码成字节流,加上长度字段
- 将数据发送到服务端
- 等待并且解析响应
服务端步骤:
- 启动服务器监听端口
- 接收连接,并且读取数据
- 将数据还原回调用信息
- 根据服务名查找该实例上注册的服务
- 利用反射执行方法调用
- 写回响应
七、面试要点总结
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 框架。
关键知识点回顾:
- 网络通信的核心是
net.Listen(服务端)和net.Dial(客户端) - 处理连接的基本流程是:读数据 -> 处理数据 -> 回写响应
- 连接池通过复用连接来避免频繁创建/销毁的开销
- RPC 的本质是"如同本地调用一般调用远程方法"
- 代理模式和反射是实现 RPC 客户端的核心技术
- 长度前缀 + JSON 是最简单的 RPC 通信协议
下一章我们将深入讲解 RPC 协议的设计与实现,包括更完善的协议设计、序列化方案选择等内容。
自测题与动手练习
自测题(合上书能答出来,才算懂):
- 单体应用也可以部署多个实例,那它和微服务多服务部署的本质区别是什么?(提示:看不同实例之间有没有网络间交互)
- goroutine 模式二、模式三都把"读"和"处理/写"拆到不同 goroutine。为什么多个 goroutine 同时
Write同一个conn会出乱子?应该怎么规避? - 连接池
Get时,空闲队列里有连接但已经超过MaxIdleTime过期了,该怎么处理?为什么sql.DB从队尾(而不是队首)取空闲连接? - RPC 客户端要告诉服务端"调用了哪个服务、哪个方法、什么参数"。Go 为什么不需要把参数类型也传过去?什么语言场景下才必须传?
- 最简 RPC 用"4 字节长度前缀 + JSON"的格式。如果去掉长度前缀、直接发 JSON,连续发两条消息服务端会怎样?这跟 TCP 的什么特性有关?
动手练习(建议真做一遍):
- 在 2.3 的 TCP 服务器基础上,把
handleConn改成模式二(读请求后交给新 goroutine 处理),体会单连接吞吐提升,同时给它加一把sync.Mutex保护conn.Write,观察是否还出现数据交错。 - 把 4.5 的连接池
MaxCap改成 2,写个并发 10 的循环抢连接,观察"获取连接超时"什么时候触发、空闲队列满时多余连接如何被关掉。 - 在 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 等工业级框架的基石。