Files

252 lines
9.1 KiB
Markdown
Raw Permalink Normal View History

2026-05-24 11:42:38 +08:00
---
tags: [gRPC, Go, Streaming, ServerStreamingClient, ClientStreamingClient, BidiStreaming, Recv, Send, ContextCancellation, GracefulShutdown]
create time: 2026-05-11 15:34
---
# Streaming Client
## 概述
客户端处理 stream 的方式跟 unary 截然不同。你不能简单地 `resp, err := client.Method()`,而是需要处理一个连续的收发消息循环。这里我们逐个拆解每种 stream 模式下客户端的正确写法。
> [!tip] Stream 不是魔法
> 本质上一根 HTTP/2 stream 就是一个 duplex channel。gRPC 只是把 Send 和 Recv 的时序通过代码组织起来。理解了这一点,就能理解所有 stream 变体的模式。
## Server Streaming Client
服务端返回多条消息,客户端逐条接收:
```go
stream, err := client.ListUsers(ctx, &pb.ListUsersRequest{})
if err != nil {
return err // call 启动失败
}
for {
resp, err := stream.Recv()
if err == io.EOF {
break // 正常结束
}
if err != nil {
return err // 真实错误
}
fmt.Println(resp.User)
}
```
要点:
1. **调用时返回 stream object**,而非单个 response
2. **Recv loop 直到收到 `io.EOF` 才算正常结束**,EOF 是服务端主动关闭发送端的信号
3. **全程只有两个 error 关注点**:initial error(call 启动时)和 final error(EOF 前),中间 Recv 成功不需要检查 err
这种模式适合列表类接口——数据量可能很大但服务端控制发送节奏,客户端不需要主动推数据。
## Client Streaming Client
客户端连续发送多条消息,服务端最后返回一条聚合结果:
```go
stream, err := client.UploadData(ctx)
if err != nil {
return err
}
chunks := splitIntoChunks(data)
for _, chunk := range chunks {
if err := stream.Send(&pb.Chunk{Data: chunk}); err != nil {
return err // 某条发送失败
}
}
result, err := stream.CloseAndRecv() // 关闭发送端并获取最终响应
if err != nil {
return err
}
fmt.Printf("uploaded %d bytes\n", result.TotalBytes)
```
要点:
1. **`Send` 循环中每次调用都可能报错**(网络断开、context 取消、服务端关闭等)
2. **`CloseAndRecv` 一步完成两件事**:关闭发送端 + 接收最终响应——减少一次 round-trip
3. **等价写法**是分开两步:先 `stream.CloseSend()` 再 `stream.Recv()`,但 CloseAndRecv 更简洁且语义更清晰
> [!question] CloseSend vs CloseAndRecv 怎么选?
>
> | 方法 | 语义 | 适用场景 |
> |------|------|---------|
> | `CloseAndRecv()` | 关发送 + 收响应,一步完成 | **大多数 client streaming 场景(推荐)** |
> | `CloseSend()` + `Recv()` | 分开执行,两步骤 | 需要在上一步之间做自定义逻辑(如日志、指标采集) |
`CloseAndRecv` 的优势在于原子性:从客户端视角看,"关闭发送端"和"获取响应"是一个操作。如果用两步写,两步之间存在微小的时间窗口——服务端可能在这期间发了响应但客户端还没调 `Recv()`,导致时序上的不确定性。虽然实际影响极小,但原子操作的语义更不容易出错。
## Bidirectional Streaming Client
双方互相发消息,收发通常是两个 goroutine:
```go
stream, err := client.ChatRoom(ctx)
if err != nil {
return err
}
done := make(chan struct{})
// recv goroutine
go func() {
defer close(done)
for {
msg, err := stream.Recv()
if err == io.EOF {
return
}
if err != nil {
log.Printf("recv error: %v", err)
return
}
display(msg.Text)
}
}()
// send goroutine
go func() {
for text := range userInputCh {
if err := stream.Send(&pb.Message{Text: text}); err != nil {
stream.CloseSend()
return
}
}
}()
<-done // 等待 recv 结束
```
要点:
1. **gRPC stream object 是线程安全的**——多个 goroutine 并发调用 Send/Recv 无需额外加锁
2. **ctx 取消会同时影响两端**——Send 和 Recv 都会立刻收到 context error,需要统一处理
3. **send goroutine 退出时务必调用 `CloseSend()`**——通知服务端发送端已关闭,否则服务端 Recv 永远阻塞等待
```mermaid
flowchart TB
subgraph "Client"
A["send goroutine<br/>Send loop"] -->|"HTTP/2 stream"| D["Server Handler"]
E["recv goroutine<br/>Recv loop"]
end
subgraph "Server"
D -->|"HTTP/2 stream"| F["reply Send"]
G["server Recv"]
end
F --> E
A -.->|"CloseSend when done"| D
style A fill:#00D866,color:#fff
style E fill:#4FC08D,color:#fff
style D fill:#FF9F43,color:#000
style F fill:#EE5A24,color:#fff
```
双向流的特殊性在于:**发送和接收是两个独立的 flow**。任何一方都可以随时发送而不阻塞对方,这也是为什么需要两个 goroutine 来管理生命周期。
## 错误分类速查表
| 错误类型 | 表现 | 处理方式 |
|----------|------|---------|
| `io.EOF` | Recv 返回 EOF | break loop,正常结束 |
| `context.Canceled` | Recv 返回 context error | 检查是否主动 cancel |
| `codes.Unavailable` | Recv/Send 返回 unavailable | 可能需要重试或 reconnect |
| `codes.ResourceExhausted` | Send 返回 resource exhausted | 背压 / 限流 / 暂停发送 |
| `codes.Internal` | Recv 返回 internal | 通常是服务端 bug,记录日志 |
当 stream 中出现非 EOF 错误后:
1. 后续所有 Send/Recv 会立刻返回同一个错误
2. 建议直接退出当前函数或返回上层处理
3. 不要尝试在同一个 stream 上恢复操作
## 批量发送优化
对于 client streaming 场景,减少 syscall 次数可以显著提升吞吐量:
```go
// 不好:每条消息都触发一次 syscall
for item := range items {
stream.Send(item) // N 次 syscall
}
// 好:聚合后批量发送
var buf []*pb.Item
for item := range items {
buf = append(buf, item)
if len(buf) >= 64 { // 批次大小按场景调优
stream.Send(&pb.Batch{Items: buf})
buf = buf[:0]
}
}
if len(buf) > 0 {
stream.Send(&pb.Batch{Items: buf})
}
```
另一种思路是利用 Protobuf 的 `repeated` 字段在服务端一次解包多条消息。将业务层的 batching 逻辑与协议层的设计对齐,可以减少网络往返次数并降低 CPU 开销。
## Context 取消与优雅退出
流式 RPC 的上下文是一个 `context.Context`——它的取消会**同时影响** Send 和 Recv:
```go
stream, err := client.ChatRoom(ctx)
if err != nil {
return err
}
// ctx cancel 后,Recv/Send 都会报错
defer stream.CloseSend() // ⚠️ defer 放这里而非 goroutine 内
```
> [!warning] Common Pitfall
> 在 bidirectional streaming 中,很多开发者会把 `CloseSend()` 放在 send goroutine 里——但如果函数先因其他原因返回,send goroutine 可能还没执行到 `CloseSend()`。**把 `defer stream.CloseSend()` 放在主函数的开头是最安全的做法**:它保证无论哪种路径退出,都会关闭发送端。
ctx 取消时的行为:
```mermaid
flowchart TD
A["ctx Cancelled"] --> B["Send 立即报错<br/>context.Canceled"]
A --> C["Recv 返回错误<br/>context.Canceled"]
A --> D["Server Handler 收到<br/>context.Canceled"]
B --> E["清理资源<br/>CloseSend + 退出"]
C --> E
D --> F["停止处理并 Return"]
style E fill:#00B6BC,color:#fff
style F fill:#FF9F43,color:#000
```
要点:
1. **不要在 Recv 循环中吞掉 `context.Canceled`**:它是主动取消的信号,不是网络异常,不需要重试
2. **`defer CloseSend()` 要放在主函数**,不在子 goroutine 里——确保一定被执行
3. **服务端收到客户端的 `CloseSend` 后**(即 `Recv` 返回 `io.EOF`),应正常结束 handler
> [!question] 如果 ctx 超时了,已经发出去的消息会不会丢?
> 不会。gRPC 底层走的是 HTTP/2,消息一旦写入操作系统内核 buffer 就算已发出。ctx 取消只是通知本地 gRPC 库:"不要再接收或发送新数据",但已在途的消息不受影响。
## 最佳实践速查
在实际项目中,这份 checklist 能帮你避开大部分坑:
1. **永远带 context**:每个 `client.Xxx(ctx, ...)` 的 ctx 应该携带 deadline 或 cancel,避免流永远挂起。
2. **defer CloseSend() 放主函数**:bidirectional stream 中确保退出路径一定关闭发送端。
3. **不要吞掉 context.Canceled**:主动取消不需要重试——重试只会制造重复流量。
4. **Recv 后检查 EOF,再检查 error**:EOF 是正常结束信号,不是错误,顺序不能反。
5. **goroutine 生命周期配对**:每开一个 recv goroutine 就要有一个 done channel 对应收取完毕。
> [!tip] 进阶方向
> - 需要重试机制?看 [[hhs/gRPC/4. 客户端开发/12-Call Options 与 Context]] 中的 retry policy
> - 需要拦截所有 stream 日志和指标?看 [[hhs/gRPC/5. 中间件与拦截器/16-日志与链路追踪]]
> - 想了解服务端对应写法?看 [[hhs/gRPC/3. 服务端实现/09-Streaming Handler]]
## 关联笔记
- [[hhs/gRPC/2. gRPC 核心篇/05-RPC 调用模式总览]] — 四种 RPC 调用模式全景对比
- [[hhs/gRPC/4. 客户端开发/11-Client 连接与 Dial]] — Stream 建立在 client connection 之上