252 lines
9.1 KiB
Markdown
252 lines
9.1 KiB
Markdown
---
|
||
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 之上
|