--- 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
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 正常关闭
return nil ✅"] C -->|"No"| E{"codes.FromError is OK?"} E -->|"Yes"| F["canceled / network issue
return ctx.Err() ⚠️"] E -->|"No"| G["真实错误
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-性能优化与压测]]