This repository has been archived on 2026-05-24. You can view files and clone it. You cannot open issues or pull requests or push a commit.
Files
all-in-kingsoft/hhs/gRPC/2. gRPC 核心篇/05-RPC 调用模式总览.md
T
2026-05-11 19:02:38 +08:00

350 lines
11 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
---
tags: [gRPC, RPC, Streaming, Go, Microservice]
create time: 2026-05-11 16:40
---
# RPC 调用模式总览
## 概述
gRPC 提供四种 RPC 调用模式,从最简单的请求-响应到完全的双向流。选择合适的模式是设计高性能 API 的第一步。搞懂它们的区别,你就知道什么时候该用简单调用、什么时候需要双向通信。
> [!question] 如果只选一种模式能走天下吗?
> 技术上可以——Unary RPC 确实能解决几乎所有问题。但强行用 Unary 实现实时推送,意味着你要轮询(polling),这会产生大量无效请求和延迟。模式选择本质上是在"延迟 vs 资源消耗"之间做 tradeoff。
## RPC 模式全景图
```mermaid
flowchart LR
A[Unary] -->|"一问一答"| B[最简单]
C[Server Stream] -->|"一问多答"| D[广播式]
E[Client Stream] -->|"多问一答"| F[收集式]
G[BiDi Stream] -->|"多问多答"| H[全双工]
style A fill:#00B6BC,color:#fff
style G fill:#EE5A24,color:#fff
```
> [!example] 各模式数据流向速览
> 左列为 Client,右列为 Server,箭头方向表示数据流动方向。
```mermaid
sequenceDiagram
participant C as Client
participant S as Server
rect rgba(0, 182, 188, 0.1)
Note over C,S: Unary — 阻塞式一次往返
C->>S: request
S-->>C: response
end
rect rgba(0, 216, 102, 0.1)
Note over C,S: Server Stream — 请求后连续响应
C->>S: request
S-->>C: response 1
S-->>C: response 2
S-->>C: ... n (EOF)
end
rect rgba(79, 195, 247, 0.1)
Note over C,S: Client Stream — 连续发送后一次性响应
C->>S: chunk 1
C->>S: chunk 2
C->>S: ... n (done)
S-->>C: result
end
rect rgba(238, 90, 36, 0.1)
Note over C,S: BiDi Stream — 双向独立通信
C->>S: msg 1
S-->>C: reply 1
C->>S: msg 2
S-->>C: reply 2
end
```
## Unary RPC(普通调用)
最经典、最常见的模式,等同于 REST 的 request-response。
**特点:**
- 一次客户端请求,一次服务端响应
- 简单、易调试、可直接映射 HTTP GET/POST
- 适合:CRUD 操作、短查询、标准 API 端点
```go
// proto 定义
// rpc GetUser(GetUserRequest) returns (GetUserResponse);
// Client side
func main() {
conn, _ := grpc.Dial("localhost:50051", grpc.WithTransportCredentials(insecure.NewCredentials()))
defer conn.Close()
client := pb.NewUserServiceClient(conn)
resp, err := client.GetUser(context.Background(), &pb.GetUserRequest{Id: 42})
if err != nil {
log.Fatal(err) // error comes from transport or server handler
}
fmt.Println(resp.Name)
}
// Server side
type Server struct {
pb.UnimplementedUserServiceServer
}
func (s *Server) GetUser(ctx context.Context, req *pb.GetUserRequest) (*pb.GetUserResponse, error) {
user := fetchFromDB(req.GetId())
return &pb.GetUserResponse{Name: user.Name}, nil
}
```
**思考题:**如果一个接口需要 30 秒才能返回结果,你应该用 Unary 还是其他模式?(提示:考虑超时和连接的持有时间)
> [!answer]+ 参考答案
> **可以用 Unary,但要注意三件事:**
>
> 1. **设置合理的 context timeout**:`context.WithTimeout`,避免无限期等待
> 2. **HTTP/2 ping keepalive**:gRPC 默认会发送 keepalive ping 防止代理(Nginx/LB)因"连接空闲"而断开
> 3. **是否真的需要阻塞等待**:如果是异步任务(如报表生成),更好的做法是——Unary 提交任务 + 轮询/回调通知结果,而非让一个 RPC 连接挂 30 秒。
>
> Streaming 并不会延长超时时间——超时由 context 控制,与使用哪种 RPC 模式无关。
## Server Streaming RPC(服务端流)
Client 发一个请求,Server 持续返回多个响应。
**经典场景:**
- 实时通知推送
- 大列表分批返回
- 日志/事件流订阅
```protobuf
rpc Subscribe(SubscribeRequest) returns (stream Event);
```
```go
// Server side - 持续发送事件(注意不要用 log.Fatal,会杀死整个进程)
func (s *Server) Subscribe(req *pb.SubscribeRequest, stream pb.UserService_SubscribeServer) error {
for _, event := range s.watchEvents(req.Topic) {
if err := stream.Send(&pb.Event{Data: event}); err != nil {
return err // channel closed or context canceled
}
}
return nil
}
```
> [!note] import 提示
> 流式示例中用到以下包:`context`, `fmt`, `io`, `log`.
> 每个文件只需 import 实际用到的即可,不必照抄。
```go
// Client side - 遍历接收事件(生产环境建议用 recover + defer 做错误恢复)
func main() {
client := pb.NewUserServiceClient(conn)
stream, err := client.Subscribe(context.Background(), &pb.SubscribeRequest{Topic: "orders"})
if err != nil {
log.Fatal(err)
}
for {
event, err := stream.Recv()
if err == io.EOF {
break // server finished sending
}
if err != nil {
return err // ⚠️ 这里用 return 而非 log.Fatal,服务端的 handler 不能 kill 进程
}
fmt.Printf("received: %s\n", event.Data)
}
}
```
> [!tip] 注意:Client 仍然可以通过 context cancel 随时中断流。这是调试流式问题时最容易忽略的一点——不是 Server 主动关了连接,而是 Client 放弃了。
## Client Streaming RPC(客户端流)
Client 持续发送多个请求,Server 在所有数据发送完毕后返回一个响应。
**经典场景:**
- 大批量数据上传
- 文件分片聚合处理
- 批量日志采集
```protobuf
rpc Upload(stream FileChunk) returns (UploadResult);
```
```go
// Client side - 流式发送数据分片
func uploadFile(client pb.FileServiceClient, chunks [][]byte) error {
stream, err := client.Upload(context.Background())
if err != nil {
return err
}
for _, chunk := range chunks {
if err := stream.Send(&pb.FileChunk{Data: chunk}); err != nil {
return err
}
}
result, err := stream.CloseAndRecv() // 结束发送,获取最终结果
return err
}
// Server side - 逐块接收后聚合(用 return 而非 Fatal,让 gRPC 框架处理错误上报)
func (s *Server) Upload(stream pb.FileService_UploadServer) error {
var buffer bytes.Buffer
for {
chunk, err := stream.Recv()
if err == io.EOF {
break // client closed send direction
}
if err != nil {
return err
}
buffer.Write(chunk.Data)
}
_, err := stream.SendAndReceive(&pb.UploadResult{Size: int32(buffer.Len())})
return err
}
```
**优势:**内存友好——不需要一次性 load 全部数据到内存中,每个 chunk 独立收发。
## Bidirectional Streaming RPC(双向流)
Client 和 Server 可以同时独立地发送消息,是全双工通信。
**经典场景:**
- 聊天室
- 实时协作编辑
- 游戏状态同步
```protobuf
rpc Chat(stream ChatMessage) returns (stream ChatMessage);
```
```go
// Server side - 转发逻辑,两端各自独立循环
func (s *Server) Chat(stream pb.ChatService_ChatServer) error {
done := make(chan struct{})
// goroutine 1: 读取客户端消息
go func() {
for {
msg, err := stream.Recv()
if err == io.EOF {
close(done)
return
}
if err != nil {
return
}
// 广播给其他 connected clients...
s.broadcast(msg)
}
}()
// goroutine 2: 定时推送服务器消息
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
for {
select {
case <-done:
return nil
case t := <-ticker.C:
stream.Send(&pb.ChatMessage{Text: fmt.Sprintf("heartbeat: %s", t)})
}
}
}
```
**关键点:**
- 两端的 Send 和 Recv 是独立的——一端 Recv 完不影响另一端继续 Send
- 必须用两个 goroutine 分别处理 recv 和 send 循环
- `io.EOF` 只表示对方的关闭,不代表己方也要停止
### 背压与心跳(生产级要点)
BiDi Stream 在长连接场景下,有两个必须考虑的问题:
```go
// 背压控制:如果 Send 堆积过多,应该限流或暂停
func (s *Server) handleStream(stream pb.ChatService_ChatServer) error {
sendCh := make(chan *pb.ChatMessage, 100) // buffer size = 100
go func() {
for msg := range sendCh {
// Send 是阻塞的——buffer 满时自动背压
stream.Send(msg)
}
}()
// ...recv loop 往 sendCh 里塞消息即可
}
```
> [!tip] 背压原理
> gRPC 的 `Send()` 是**有缓冲阻塞**的。当 internal buffer 写满时,发送端会自动 pause——这就是 HTTP/2 Flow Control 提供的天然背压机制,不需要手动实现。但你应该设置合理的 buffer size,过大浪费内存,过小影响吞吐。
> [!note] Keepalive 配置示例
> ```go
> conn, _ := grpc.Dial(addr,
> grpc.WithKeepaliveParams(keepalive.ClientParameters{
> Time: 10 * time.Second, // ping interval
> Timeout: 20 * time.Second, // wait for ping ack
> PermitWithoutStream: true, // 即使无活跃 RPC 也发 ping
> }),
> )
> ```
> 这对穿越 Nginx / AWS ALB 等负载均衡器至关重要——它们通常会对空闲连接执行 tcp idle timeout 断开。
> [!warning] 复杂度警告
> BiDi Streaming 是最强大但也最容易出错的模式。你必须同时处理:context cancel、io.EOF、网络异常、心跳保活、背压(backpressure)。生产环境中除非必要,否则优先考虑其他三种模式。
## 模式选型决策指南
```mermaid
flowchart TD
Start{是否需要<br/>实时交互?}
Start -->|否| Simple{单次<br/>请求?}
Start -->|是| BiDi{高频<br/>交互?}
Simple -->|是| U[Unary RPC<br/>最简单]
Simple -->|否| SS[Server Stream<br/>一次请求多次返回]
BiDi -->|是| BD[Bidirectional Stream<br/>全双工通信]
BiDi -->|否| CS{数据量<br/>超大?}
CS -->|是| CB[Client Stream<br/>分批上传]
CS -->|否| SS
style U fill:#00D866,color:#fff
style BD fill:#FF6B35,color:#fff
```
## 性能对比
| 维度 | Unary | Server Stream | Client Stream | BiDi Stream |
|------|-------|---------------|---------------|-------------|
| **RTT** | 1 | 1+N | M+1 | M+N |
| **实现复杂度** | ⭐ | ⭐⭐ | ⭐⭐ | ⭐⭐⭐⭐ |
| **内存占用** | 高(一次加载完整响应) | 低(逐条处理) | 低(分片发送) | 低(双向流控) |
| **超时风险** | 高(连接全程持有一直到完整响应) | 低(请求已发出,可取消) | 中(大文件需合理 timeout) | 中(长连接需 keepalive) |
| **典型场景** | CRUD、短查询 | 订阅推送、列表分页 | 大文件上传、批量采集 | 聊天室、实时协作、游戏同步 |
## 与 HTTP 方法的映射类比
| gRPC 模式 | 近似的 HTTP 模式 |
|-----------|-----------------|
| Unary | GET / POST |
| Server Stream | SSE (Server-Sent Events) |
| Client Stream | Multipart Upload |
| BiDi Stream | WebSocket |
> [!note] 类比 ≠ 等价
> 这些只是功能层面的类比。gRPC 是二进制协议且基于 HTTP/2,行为特征和 HTTP 层语义有本质差异——比如 HTTP/2 的多路复用让 gRPC Stream 比 WebSocket 更高效。
## 关联笔记
- [[hhs/gRPC/2. gRPC 核心篇/06-Service 定义与代码生成]]
- [[hhs/gRPC/3. 服务端实现/09-Streaming Handler]]