vault backup: 2026-05-13 11:17:21

This commit is contained in:
hhs
2026-05-13 11:17:21 +08:00
parent e8ee0c675c
commit 93daec5d10
6 changed files with 597 additions and 55 deletions
@@ -1,6 +1,7 @@
---
tags: [gRPC, RPC, Streaming, Go, Microservice]
tags: [gRPC, RPC, Streaming, Go, Microservice, API Design]
create time: 2026-05-11 16:40
update time: 2026-05-13 00:00
---
# RPC 调用模式总览
@@ -16,17 +17,18 @@ gRPC 提供四种 RPC 调用模式,从最简单的请求-响应到完全的双
```mermaid
flowchart LR
A[Unary] -->|"一问一答"| B[最简单]
C[Server Stream] -->|"一问多答"| D[广播式]
E[Client Stream] -->|"多问一答"| F[收集式]
G[BiDi Stream] -->|"多问多答"| H[全双工]
U["Unary\n一问一答"] --> S1["简单 · 阻塞 · 一次往返"]
SS["Server Stream\n一问多答"] --> S2["广播式 · 服务端推送"]
CS["Client Stream\n多问一答"] --> S3["收集式 · 分批上传"]
BD["BiDi Stream\n多问多答"] --> S4["全双工 · 独立收发"]
style A fill:#00B6BC,color:#fff
style G fill:#EE5A24,color:#fff
style U fill:#00B6BC,color:#fff
style SS fill:#00D866,color:#fff
style CS fill:#4FC3F7,color:#fff
style BD fill:#EE5A24,color:#fff
```
> [!example] 各模式数据流向速览
> 左列为 Client,右列为 Server,箭头方向表示数据流动方向。
## RPC 数据流示意图
```mermaid
sequenceDiagram
@@ -34,13 +36,13 @@ sequenceDiagram
participant S as Server
rect rgba(0, 182, 188, 0.1)
Note over C,S: Unary — 阻塞式一次往返
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 — 请求后连续响应
Note over C,S: Server Stream — 一次请求,多次响应
C->>S: request
S-->>C: response 1
S-->>C: response 2
@@ -48,7 +50,7 @@ sequenceDiagram
end
rect rgba(79, 195, 247, 0.1)
Note over C,S: Client Stream — 连续发送后一次性响应
Note over C,S: Client Stream — 多次请求,一次响应
C->>S: chunk 1
C->>S: chunk 2
C->>S: ... n (done)
@@ -142,27 +144,37 @@ func (s *Server) Subscribe(req *pb.SubscribeRequest, stream pb.UserService_Subsc
> 每个文件只需 import 实际用到的即可,不必照抄。
```go
// Client side - 遍历接收事件(生产环境建议用 recover + defer 做错误恢复)
func main() {
client := pb.NewUserServiceClient(conn)
stream, err := client.Subscribe(context.Background(), &pb.SubscribeRequest{Topic: "orders"})
// Client side - 遍历接收事件(生产环境建议加 context timeout 和 recover)
func subscribeEvents(client pb.UserServiceClient, topic string) error {
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
stream, err := client.Subscribe(ctx, &pb.SubscribeRequest{Topic: topic})
if err != nil {
log.Fatal(err)
return err // dial failed — 上层 main 可以做 Fatal
}
defer func() {
if r := recover(); r != nil {
log.Printf("panic recovered: %v", r)
}
}()
for {
event, err := stream.Recv()
if err == io.EOF {
break // server finished sending
}
if err != nil {
return err // ⚠️ 这里用 return 而非 log.Fatal,服务端的 handler 不能 kill 进程
log.Printf("stream error: %v", err) // ⚠️ 这里用 return/log,不能用 Fatal
return err
}
fmt.Printf("received: %s\n", event.Data)
}
return nil
}
```
> [!tip] 注意:Client 仍然可以通过 context cancel 随时中断流。这是调试流式问题时最容易忽略的一点——不是 Server 主动关了连接,而是 Client 放弃了。
> [!tip] Context Cancel 与 Recv 的关系
> Client 仍可通过 cancel 随时中断流——服务端 `Recv()` 将返回一个 context canceled 错误。这也是调试流式问题时最容易忽略的一点:**不是 Server 主动关了连接,而是 Client 放弃了**。
## Client Streaming RPC(客户端流)
@@ -178,19 +190,31 @@ rpc Upload(stream FileChunk) returns (UploadResult);
```
```go
// Client side - 流式发送数据分片
// Client side - 流式发送数据分片(注意 defer CleanupSend 做资源清理)
func uploadFile(client pb.FileServiceClient, chunks [][]byte) error {
stream, err := client.Upload(context.Background())
ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
defer cancel()
stream, err := client.Upload(ctx)
if err != nil {
return err
}
defer func() {
if r := recover(); r != nil {
stream.CloseSend() // panic 时确保 send 方向关闭
}
}()
for _, chunk := range chunks {
if err := stream.Send(&pb.FileChunk{Data: chunk}); err != nil {
return err
}
}
result, err := stream.CloseAndRecv() // 结束发送,获取最终结果
return err
if err != nil {
return err
}
fmt.Printf("upload complete, size: %d\n", result.Size)
return nil
}
// Server side - 逐块接收后聚合(用 return 而非 Fatal,让 gRPC 框架处理错误上报)
@@ -206,13 +230,34 @@ func (s *Server) Upload(stream pb.FileService_UploadServer) error {
}
buffer.Write(chunk.Data)
}
_, err := stream.SendAndReceive(&pb.UploadResult{Size: int32(buffer.Len())})
result, err := stream.SendAndReceive(&pb.UploadResult{Size: int32(buffer.Len())})
return err
}
```
> [!tip] CloseAndRecv vs CloseSend + Recv
> `CloseAndRecv()` 是 Client Stream 中的便捷方法——它同时完成"关闭发送方向"和"读取响应"两步。等价于先 `stream.CloseSend()` 再 `stream.Recv()`。在 Client Stream 中两者效果相同,但 `CloseAndRecv` 更简洁、出错概率更低。
**优势:**内存友好——不需要一次性 load 全部数据到内存中,每个 chunk 独立收发。
## 错误处理模式总结
四种模式共用同一套错误处理原则:
| 场景 | Unary | Server Stream | Client Stream | BiDi Stream |
|------|:-----:|:-------------:|:-------------:|:-----------:|
| **Stream 创建失败(Client)** | — | `log.Fatal` ✅ | `log.Fatal` ✅ | `log.Fatal` ✅ |
| **Recv 返回 io.EOF** | — | `break` ✅ | `break` ✅ | `close(done)` ✅ |
| **Recv 返回其他 error** | `return/log` ✅ | `return/log` ✅ | `return/log` ✅ | 分别处理两端 ✅ |
| **Send 返回 error** | — | `return` ✅ | `return` ✅ | `return` ✅ |
| **Context canceled** | RPC 自动取消 | 同 error ✅ | 同 error ✅ | `select <-ctx.Done()` ✅ |
> [!important] 黄金法则
> 1. **只有连接建立阶段的 dial/send 错误才能用 `log.Fatal`**——handler 内部永远用 return
> 2. **io.EOF 不是错误**——它表示对方完成了发送方向,是正常退出信号
> 3. **每次 stream recv/send 都要检查 error**,哪怕代码看起来"不可能失败"
> 4. **context timeout 对所有模式生效**——Streaming 不会因为"流式"就自动获得更长超时
## Bidirectional Streaming RPC(双向流)
Client 和 Server 可以同时独立地发送消息,是全双工通信。
@@ -227,8 +272,9 @@ rpc Chat(stream ChatMessage) returns (stream ChatMessage);
```
```go
// Server side - 转发逻辑,两端各自独立循环
// Server side - 转发逻辑,两端各自独立循环(增加 context cancel 处理)
func (s *Server) Chat(stream pb.ChatService_ChatServer) error {
ctx := stream.Context()
done := make(chan struct{})
// goroutine 1: 读取客户端消息
@@ -240,6 +286,7 @@ func (s *Server) Chat(stream pb.ChatService_ChatServer) error {
return
}
if err != nil {
log.Printf("recv error: %v", err)
return
}
// 广播给其他 connected clients...
@@ -252,32 +299,41 @@ func (s *Server) Chat(stream pb.ChatService_ChatServer) error {
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return ctx.Err() // client disconnected → clean exit
case <-done:
return nil
case t := <-ticker.C:
stream.Send(&pb.ChatMessage{Text: fmt.Sprintf("heartbeat: %s", t)})
if sendErr := stream.Send(&pb.ChatMessage{Text: fmt.Sprintf("heartbeat: %s", t)}); sendErr != nil {
log.Printf("send heartbeat error: %v", sendErr)
return sendErr
}
}
}
}
```
**关键点:**
- 两端的 Send 和 Recv 是独立的——一端 Recv 完不影响另一端继续 Send
- 必须用两个 goroutine 分别处理 recv 和 send 循环
- `io.EOF` 只表示对方的关闭,不代表己方也要停止
- 两端的 Send 和 Recv 是独立的——一端 Recv 完(EOF)不影响另一端继续 Send
- 必须用两个 goroutine 分别处理 recv 和 send 循环——单线程无法同时读写
- `io.EOF` 只表示对方的关闭,不代表己方也要停止发送
- **Context Cancel 优先级高于 EOF**:Client 断开连接时,`stream.Context().Done()` 先触发。生产代码中永远要检查它
- 建议在 handler 中使用 `defer stream.SendAndClose(...)` 或 defer 清理逻辑确保资源释放
### 背压与心跳(生产级要点)
BiDi Stream 在长连接场景下,有两个必须考虑的问题:
BiDi Stream 在长连接场景下,有三个必须考虑的问题:context cancel、io.EOF、背压和心跳保活。
```go
// 背压控制:如果 Send 堆积过多,应该限流或暂停
// 背压控制:如果 Send 堆积过多,buffer 满时自动阻塞发送端
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)
// Send 是阻塞的——buffer 满时自动触发背压
if err := stream.Send(msg); err != nil {
return // connection lost, goroutine exits cleanly
}
}
}()
// ...recv loop 往 sendCh 里塞消息即可
@@ -285,7 +341,7 @@ func (s *Server) handleStream(stream pb.ChatService_ChatServer) error {
```
> [!tip] 背压原理
> gRPC 的 `Send()` 是**有缓冲阻塞**的。当 internal buffer 写满时,发送端会自动 pause——这就是 HTTP/2 Flow Control 提供的天然背压机制,不需要手动实现。但你应该设置合理的 buffer size,过大浪费内存,过小影响吞吐。
> gRPC 的 `Send()` 是**有缓冲阻塞**的。当 internal buffer 写满时,发送端会自动 pause——这就是 HTTP/2 Flow Control 提供的天然背压机制,不需要手动实现。但你应该设置合理的 buffer size,过大浪费内存,过小影响吞吐。建议从 100 起步,根据实际监控调整。
> [!note] Keepalive 配置示例
> ```go
@@ -293,35 +349,41 @@ func (s *Server) handleStream(stream pb.ChatService_ChatServer) error {
> grpc.WithKeepaliveParams(keepalive.ClientParameters{
> Time: 10 * time.Second, // ping interval
> Timeout: 20 * time.Second, // wait for ping ack
> PermitWithoutStream: true, // 即使无活跃 RPC 也发 ping
> PermitWithoutStream: true, // 即使无活跃 RPC 也发 ping
> }),
> )
> ```
> 这对穿越 Nginx / AWS ALB 等负载均衡器至关重要——它们通常会对空闲连接执行 tcp idle timeout 断开。
> 这对穿越 Nginx / AWS ALB 等负载均衡器至关重要——它们通常会对空闲连接执行 TCP idle timeout 断开。服务端也需要配置类似的 keepalive,否则 Server→Client 方向的心跳缺失会导致 Client 误判连接死亡。
> [!warning] 复杂度警告
> BiDi Streaming 是最强大但也最容易出错的模式。你必须同时处理:context cancel、io.EOF、网络异常、心跳保活、背压(backpressure)。生产环境中除非必要,否则优先考虑其他三种模式。
> [!question] 为什么需要 PermitWithoutStream?
> 因为某些场景中连接处于"空闲状态"——没有正在进行的 RPC 调用——此时 LB 会因为检测到 TCP 层无任何流量而主动断开连接。设置 `PermitWithoutStream: true` 确保即使没有活跃流,gRPC 仍会持续发送 ping 包维持连接。
## 模式选型决策指南
```mermaid
flowchart TD
Start{是否需要<br/>实时交互?}
Start -->|否| Simple{单次<br/>请求?}
Start -->|是| BiDi{高频<br/>交互?}
Start["是否需要\n实时交互?"] -->|否| Simple["单次请求?"]
Start -->|是| BiDi["高频交互?"]
Simple -->|是| U[Unary RPC<br/>最简单]
Simple -->|否| SS[Server Stream<br/>一次请求多次返回]
Simple -->|是| U["Unary RPC\n最简单 ⭐"]
Simple -->|否| SS["Server Stream\n一次请求多次返回"]
BiDi -->|是| BD[Bidirectional Stream<br/>全双工通信]
BiDi -->|否| CS{数据量<br/>超大?}
CS -->|是| CB[Client Stream<br/>分批上传]
BiDi -->|是| BD["Bidirectional Stream\n全双工通信 🔥"]
BiDi -->|否| CS["数据量超大?"]
CS -->|是| CB["Client Stream\n分批上传"]
CS -->|否| SS
style U fill:#00D866,color:#fff
style BD fill:#FF6B35,color:#fff
```
> [!tip] 选型原则
> **默认选 Unary**。只有在三种情况下考虑其他模式:(1) 需要实时推送——用 Server Stream;(2) 数据量大且需流式发送——用 Client Stream;(3) 双向高频交互——用 BiDi Stream。永远不要为了炫技而选择更复杂的模式。
## 性能对比
| 维度 | Unary | Server Stream | Client Stream | BiDi Stream |