Files
cs-note/hhs/gRPC/3. 服务端实现/09-Streaming Handler.md
T
2026-05-24 11:42:38 +08:00

522 lines
17 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, Go, Streaming, RPC]
create time: 2026-05-11 16:30
---
# Streaming Handler
## 概述
本文档讲解 gRPC 三种 streaming RPC handler 的编写模式——**Server Streaming**、**Client Streaming** 和 **Bidirectional Streaming**。核心主题是循环内的生命周期管理:如何正确感知 client 断开、如何在大数据量下避免 goroutine 泄漏、以及如何优雅地处理背压。
> [!warning] Streaming handler 最大的陷阱:goroutine 泄漏
> 如果你在 loop 里 spawn 了子 goroutine 但没有正确的退出机制,一旦 client 断开连接,server 端的 goroutine 会永远无法回收。**下面每一个例子都会标注如何避免这个问题。**
## Server Streaming Handler
Server streaming:client 发一次请求,server 持续返回多条响应。
```go
func (s *server) ListUsers(req *pb.ListUsersRequest, stream pb.UserService_ListUsersServer) error {
users := getAllUsers() // 从 DB / cache 获取
for _, u := range users {
if err := stream.Send(&pb.ListUsersResponse{User: u}); err != nil {
// client 已经断开或主动取消
s.logger.Warn("send failed", "error", err)
return status.Error(codes.Internal, err.Error())
}
}
return nil // 正常结束:数据全部发送完毕
}
```
### 核心要点
- Server 通过 `stream.Send()` 逐条发送,每次 Send 都会编码 Protobuf 并通过 HTTP/2 frame 推送到 client
- **`stream.Context()` 是唯一可靠的退出信号** —— 当 client cancel 或网络中断时,context 会被取消
- 正常结束(遍历完全部数据)返回 `nil`;断连时 context 报错
- 不要在 loop 之外做 cleanup,否则即使 send 失败也会执行不必要的清理
### 大数据量场景下的超时检测
上面的简单示例只适用于结果集可控的场景。如果用户量大(比如百万级),需要在循环中**主动检查 context** 提前退出:
```go
func (s *server) ListUsers(req *pb.ListUsersRequest, stream pb.UserService_ListUsersServer) error {
ctx := stream.Context()
results := hugeQuery(ctx) // 可能返回大量数据
for i, u := range results {
// 每发送 100 条检查一下 context
if i >= 100 && len(results) >= 1000 {
select {
case <-ctx.Done():
return ctx.Err() // client 已取消
default:
}
}
if err := stream.Send(&pb.ListUsersResponse{User: u}); err != nil {
return status.Error(codes.Internal, err.Error())
}
}
return nil
}
```
这里有个设计取舍:**每次 Send 都检查 vs 每隔 N 次检查**。全检查保证响应快但多了 select 开销;间隔检查减少 CPU 消耗但有延迟。根据业务对"取消感知的灵敏度"来决定。
### 数据流转时序
```mermaid
sequenceDiagram
participant C as Client
participant S as Server Handler
participant DB as Data Source
C->>S: Send(ListUsersRequest)
S->>DB: Query All Users
DB-->>S: []User
loop For Each User
alt Client Still Connected
S->>S: Check Context
S->>C: Send(User)
else Context Cancelled
S->>S: Return ctx.Err()
end
end
alt Normal Completion
C->>S: (no more data)
S-->>C: Return nil (EOF)
else Send Failed
S-->>C: Return error Status
end
```
## Client Streaming Handler
Client streaming:client 持续发送多条请求,server 汇总后返回一次响应。
```go
func (s *server) UploadData(stream pb.UserService_UploadDataServer) (*pb.UploadResult, error) {
ctx := stream.Context()
var totalBytes int64
var processedItems []*Item
for {
chunk, err := stream.Recv()
if err == io.EOF {
break // client 结束了发送
}
if err != nil {
// client 主动取消或网络异常 —— 不算业务错误,直接退出
if ctx.Err() != nil {
return nil, ctx.Err()
}
return nil, status.Errorf(codes.InvalidArgument, "recv failed: %v", err)
}
totalBytes += int64(len(chunk.Data))
processedItems = append(processedItems, chunk.ToItem())
// TODO: 批量写入 DB,不要每来一条就 flush
}
return &pb.UploadResult{TotalBytes: totalBytes, Count: int32(len(processedItems))}, nil
}
```
### 核心要点
- 用 `io.EOF` 判断 client 结束发送(注意:是 `io.EOF`,不是其他 error)
- **Recv 的 error 可能来自两种情况**:client 正常关闭(`io.EOF`)或 client 取消/断连。后者不应视为业务错误,而应直接通过 `ctx.Err()` 退出
- Recv 循环必须在 EOF 之前 `return` 或 `break`,否则会无限阻塞
- 所有临时状态放在 loop 内,返回前一次性聚合结果
### 内存安全注意事项
当 client 发送的数据量不可控时,**必须限制内存使用**:
```go
const maxBufferSize = 100 * 1024 * 1024 // 100 MB limit
func (s *server) UploadData(stream pb.UserService_UploadDataServer) (*pb.UploadResult, error) {
var totalBytes int64
var items []*Item // ⚠️ 生产环境应写流到磁盘而非放内存
for {
chunk, err := stream.Recv()
if err == io.EOF {
break
}
if err != nil {
return nil, status.Errorf(codes.InvalidArgument, "recv failed: %v", err)
}
totalBytes += int64(len(chunk.Data))
if totalBytes > maxBufferSize {
return nil, status.Errorf(codes.ResourceExhausted,
"exceeds max upload size: %d bytes", maxBufferSize)
}
items = append(items, chunk.ToItem())
}
return aggregateResult(items), nil
}
```
> [!tip] 何时该用临时文件?
> 如果单条消息超过 10MB 或者总量难以预估,直接将 `chunk.Data` 写到一个 `os.File` 中比存在内存切片里安全得多。handler 返回前一次性处理临时文件即可。
## Bidirectional Streaming Handler
Bidirectional streaming 是最复杂的模式:收和发两个方向完全独立,各自有自己的生命周期。
```go
func (s *server) ChatRoom(stream pb.UserService_ChatRoomServer) error {
ctx := stream.Context()
inboundCh := make(chan *pb.ChatMessage, 100)
errCh := make(chan error, 1) // 错误传播通道
// ---- 接收方向(后台 goroutine)----
go func() {
defer close(inboundCh)
for {
msg, err := stream.Recv()
if err != nil {
errCh <- err // EOF / cancel / network error
return
}
select {
case inboundCh <- msg:
default:
// channel full, drop message
}
}
}()
// ---- 发送方向(主循环)----
for {
select {
case <-ctx.Done():
return ctx.Err() // client 断了
case err := <-errCh:
return err // recv 出错,主循环退出
case envelope := <-inboundCh:
envelope := transform(envelope) // 业务转换
if err := stream.Send(envelope); err != nil {
return err
}
}
}
}
```
### 核心架构
```mermaid
flowchart TB
subgraph "Handler Goroutine"
S["Send Loop<br/>for-select"]
end
subgraph "Recv Goroutine"
R["Recv Loop"] -->|"inboundCh buffered"| B["Business Logic"]
end
R --> stream
S --> stream
B -->|"outboundCh buffered"| S
stream["HTTP/2 Stream"]
style S fill:#4FC08D,color:#fff
style R fill:#00B6BC,color:#fff
style B fill:#FF9F43,color:#000
```
### 关键设计原则
1. **Context.Done() 是唯一可靠的退出信号** — 不管哪个方向出错,最终都要通过 context 来协调
2. **收和发是两个独立的生命周期** — recv 在一个 goroutine,send 在主循环
3. **务必给 channel 设置 buffer**,否则 sender 堵住时 receiver 也会饿死
### 背压策略选择
channel full 时如何处理,决定了系统的韧性:
| 策略 | 实现方式 | 适用场景 | 风险 |
|------|---------|---------|------|
| **Drop** | `select + default`(如上文) | 实时性优先,丢消息可接受 | 丢失关键消息 |
| **Block** | 普通 `ch <- msg` | 吞吐量重要,允许等待 | receiver 被拖慢 |
| **Backpressure** | `select + ctx.Done()` | 需要优雅降级 | 代码略复杂 |
| **Reject** | 统计丢弃数,超标后返回 error | 有 SLA 要求的系统 | 直接中断连接 |
**推荐**大多数场景用 Backpressure 模式——超时未消费则丢弃,连续丢弃过多时通过 errCh 通知主循环终止:
```go
func (s *server) ChatRoom(stream pb.UserService_ChatRoomServer) error {
ctx := stream.Context()
inboundCh := make(chan *pb.ChatMessage, 100)
errCh := make(chan error, 1)
go func() {
defer close(inboundCh)
drops := 0
for {
msg, err := stream.Recv()
if err != nil {
errCh <- err
return
}
select {
case inboundCh <- msg:
case <-time.After(10 * time.Millisecond):
drops++
if drops > 1000 { // 连续丢弃过多,通知主循环终止
errCh <- status.Errorf(codes.ResourceExhausted,
"too many messages dropped")
return
}
}
}
}()
// 主循环同上,新增对 errCh 的监听即可
for {
select {
case <-ctx.Done():
return ctx.Err()
case err := <-errCh:
return err
case envelope := <-inboundCh:
if err := stream.Send(envelope); err != nil {
return err
}
}
}
}
```
> [!tip] bidirectional streaming 的设计直觉
> 想象成一个管道:一端进水(Recv)、一端出水(Send)。中间的业务逻辑可以是任意的转换、过滤、聚合。但管壁(context)不能漏。
## 错误处理对照表
streaming handler 中的错误来源复杂,不同错误的处理方式也不同:
| 错误来源 | stream.Recv 返回值 | 处理方式 |
|----------|-------------------|---------|
| Client 正常关闭连接 | `io.EOF` | `break` 或 `return nil` |
| Client 主动取消请求 | `context.Canceled`(包装为普通 error) | 记录日志,`break` |
| Network 中断 | 非 nil error(含 canceled) | 记录日志,`break` |
| 数据格式错误 | 具体 error | 原样返回给 client |
| 业务校验失败 | `status.Errorf(codes.InvalidArgument, ...)` | 返回具体错误码,终止循环 |
| 数据量超限 | `status.Errorf(codes.ResourceExhausted, ...)` | 返回错误码,终止循环 |
> [!note] 关于 Recv 返回 error 的细节
> gRPC-go 会把 `context.Canceled` 包装成普通 error 返回,不会返回 `io.EOF`。所以判断 client 断开应该同时检查两种 case。
```mermaid
flowchart TD
A["Recv() returns error?"] -->|"No"| B["继续处理消息"]
A -->|"Yes"| C{"err == io.EOF?"}
C -->|"Yes"| D["client 正常关闭<br/>return nil ✅"]
C -->|"No"| E{"codes.FromError is OK?"}
E -->|"Yes"| F["canceled / network issue<br/>return ctx.Err() ⚠️"]
E -->|"No"| G["真实错误<br/>return status.Errorf ❌"]
```
## Context 超时传播
Streaming handler 中,context 超时由 client 侧的 WithTimeout 控制,server 必须主动感知并退出:
```go
func (s *server) LongRunningOp(req *pb.LongRequest, stream pb.Service_LongRunningOpServer) error {
ctx := stream.Context()
for i := 0; i < 100; i++ {
// 每一步都检查 context
select {
case <-ctx.Done():
return ctx.Err() // client 取消了
default:
}
result := doStep(i)
if err := stream.Send(result); err != nil {
return err // send 失败也立即退出
}
time.Sleep(100 * time.Millisecond)
}
return nil
}
```
这个模式的精髓是 **"每一次 IO 操作前都检查 context"**。只要有一次遗漏,就可能变成孤儿 goroutine。
## 关闭与连接检测
在 streaming 场景中,理解"如何感知对端断开"比"如何发消息"更重要。
### Client 关闭的信号路径
```mermaid
flowchart LR
A["Client 正常关闭发送端"] -->|"Recv() → io.EOF"| B["handler return nil ✅"]
C["Client Cancel context"] -->|"Recv() → ctx.Err()"| D["handler return ctx.Err() ⚠️"]
E["网络断开"] -->|"Recv() → connection error"| F["handler return error ⚠️"]
G["不读 response"] -->|"Send() → stream error"| H["handler return error ❌"]
```
| Client 行为 | Server 感知现象 | handler 处理 |
|-------------|----------------|-------------|
| 正常关闭发送端 | `Recv()` 返回 `io.EOF` | `return nil` ✅ |
| Cancel context | `Recv()` 返回 context canceled | `return ctx.Err()` ⚠️ |
| 网络断开 | `Recv()` 返回连接相关 error | `return err` ⚠️ |
| 不读 response | `Send()` 返回 stream error | `return err` ❌ |
**核心原则**:无论哪种情况,handler 的正确做法都是尽快 return,让 gRPC 内部清理资源。不需要手动 close channel 或做其他清理。
### 为什么不存在 CloseSend 通知 channel
有些开发者误以为 gRPC 提供了类似 `CloseSend() chan struct{}` 的通知通道。实际上 **gRPC-Go 并没有这样的 API**。Server 只能通过以下两种方式检测 client 停止发送:
1. **`Recv() io.EOF`** —— client 正常调用 `CloseSend()` 后
2. **`Recv() error`** —— context 取消或网络异常
任何依赖"某个 channel 被关闭"来判断的模式都是不可靠的。
## 资源清理模式
Streaming handler 中经常需要订阅外部事件源(DB watcher、消息队列等),资源清理是关键:
```go
func (s *server) WatchEvents(req *pb.WatchRequest, stream pb.Service_WatchServer) error {
ctx := stream.Context()
// 订阅事件流 —— 必须传 context
ch, cancel := s.eventBroker.Subscribe(ctx, req.GetFilter())
defer cancel() // handler 返回时自动释放订阅
for {
select {
case <-ctx.Done():
return nil // client 断了
case event, ok := <-ch:
if !ok {
return nil // broker 关闭了 channel
}
if err := stream.Send(event); err != nil {
return err // send 失败
}
}
}
}
```
要点总结:
1. `Subscribe(ctx, ...)` 一定要传 context,不能用 `context.Background()`
2. `defer cancel()` 确保 handler 退出时释放订阅
3. select 中始终优先检测 `ctx.Done()`
```mermaid
sequenceDiagram
participant C as Client
participant H as Handler
participant EB as Event Broker
participant S as Stream
H->>EB: Subscribe(ctx, filter)
EB-->>H: ch, cancel
loop Until Disconnected
alt Client Still Connected
H->>H: select { ctx.Done(), <-ch }
EB->>H: event
H->>S: Send(event)
else Client Disconnected
H->>H: ctx.Done() triggers
H->>H: defer cancel()
H->>EB: Unsubscribe via cancel
H-->>C: Return nil
end
end
```
## 生产实践要点
写完 handler 只是第一步,生产环境还需要考虑以下方面:
### 背压与速率限制
- **Server Streaming** 中如果 client 读取速度慢于 server 发送速度,gRPC 会在 TCP/HTTP2 层积压数据。建议:监控 Send 耗时,如果 P99 持续升高则降低发送频率或加入 rate limiter
- **Client Streaming** 中可通过 `maxReceiveMessageSize` dial option 限制最大消息大小
### 监控指标
在 streaming handler 中埋点三个关键指标:
```go
// 1. handler 存活时长
duration := time.Since(start)
histogram.Record("grpc.stream.duration", duration.Milliseconds())
// 2. 发送/接收消息计数
counter.Increment("grpc.stream.messages_sent")
// 3. 非正常终止次数
counter.Increment("grpc.stream.errors", errorCode)
```
### 优雅关闭(Graceful Stop)
当 server 要 shutdown 时,不能直接 kill 正在执行的 streaming handler。正确流程:
1. 调用 `grpc.Server.GracefulStop()` —— 拒绝新连接
2. 等待已有 handler 自然返回(可以设 deadline)
3. 超时后强制 shutdown
```go
srv.GracefulStop() // 停止接收新请求
// 或带超时的优雅关闭:
stopCh := make(chan struct{})
go func() {
srv.GracefulStop()
close(stopCh)
}()
select {
case <-stopCh:
log.Println("all handlers finished")
case <-time.After(30 * time.Second):
srv.Stop() // 强制杀
log.Println("force stopped after timeout")
}
```
## 常见问题 Checklist
写完 streaming handler 后逐项核对:
- [ ] **Recv 错误分清了 `io.EOF` 和 `ctx.Err()` 吗?**(正常结束 vs client 取消的处理不同)
- [ ] Context 取消了吗?(`ctx.Done()` 在所有关键路径被检测)
- [ ] Send 失败正确处理了吗?(不会因为 send 错误继续无效循环)
- [ ] goroutine 会泄漏吗?(所有 spawn 的 goroutine 都有明确退出条件,bidirectional 中 recv goroutine 的错误能传播到主循环)
- [ ] Channel 有 buffer 吗?(bidirectional 模式下 sender/receiver 互相不饿死)
- [ ] 内存有上限吗?(client streaming 中对累积的消息数量做了限制)
- [ ] **消息大小有限制吗?**(gRPC 默认 4MB 限制,超出会直接报错)
- [ ] Subscribe 传了 context 吗?(外部事件源随 handler 一起退出)
- [ ] 大数据集做了间隔检测吗?(百万级数据不会无意义地每个元素都 select)
- [ ] 背压策略选对了吗?(drop/block/backpressure/reject 符合业务需求)
## 关联笔记
- [[hhs/gRPC/3. 服务端实现/08-Server 搭建与注册]]
- [[hhs/gRPC/3. 服务端实现/10-健康检查与反射]]
- [[hhs/gRPC/4. 客户端开发/13-Streaming Client]]
- [[hhs/gRPC/6. 工程实践篇/20-性能优化与压测]]