kratos中实现服务端流式响应需使用grpc的serverstreaming模式:在.proto文件中定义返回stream类型的方法,生成代码后,在service层通过srv.send()逐条发送响应,确保在单goroutine内完成,不可并发调用。

Kratos框架服务端流式响应怎么实现
要在Kratos中实现服务端向客户端持续推送数据(如日志实时输出、大文件分块传输、事件流通知),必须使用gRPC的ServerStreaming模式,不能依赖HTTP长轮询或SSE——因为Kratos的HTTP层默认不支持原生流式响应,而gRPC层从协议设计上就天然支持单请求多响应。
定义proto并生成ServerStreaming接口
在.proto文件中声明一个返回stream类型的RPC方法:
```protobuf
service Greeter {
rpc SayHelloStream (HelloRequest) returns (stream HelloReply);
}
注意:返回类型必须是stream HelloReply,不是HelloReply;请求参数仍为单值,这是ServerStreaming的特征——一请求,多响应。
【必须用goctl或protoc-gen-go重新生成代码】,否则GreeterServer接口里不会出现SayHelloStream方法签名,也不会生成对应的Greeter_SayHelloStreamServer流式服务端对象。
在Service层编写流式逻辑
打开生成的service/greeter.go,找到SayHelloStream方法实现位置:
```go
func (s *GreeterService) SayHelloStream(ctx context.Context, req *v1.HelloRequest, srv v1.Greeter_SayHelloStreamServer) error {
for i := 0; i reply := &v1.HelloReply{Message: fmt.Sprintf("Hello %s, count %d", req.GetName(), i)}
if err := srv.Send(reply); err != nil {
return err
}
time.Sleep(1 * time.Second)
}
return nil
}
【srv.Send()必须在同一个goroutine内调用,不可并发写入】。gRPC流式服务端对象不是线程安全的,如果在多个goroutine里同时调用srv.Send(),会触发panic: “send on closed channel”。
这一步操作起来很简单,直接把循环和Send写进去就行,但务必确保整个发送过程在当前函数协程中完成——不要起新goroutine去Send。
启动gRPC服务并暴露端口
确认internal/server/grpc.go中已启用gRPC服务:
```go
srv := grpc.NewServer(
grpc.Address(":9000"),
grpc.Middleware(
recover.Recover(),
tracing.Server(),
),
)
v1.RegisterGreeterServer(srv, service.NewGreeterService())
return srv
注意端口不能与HTTP服务冲突;若已有HTTP服务占用了:8000,gRPC必须换端口(如:9000)——Kratos默认不允许多协议共用同一监听地址。
启动后,可用grpcurl -plaintext -d '{"name":"kratos"}' localhost:9000 v1.Greeter/SayHelloStream验证流式响应是否逐条返回。
前端或客户端如何消费流式响应
方法一:用grpcurl命令行工具(调试首选)
```bash
grpcurl -plaintext -d '{"name":"test"}' localhost:9000 v1.Greeter/SayHelloStream
方法二:Go客户端调用(生产环境)
```go
conn, _ := grpc.Dial("localhost:9000", grpc.WithTransportCredentials(insecure.NewCredentials()))
defer conn.Close()
client := v1.NewGreeterClient(conn)
stream, _ := client.SayHelloStream(context.Background(), &v1.HelloRequest{Name: "client"})
for {
resp, err := stream.Recv()
if err == io.EOF { break }
if err != nil { log.Fatal(err) }
log.Println(resp.Message)
}
Recv()会阻塞直到下一条消息到达或流关闭;收到io.EOF表示服务端已结束发送,此时循环应退出。











