522 lines
17 KiB
Markdown
522 lines
17 KiB
Markdown
---
|
||
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-性能优化与压测]]
|