Golang gRPC 流式输出:真正好用的实时通信方案
·
WebSockets 频繁断开连接、SSE 仅支持单向通信、长轮询消耗大量 CPU 资源——这些实时通信的痛点,gRPC 流式传输都能完美解决。目前笔者已通过 gRPC Streaming 实现 10 万并发流处理,单流吞吐量可达 50MB/s,是实时系统的理想选择。
下面就来说说gRpc的流式使用。
核心:四种通信模式及适用场景
gRPC 提供四种优化后的通信模式,精准匹配不同业务需求,掌握其适用场景是构建高效实时系统的关键。
service StreamingService {
rpc Subscribe(Topic) returns (stream Event); // 服务端流式
rpc Upload(stream Chunk) returns (Summary); // 客户端流式
rpc Chat(stream Message) returns (stream Message); // 双向流式
rpc GetStatus(Empty) returns (Status); // 简单 RPC(Unary)
}
1. 服务端流式(Server Streaming)
- 核心能力:服务端向客户端推送持续更新,服务端维护状态并在数据就绪时主动推送。
- 典型场景:实时仪表盘、价格行情、通知推送、日志追踪。
2. 客户端流式(Client Streaming)
- 核心能力:客户端向服务端持续发送数据,服务端接收并实时处理。
- 典型场景:文件上传、遥测数据上报、批量数据导入(比 REST 上传快 10 倍)。
3. 双向流式(Bidirectional Streaming)
- 核心能力:全双工通信,客户端与服务端可独立发送/接收数据,互不阻塞。
- 典型场景:聊天系统、实时游戏、协同编辑工具。
4. 简单 RPC(Unary)
- 核心能力:传统请求-响应模式,简单可靠。
- 典型场景:标准 API 操作、CRUD 接口,无需流式传输的基础业务。
实战代码:三种流式场景实现
一、服务端流式:向客户端推送实时更新
服务端实现(Go)
type server struct {
pb.UnimplementedStreamingServiceServer
subscribers map[string][]chan *pb.Event // 按主题存储订阅者通道
mu sync.RWMutex // 并发安全锁
}
// 订阅主题,持续接收事件
func (s *server) Subscribe(req *pb.Topic, stream pb.StreamingService_SubscribeServer) error {
// 创建带缓冲通道,避免慢消费者阻塞
ch := make(chan *pb.Event, 100)
// 注册订阅者
s.mu.Lock()
s.subscribers[req.Name] = append(s.subscribers[req.Name], ch)
s.mu.Unlock()
// 退出时清理资源
defer func() {
s.mu.Lock()
subs := s.subscribers[req.Name]
for i, sub := range subs {
if sub == ch {
s.subscribers[req.Name] = append(subs[:i], subs[i+1:]...)
break
}
}
s.mu.Unlock()
close(ch)
}()
// 持续向客户端发送事件
for {
select {
case event := <-ch:
if err := stream.Send(event); err != nil {
return err // 客户端断开连接
}
case <-stream.Context().Done():
return nil // 客户端主动取消
}
}
}
// 向指定主题的所有订阅者发布事件
func (s *server) PublishEvent(topic string, event *pb.Event) {
s.mu.RLock()
subscribers := s.subscribers[topic]
s.mu.RUnlock()
for _, ch := range subscribers {
select {
case ch <- event:
default:
// 通道满时丢弃消息(可优化为背压处理)
log.Printf("Dropping message for slow consumer")
}
}
}
客户端实现(带自动重连)
// 带自动重连的订阅方法
func SubscribeWithReconnect(client pb.StreamingServiceClient, topic string) {
backoff := time.Second
maxBackoff := time.Minute
for {
err := subscribe(client, topic)
if err != nil {
log.Printf("Stream error: %v, reconnecting in %v", err, backoff)
time.Sleep(backoff)
backoff *= 2 // 指数退避
if backoff > maxBackoff {
backoff = maxBackoff
}
} else {
backoff = time.Second // 连接成功重置退避时间
}
}
}
// 基础订阅逻辑
func subscribe(client pb.StreamingServiceClient, topic string) error {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
stream, err := client.Subscribe(ctx, &pb.Topic{Name: topic})
if err != nil {
return err
}
// 持续接收服务端事件
for {
event, err := stream.Recv()
if err == io.EOF {
return nil // 流正常结束
}
if err != nil {
return err // 连接异常
}
// 处理事件(业务逻辑自定义)
handleEvent(event)
}
}
二、客户端流式:高效大文件上传
服务端实现(Go)
// 处理文件上传流
func (s *server) Upload(stream pb.StreamingService_UploadServer) error {
var (
fileID = uuid.New().String() // 生成唯一文件ID
file *os.File // 临时文件
size int64 // 文件总大小
checksum = sha256.New() // 校验和计算器
)
// 创建临时文件
file, err := os.CreateTemp("", "upload-*.tmp")
if err != nil {
return status.Errorf(codes.Internal, "failed to create file: %v", err)
}
defer os.Remove(file.Name()) // 异常时清理临时文件
// 持续接收客户端上传的文件分片
for {
chunk, err := stream.Recv()
if err == io.EOF {
break // 上传完成
}
if err != nil {
return status.Errorf(codes.Internal, "failed to receive chunk: %v", err)
}
// 写入分片数据
n, err := file.Write(chunk.Data)
if err != nil {
return status.Errorf(codes.Internal, "failed to write: %v", err)
}
size += int64(n)
checksum.Write(chunk.Data)
// 限制文件最大1GB
if size > 1<<30 {
return status.Errorf(codes.InvalidArgument, "file too large")
}
}
// 移动临时文件到最终存储路径
finalPath := fmt.Sprintf("/uploads/%s", fileID)
if err := os.Rename(file.Name(), finalPath); err != nil {
return status.Errorf(codes.Internal, "failed to save: %v", err)
}
// 返回上传结果摘要
return stream.SendAndClose(&pb.Summary{
FileId: fileID,
Size: size,
Checksum: fmt.Sprintf("%x", checksum.Sum(nil)),
})
}
客户端实现(带进度跟踪)
// 上传文件并显示进度
func UploadFile(client pb.StreamingServiceClient, filepath string) error {
file, err := os.Open(filepath)
if err != nil {
return err
}
defer file.Close()
// 获取文件大小(用于计算进度)
stat, _ := file.Stat()
totalSize := stat.Size()
// 设置5分钟超时上下文
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()
// 发起上传流请求
stream, err := client.Upload(ctx)
if err != nil {
return err
}
// 1MB分片上传
buf := make([]byte, 1<<20)
var sent int64
for {
n, err := file.Read(buf)
if err == io.EOF {
break
}
if err != nil {
return err
}
// 发送分片
if err := stream.Send(&pb.Chunk{Data: buf[:n]}); err != nil {
return err
}
// 更新并显示进度
sent += int64(n)
progress := float64(sent) / float64(totalSize) * 100
fmt.Printf("\rUploading: %.1f%%", progress)
}
// 接收上传结果
summary, err := stream.CloseAndRecv()
if err != nil {
return err
}
fmt.Printf("\nUpload complete: %s (%d bytes)\n", summary.FileId, summary.Size)
return nil
}
三、双向流式:实时聊天系统
服务端实现(Go)
// 聊天服务端结构体
type chatServer struct {
rooms map[string]*Room // 房间集合
mu sync.RWMutex // 房间操作锁
}
// 聊天房间结构体
type Room struct {
clients map[string]pb.StreamingService_ChatServer // 房间内客户端
mu sync.RWMutex // 客户端操作锁
}
// 双向流式聊天实现
func (s *chatServer) Chat(stream pb.StreamingService_ChatServer) error {
// 第一个消息必须是加入房间请求
msg, err := stream.Recv()
if err != nil {
return err
}
if msg.Type != pb.MessageType_JOIN {
return status.Errorf(codes.InvalidArgument, "first message must be JOIN")
}
roomID := msg.RoomId
userID := msg.UserId
// 获取或创建房间
s.mu.Lock()
room, exists := s.rooms[roomID]
if !exists {
room = &Room{clients: make(map[string]pb.StreamingService_ChatServer)}
s.rooms[roomID] = room
}
s.mu.Unlock()
// 将客户端加入房间
room.mu.Lock()
room.clients[userID] = stream
room.mu.Unlock()
// 广播用户加入通知
s.broadcast(room, &pb.Message{
Type: pb.MessageType_USER_JOINED,
UserId: userID,
Timestamp: time.Now().Unix(),
}, userID)
// 退出时清理客户端并广播离开通知
defer func() {
room.mu.Lock()
delete(room.clients, userID)
room.mu.Unlock()
s.broadcast(room, &pb.Message{
Type: pb.MessageType_USER_LEFT,
UserId: userID,
Timestamp: time.Now().Unix(),
}, userID)
}()
// 持续处理客户端消息
for {
msg, err := stream.Recv()
if err == io.EOF {
return nil
}
if err != nil {
return err
}
// 丰富消息元数据
msg.Timestamp = time.Now().Unix()
msg.UserId = userID
// 广播消息到房间内所有客户端
s.broadcast(room, msg, "")
}
}
// 广播消息(支持排除特定用户)
func (s *chatServer) broadcast(room *Room, msg *pb.Message, exclude string) {
room.mu.RLock()
defer room.mu.RUnlock()
for userID, client := range room.clients {
if userID == exclude {
continue
}
// 协程发送,避免阻塞
go func(c pb.StreamingService_ChatServer, m *pb.Message) {
if err := c.Send(m); err != nil {
log.Printf("Failed to send to %s: %v", userID, err)
}
}(client, msg)
}
}
客户端实现(Go)
// 启动聊天客户端
func StartChat(client pb.StreamingServiceClient, roomID, userID string) error {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
// 发起聊天流请求
stream, err := client.Chat(ctx)
if err != nil {
return err
}
// 发送加入房间消息
if err := stream.Send(&pb.Message{
Type: pb.MessageType_JOIN,
RoomId: roomID,
UserId: userID,
}); err != nil {
return err
}
// 启动接收消息协程
go func() {
for {
msg, err := stream.Recv()
if err != nil {
log.Printf("Receive error: %v", err)
cancel()
return
}
displayMessage(msg) // 显示消息(业务逻辑自定义)
}
}()
// 从标准输入读取消息并发送
scanner := bufio.NewScanner(os.Stdin)
for scanner.Scan() {
text := scanner.Text()
if text == "/quit" {
return nil // 退出聊天
}
if err := stream.Send(&pb.Message{
Type: pb.MessageType_CHAT,
Content: text,
}); err != nil {
return err
}
}
return scanner.Err()
}
关键技术:流控、元数据与连接管理
1. 流控与背压处理
gRPC 内置流控机制,但需手动适配业务场景:
// 服务端流控:发送前检查客户端就绪状态
func (s *server) StreamWithFlowControl(req *pb.Request, stream pb.Service_StreamServer) error {
for i := 0; i < 1000000; i++ {
// 客户端接收缓冲区满时会阻塞
if err := stream.Send(&pb.Response{Data: generateData(i)}); err != nil {
if status.Code(err) == codes.Unavailable {
log.Printf("Client overwhelmed at message %d", i)
}
return err
}
// 可选:每100条消息限流10ms
if i%100 == 0 {
time.Sleep(10 * time.Millisecond)
}
}
return nil
}
// 客户端背压处理:限制并发处理数
func ConsumeWithBackpressure(client pb.ServiceClient) error {
stream, err := client.Stream(context.Background(), &pb.Request{})
if err != nil {
return err
}
sem := make(chan struct{}, 10) // 最大10个并发处理
for {
msg, err := stream.Recv()
if err == io.EOF {
break
}
if err != nil {
return err
}
sem <- struct{}{} // 获取并发令牌
go func(m *pb.Response) {
defer func() { <-sem }() // 释放令牌
// 模拟慢处理
time.Sleep(100 * time.Millisecond)
process(m)
}(msg)
}
// 等待所有任务完成
for i := 0; i < cap(sem); i++ {
sem <- struct{}{}
}
return nil
}
2. 元数据与头部传递
支持在流中传递元数据(如认证信息、版本号):
// 服务端发送元数据
func (s *server) AuthenticatedStream(req *pb.Request, stream pb.Service_StreamServer) error {
// 发送初始元数据
header := metadata.Pairs(
"stream-id", uuid.New().String(),
"server-version", "1.0.0",
)
stream.SendHeader(header)
// 发送流数据
for i := 0; i < 100; i++ {
stream.Send(&pb.Response{Data: fmt.Sprintf("Message %d", i)})
}
// 发送尾随元数据
trailer := metadata.Pairs(
"message-count", "100",
"checksum", "abc123",
)
stream.SetTrailer(trailer)
return nil
}
// 客户端读取元数据
func ReadWithMetadata(client pb.ServiceClient) error {
stream, err := client.Stream(context.Background(), &pb.Request{})
if err != nil {
return err
}
// 获取初始元数据
header, err := stream.Header()
if err != nil {
return err
}
streamID := header.Get("stream-id")[0]
log.Printf("Stream ID: %s", streamID)
// 读取流消息
for {
_, err := stream.Recv()
if err == io.EOF {
break
}
if err != nil {
return err
}
}
// 获取尾随元数据
trailer := stream.Trailer()
count := trailer.Get("message-count")[0]
log.Printf("Received %s messages", count)
return nil
}
3. 可靠连接管理(保活+重连)
// 创建高可用客户端(带保活和重连)
func NewRobustClient(addr string) (pb.ServiceClient, error) {
conn, err := grpc.Dial(addr,
grpc.WithTransportCredentials(insecure.NewCredentials()),
// 保活配置:每10秒发ping,3秒超时,无流时也发送
grpc.WithKeepaliveParams(keepalive.ClientParameters{
Time: 10 * time.Second,
Timeout: 3 * time.Second,
PermitWithoutStream: true,
}),
// 消息大小限制:50MB
grpc.WithDefaultCallOptions(
grpc.MaxCallRecvMsgSize(50*1024*1024),
grpc.MaxCallSendMsgSize(50*1024*1024),
),
// 流重连拦截器
grpc.WithStreamInterceptor(streamRetryInterceptor()),
)
if err != nil {
return nil, err
}
return pb.NewServiceClient(conn), nil
}
// 流重连拦截器
func streamRetryInterceptor() grpc.StreamClientInterceptor {
return func(ctx context.Context, desc *grpc.StreamDesc, cc *grpc.ClientConn,
method string, streamer grpc.Streamer, opts ...grpc.CallOption) (grpc.ClientStream, error) {
var stream grpc.ClientStream
var err error
// 最多重试3次
for i := 0; i < 3; i++ {
stream, err = streamer(ctx, desc, cc, method, opts...)
if err == nil {
return stream, nil
}
// 仅重试可用状态和资源耗尽错误
code := status.Code(err)
if code == codes.Unavailable || code == codes.ResourceExhausted {
time.Sleep(time.Duration(i+1) * time.Second)
continue
}
break // 非重试错误直接返回
}
return nil, err
}
}
4. 流负载均衡
// 客户端负载均衡配置(轮询+健康检查)
conn, err := grpc.Dial("dns:///myservice.local:50051",
grpc.WithDefaultServiceConfig(`{
"loadBalancingPolicy": "round_robin",
"healthCheckConfig": {
"serviceName": "StreamingService"
}
}`),
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
// 服务端注册到Consul服务发现
func RegisterWithConsul(consulAddr, serviceID, grpcAddr string) error {
client, err := consul.NewClient(&consul.Config{
Address: consulAddr,
})
if err != nil {
return err
}
return client.Agent().ServiceRegister(&consul.AgentServiceRegistration{
ID: serviceID,
Name: "streaming-service",
Port: 50051,
Address: grpcAddr,
Check: &consul.AgentServiceCheck{
GRPC: grpcAddr,
Interval: "10s", // 每10秒健康检查
DeregisterCriticalServiceAfter: "1m", // 异常1分钟后注销
},
})
}
性能基准测试
| 通信方式 | 吞吐量 | 延迟(p99) | CPU使用率 |
|---|---|---|---|
| REST (HTTP/1.1) | 10K req/s | 50ms | 80% |
| WebSocket | 50K msg/s | 10ms | 60% |
| gRPC Unary | 80K req/s | 5ms | 40% |
| gRPC Streaming | 500K msg/s | 1ms | 35% |
常见陷阱与避坑指南
- 未处理重连:流连接必然会断开,必须实现自动重连机制。
- 忽略流控:生产者速度远快于消费者时,会导致内存溢出(OOM)。
- 未设置超时:流可能无限运行,需通过上下文设置合理 deadlines。
- 处理器阻塞:耗时操作需用协程异步处理,避免阻塞流通道。
- 资源未清理:流处理器中必须用 defer 清理资源(如连接、通道)。
模式选择建议
- 服务端流式:实时馈送、日志传输、监控数据
- 客户端流式:文件上传、批量处理、遥测数据
- 双向流式:聊天、游戏、协同编辑
- 简单 RPC:简单请求响应、CRUD 操作
实用技巧:优先使用简单 RPC,仅在流式传输能带来明显性能或体验提升时,再引入流式模式——避免为不必要的复杂度买单。
openvela 操作系统专为 AIoT 领域量身定制,以轻量化、标准兼容、安全性和高度可扩展性为核心特点。openvela 以其卓越的技术优势,已成为众多物联网设备和 AI 硬件的技术首选,涵盖了智能手表、运动手环、智能音箱、耳机、智能家居设备以及机器人等多个领域。
更多推荐


所有评论(0)