Files
2026-05-24 11:42:38 +08:00

9.1 KiB
Raw Permalink Blame History

tags, create time
tags create time
gRPC
Go
Streaming
ServerStreamingClient
ClientStreamingClient
BidiStreaming
Recv
Send
ContextCancellation
GracefulShutdown
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

服务端返回多条消息,客户端逐条接收:

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

客户端连续发送多条消息,服务端最后返回一条聚合结果:

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:

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 永远阻塞等待
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 次数可以显著提升吞吐量:

// 不好:每条消息都触发一次 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:

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 取消时的行为:

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 之上