gRPC 四种流式通信模型:Unary、Server/Client Streaming 与 Bidirectional 实战

深入讲解 gRPC 四种通信模型:Unary 一元调用、Server Streaming 服务端流、Client Streaming 客户端流、Bidirectional Streaming 双向流。结合 Go 代码演示断线重连、流控 backpressure、拦截器、错误处理,并给出流式 gRPC 在生产环境的选型与避坑。

导语:gRPC 不止有一问一答

很多开发者以为 gRPC 只是"更快的 REST",是 HTTP/2 上的 JSON 换 Protobuf。但实际上 gRPC 最强大的能力来自 HTTP/2 的多路复用与全双工流式传输——它支持四种 RPC 调用模型,能承载从简单查询到实时数据推送、流式上传下载的完整 spectrum。

一句话总结:gRPC 真正的威力在"流"——四模型里,双向流是实时交互的顶配,但每个流都要放在 HTTP/2 + 服务端逻辑的约束下设计。


1. 四种通信模型总览

1.1 Proto 定义

syntax = "proto3";

package chat;

service ChatService {
    // 1. Unary:一元调用,请求-响应各一个
    rpc SendMessage (Message) returns (MessageReply);

    // 2. Server Streaming:客户端发一个,服务端回一堆(推送/订阅)
    rpc Subscribe (SubscribeReq) returns (stream Event);

    // 3. Client Streaming:客户端发一堆,服务端回一个(上传/聚合)
    rpc Upload (stream Chunk) returns (UploadReply);

    // 4. Bidirectional Streaming:两端都流式,全双工
    rpc Chat (stream ChatMsg) returns (stream ChatMsg);
}

message Message { string text = 1; }
message MessageReply { string reply = 1; }
message SubscribeReq { string topic = 1; }
message Event { bytes data = 1; }
message Chunk { bytes payload = 1; }
message UploadReply { int64 size = 1; }
message ChatMsg { string sender = 1; string text = 2; }

1.2 四模型对比

模型客户端发送服务端响应典型场景
Unary1 个请求1 个响应普通 RPC、查询、写操作
Server Streaming1 个请求多个响应推送、订阅、监听、分页拉流
Client Streaming多个请求1 个响应大文件上传、批量聚合、日志收集
Bidirectional多个请求多个响应聊天、实时协作、双向数据流

一句话总结:选型看数据流向——客户端反复要结果用服务端流,客户端批量提交用客户端流,两边都活跃用双向流,否则就用最简单的 Unary。


2. Unary 与服务端流(最常用)

2.1 服务端流式拉取

// 服务端:订阅推送(stream)
func (s *chatServer) Subscribe(req *SubscribeReq, stream chat.ChatService_SubscribeServer) error {
    events := make(chan *Event)   // 假设业务从这里拿到事件
    defer close(events)

    go s.produceEvents(req.Topic, events)

    for ev := range events {
        // 每次 Send 发送一个流消息
        if err := stream.Send(ev); err != nil {
            return err            // 客户端断开/网络异常
        }
    }
    return nil
}

2.2 客户端消费

func subscribe() error {
    conn, _ := grpc.NewClient("localhost:50051", grpc.WithTransportCredentials(...))
    defer conn.Close()
    client := chat.NewChatServiceClient(conn)

    stream, err := client.Subscribe(ctx, &SubscribeReq{Topic: "stock.tsla"})
    if err != nil { return err }

    for {
        ev, err := stream.Recv()
        if errors.Is(err, io.EOF) {
            break                      // 服务端正常结束
        }
        if err != nil { return err }   // 异常断开
        handleEvent(ev)
    }
    return nil
}

一句话总结:服务端流的本质是"一个请求、循环 Send",客户端循环 Recv 直到 EOF——推送、订阅、大表分页都能用它实现。


3. 客户端流与服务端流接收

3.1 客户端流式上传

// 服务端
func (s *ClientServer) Upload(stream chat.ChatService_UploadServer) error {
    var total int64
    for {
        chunk, err := stream.Recv()
        if errors.Is(err, io.EOF) {
            break                     // 客户端发完
        }
        if err != nil { return err }
        total += int64(len(chunk.Payload))
        // 逐块写入存储……
    }
    return stream.SendAndClose(&UploadReply{Size: total})
}

// 客户端
stream, _ := client.Upload(ctx)
// 循环发
for _, chunk := range fileChunks(file) {
    if err := stream.Send(&Chunk{Payload: chunk}); err != nil {
        return err
    }
}
// 发完关闭发送端,等待服务端回执
reply, err := stream.CloseAndRecv()
_ = reply.Size

一句话总结:客户端流" 客户端循环 Send、结束 CloseAndRecv、服务端 Recv 到 EOF 再 SendAndClose"——是大文件上传和批量提交的标配。


4. 双向流:实时对话的顶配

4.1 双向聊天

// 服务端:双向流
func (s *ChatServer) Chat(stream chat.ChatService_ChatServer) error {
    for {
        in, err := stream.Recv()      // 收一条客户端消息
        if errors.Is(err, io.EOF) {
            return nil
        }
        if err != nil { return err }

        // 处理并回复(可并发 Send)
        if err := stream.Send(&ChatMsg{
            Sender: "server",
            Text:  "收到: " + in.Text,
        }); err != nil {
            return err
        }
    }
}

4.2 双向流的并发约束

关键点 1:Send 与 Recv 由同一 stream 对象完成,但底下的 Session 是 HTTP/2 全双工。

关键点 2:Recv 和 Send 可以各自在一个 goroutine 里独立跑,互不阻塞。

关键点 3:要控制流控背压——如果服务端 Send 太快、客户端 Recv 太慢,
  底层 TCP 缓冲会积压,最终触发阻塞,需要配合并发限制/令牌桶。

典型模式:
  - 消息 → worker 池(限并发)
  - worker 结果 → channel
  - 单独发送 goroutine 从 channel 取出 Send

一句话总结:双向流让 send/receive 双工并行,但真正的挑战是流控——用 channel + worker 池平衡收发速率,避免背压死锁。


5. 拦截器、错误处理与重连

5.1 拦截器(Interceptor)

// 服务端一元拦截器:统一鉴权 / 日志 / 限流
func unaryAuthInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
    md, ok := metadata.FromIncomingContext(ctx)
    if !ok { return nil, status.Error(codes.Unauthenticated, "no metadata") }
    if token := md.Get("authorization"); len(token) == 0 {
        return nil, status.Error(codes.Unauthenticated, "missing token")
    }
    // 鉴权通过,进入实际 handler
    return handler(ctx, req)
}

// 注册
grpc.NewServer(grpc.UnaryInterceptor(unaryAuthInterceptor))

5.2 错误用 gRPC status 表达

// 返回结构化错误码(客户端能精确处理)
return nil, status.Error(codes.NotFound, "user id not found")

// 客户端侧判断
if st, ok := status.FromError(err); ok {
    switch st.Code() {
    case codes.NotFound:
        // 未找到逻辑
    case codes.Unavailable:
        retry()   // 可重试
    }
}

5.3 流式重连(生产关键)

// gRPC 客户端默认带自动重连(当 resolver 配置时)
// 但流式连接中断后,需要手动重建 stream:
for {
    stream, err := newChatStream(ctx)
    if err != nil {
        time.Sleep(backoff())   // 指数退避
        continue
    }
    err = consumeStream(stream) // 阻塞,断开时返回错误
    // 记录重连次数、保留进度、指数退避
}

一句话总结:拦截器做横切(鉴权/限流/日志),status 编码错误语义,流式重连配指数退避循环——三者是生产级 gRPC 的护栏。


6. 生产选型与避坑

场景推荐模型注意
常规 CRUD / 查询Unary最简单,带超时
实时行情/事件订阅Server Streaming客户端自动重连 + 断点续传
大文件/批量上传Client Streaming分块 + backpressure
实时聊天/协作Bidirectional并发控制 + 心跳 Ping
海量日志上报Client Streaming批量 + 限流

生产避坑清单:

□ Protobuf 字段变更是非破坏的,但别复用字段号
□ 流式接口一定设超时/心跳,防 TCP 半开连接
□ 别把大对象一次性放内存,改用流式传
□ 记住 Recv 返回 io.EOF 才是正常结束
□ 用 load-balancer / resolver 做到连接级复用
□ 服务端 Send 阻塞时用并发 worker 避免卡死
□ 监控流式打开时间长度和消息速率,防泄漏

7. 总结

gRPC 四模型从简到繁覆盖了几乎全部通信需求:

数据流发起方典型关键 API
一元客户端查询Call / Recv
服务端流客户端一次服务端多次推送Send(server) / Recv(client, until EOF)
客户端流客户端多次服务端一次上传Send(client) / CloseAndRecv
双向流两边多次聊天双 for 循环 + 并发控制

终极建议:能用 Unary 别上流(简单优先),真需要实时再升级到流式,双向流一定要配好并发、背压和重连三个护栏。把 HTTP/2 的复用优势吃掉,你的 RPC 层才算真正用好了 gRPC 这把利器。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「golang」更多文章

  1. Go defer 与常见陷阱深度:闭包捕获、执行顺序、性能与资源管理
  2. Go database/sql 实战:连接池、事务、批量插入与常见坑
  3. Go HTTP/2 连接池与复用实战:Transport 调优、Keep-Alive、连接泄漏排查