在gRPC的四种通信模式中,服务端流处理指的是客户端发送一次请求,服务器可以连续返回多条消息,直到主动结束流。Go语言通过protobuf编译器生成的代码,把这种模式变成了接口方法里的stream返回值,开发者只需要实现发送循环即可。这种方式非常适合数据量不定、需要逐步推送结果的业务,比如日志订阅、批量任务进度通知。

一、定义服务端流的proto文件
要实现服务端流,首先要在proto文件中使用stream关键字标记返回类型。这意味着该方法不会一次性返回单个消息,而是返回一个流对象,服务端可以多次调用Send方法推送数据。下面给出一个最小可用的示例,定义了一个订阅日志的接口。
注意service中 returns 前面的 stream 不能省略,否则生成的是普通一元RPC。同时请求参数可以是空消息,也可以带过滤条件,这取决于业务设计。下面代码中客户端发一个主题名,服务端不断推送该主题的日志内容。
syntax = "proto3";
package logpb;
option go_package = "ippipp.com/logpb;logpb";
service LogService {
// 服务端流:请求一次,返回多条日志
rpc Subscribe(LogRequest) returns (stream LogResponse);
}
message LogRequest {
string topic = 1;
}
message LogResponse {
string line = 1;
int64 ts = 2;
}
二、生成Go代码与接口理解
使用protoc-gen-go和protoc-gen-go-grpc插件编译上述proto后,会得到一个LogServiceServer接口,其中Subscribe方法签名类似:Subscribe(*LogRequest, LogService_SubscribeServer) error。第二个参数是一个流发送接口,它内嵌了grpc.ServerStream,并提供了Send(*LogResponse) error方法。
生成的LogService_SubscribeServer接口只负责服务端往客户端写数据,客户端通过LogService_SubscribeClient的Recv()方法读取。理解这一点很关键:服务端流不是双向聊天,而是单方向的持续下发。如果业务需要客户端也发多条消息,那就应该用双向流而不是服务端流。
// 生成代码中的关键接口(示意)
type LogService_SubscribeServer interface {
Send(*LogResponse) error
grpc.ServerStream
}
type LogServiceServer interface {
Subscribe(*LogRequest, LogService_SubscribeServer) error
}
三、Golang服务端实现
在服务端,我们需要实现Subscribe方法。核心逻辑是:根据请求中的topic,循环构造消息并调用Send。为了能随时停止,通常要监听context的Done信号,避免客户端断开后服务端还在傻发。下面示例用time.Ticker模拟日志产生,每500毫秒推一条。
这里有几个实践细节。第一,Send方法在网络异常时会返回错误,应当直接return error让gRPC框架处理断开。第二,如果循环里做了阻塞操作,一定要用select兼顾ctx.Done(),否则客户端取消订阅时服务端协程会泄漏。第三,不要在流里返回业务错误码后用Send继续发,应该结束方法。
package main
import (
"context"
"log"
"time"
"ippipp.com/logpb"
)
type server struct{}
func (s *server) Subscribe(req *logpb.LogRequest, stream logpb.LogService_SubscribeServer) error {
ticker := time.NewTicker(500 * time.Millisecond)
defer ticker.Stop()
count := 0
for {
select {
case <-stream.Context().Done():
// 客户端断开或取消
log.Println("client gone:", stream.Context().Err())
return nil
case t := <-ticker.C:
count++
resp := &logpb.LogResponse{
Line: "topic " + req.Topic + " log " + string(rune(count)),
Ts: t.Unix(),
}
if err := stream.Send(resp); err != nil {
return err
}
if count >= 10 {
// 模拟推送完毕
return nil
}
}
}
}
四、客户端接收流数据
客户端调用Subscribe会拿到一个ClientStream,然后反复调用Recv()直到收到io.EOF。Recv是阻塞的,所以一般放在for循环里。当服务端正常结束流,gRPC会返回io.EOF,客户端据此退出循环。如果服务端返回error,Recv也会返回该error。
在实际工程中,客户端常把流读取放到独立goroutine,主逻辑通过channel拿数据,这样既能处理流也能响应退出信号。下面代码展示了最基础的同步接收方式,便于理解整体流程。
package main
import (
"context"
"io"
"log"
"time"
"ippipp.com/logpb"
"google.golang.org/grpc"
)
func main() {
conn, err := grpc.Dial("127.0.0.1:50051", grpc.WithInsecure())
if err != nil {
log.Fatal(err)
}
defer conn.Close()
client := logpb.NewLogServiceClient(conn)
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
stream, err := client.Subscribe(ctx, &logpb.LogRequest{Topic: "api"})
if err != nil {
log.Fatal(err)
}
for {
resp, err := stream.Recv()
if err == io.EOF {
log.Println("stream closed by server")
break
}
if err != nil {
log.Fatal(err)
}
log.Printf("recv: %s @ %d", resp.Line, resp.Ts)
}
}
五、性能与避坑要点
服务端流在Go里基于HTTP/2多路复用,一个TCP连接可并发多个流,因此不要为每个流新建连接。若消息频率极高,需注意Send的背压:客户端消费慢时,gRPC会缓冲一定量然后阻塞Send,从而自然限速,这比自己写队列更安全。
常见误区是把服务端流当消息队列用,长期挂着几万条流而不管心跳。正确做法是设置最大连接流数、利用keepalive参数探测死连接,并且在proto设计上让每条流都有明确生命周期。另外,proto里如果忘记写stream,生成的代码完全不含Send循环能力,编译不会报错但运行就是一元调用,这种问题要靠代码评审发现。
| 对比项 | 一元RPC | 服务端流 |
|---|---|---|
| 返回次数 | 一次 | 多次 |
| 适用场景 | 请求响应明确 | 实时推送、分批下发 |
| 客户端读取 | 直接拿返回值 | 循环Recv到EOF |
六、小结
用Golang实现gRPC服务端流,本质就是在proto里加stream,在服务端实现里循环调用Send,在客户端循环Recv。配合context取消和错误返回,就能写出稳定的推送服务。它比客户端轮询更省带宽,也比自己维护长轮询简单。掌握这种模式后,日志、监控、行情等场景的实时性需求都能优雅落地。