gRPC流式通信模式双向流与服务器推送实现详解

gRPC四种通信模式与流式语义对比

gRPC基于HTTP/2协议提供了四种通信模式:Unary(一元调用)、Server Streaming(服务端流)、Client Streaming(客户端流)、Bidirectional Streaming(双向流)。Unary模式类似传统RPC,一问一答;Server Streaming适用于服务端逐步推送数据的场景(如实时日志、股票行情);Client Streaming适用于客户端分批上传数据的场景(如文件分片上传);Bidirectional Streaming支持双方同时读写,是实现全双工交互式通信的基础。

流式通信的核心优势在于:降低首字节延迟(TTFB),服务端无需等待完整请求即可开始响应;支持长连接场景,避免频繁建立连接的开销;天然适配背压(Backpressure),慢消费端可自动控制数据流速。

Proto文件定义与服务生成

// chat.proto - 双向流式聊天服务
syntax = "proto3";

package chat;

service ChatService {
  rpc ChatStream (stream ChatMessage) returns (stream ChatMessage);
  rpc Subscribe (Subscription) returns (stream Notification);
  rpc UploadFile (stream FileChunk) returns (UploadResult);
}

message ChatMessage {
  string user = 1;
  string content = 2;
  int64 timestamp = 3;
}

message Subscription {
  string topic = 1;
  string client_id = 2;
}

message Notification {
  string topic = 1;
  string payload = 2;
  int64 timestamp = 3;
}

message FileChunk {
  string filename = 1;
  bytes data = 2;
  int32 sequence = 3;
  bool is_last = 4;
}

message UploadResult {
  bool success = 1;
  string file_url = 2;
  int32 total_chunks = 3;
}

生成Go代码:

protoc --go_out=. --go-grpc_out=. chat.proto

Go语言双向流服务端实现

package main

import (
    "io"
    "log"
    "sync"
    "time"

    pb "./chat"
)

type ChatRoom struct {
    clients map[string]pb.ChatService_ChatStreamServer
    mu      sync.RWMutex
}

var room = &ChatRoom{
    clients: make(map[string]pb.ChatService_ChatStreamServer),
}

func (s *server) ChatStream(stream pb.ChatService_ChatStreamServer) error {
    var userID string

    for {
        msg, err := stream.Recv()
        if err == io.EOF {
            room.mu.Lock()
            delete(room.clients, userID)
            room.mu.Unlock()
            return nil
        }
        if err != nil {
            log.Printf("Stream error: %v", err)
            return err
        }

        if userID == "" {
            userID = msg.User
            room.mu.Lock()
            room.clients[userID] = stream
            room.mu.Unlock()
        }

        room.mu.RLock()
        for uid, clientStream := range room.clients {
            if uid != userID {
                if err := clientStream.Send(&pb.ChatMessage{
                    User:      msg.User,
                    Content:   msg.Content,
                    Timestamp: time.Now().Unix(),
                }); err != nil {
                    log.Printf("Send to %s failed: %v", uid, err)
                }
            }
        }
        room.mu.RUnlock()
    }
}

双向流的错误处理需要注意:Send和Recv可能在不同goroutine中执行,任何一方出错都会导致整个流断开。服务端应在检测到错误后主动清理客户端注册信息。

服务端流推送与客户端消费

服务端流模式适合实时通知推送场景。客户端发送一次订阅请求,服务端持续推送数据:

func (s *server) Subscribe(req *pb.Subscription, stream pb.ChatService_SubscribeServer) error {
    ticker := time.NewTicker(2 * time.Second)
    defer ticker.Stop()

    for {
        select {
        case <-stream.Context().Done():
            log.Printf("Client %s disconnected", req.ClientId)
            return stream.Context().Err()
        case t := <-ticker.C:
            notification := &pb.Notification{
                Topic:     req.Topic,
                Payload:   fmt.Sprintf("Update at %s", t.Format(time.RFC3339)),
                Timestamp: t.Unix(),
            }
            if err := stream.Send(notification); err != nil {
                return err
            }
        }
    }
}

客户端消费服务端流:

func subscribeNotifications(client pb.ChatServiceClient, topic string) {
    stream, err := client.Subscribe(context.Background(), &pb.Subscription{
        Topic:    topic,
        ClientId: "client-001",
    })
    if err != nil {
        log.Fatal(err)
    }

    for {
        notification, err := stream.Recv()
        if err == io.EOF {
            break
        }
        if err != nil {
            log.Printf("Receive error: %v", err)
            break
        }
        log.Printf("[%s] %s", notification.Topic, notification.Payload)
    }
}

客户端流文件上传与背压控制

func (s *server) UploadFile(stream pb.ChatService_UploadFileServer) error {
    var totalSize int64
    var filename string

    for {
        chunk, err := stream.Recv()
        if err == io.EOF {
            return stream.SendAndClose(&pb.UploadResult{
                Success:     true,
                FileUrl:     fmt.Sprintf("/uploads/%s", filename),
                TotalChunks: int32(totalSize),
            })
        }
        if err != nil {
            return err
        }

        filename = chunk.Filename
        totalSize += int64(len(chunk.Data))

        if err := appendToFile(filename, chunk.Data); err != nil {
            return err
        }
    }
}

gRPC流式通信的背压机制由HTTP/2的流控层自动处理。当消费端处理速度跟不上生产端时,HTTP/2 WINDOW_UPDATE帧会自动调节数据流速。在Go实现中,Send操作在缓冲区满时自动阻塞,Recv操作在无数据时自动阻塞,开发者无需手动实现流控逻辑。生产环境需注意设置合理的MaxRecvMsgSize和MaxSendMsgSize(默认4MB),大文件上传场景应适当增大或采用分片策略。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/grpc-liu-shi-tong-xin-mo-shi-shuang-xiang-liu-yu-fu-wu-qi/

(0)
小编小编
上一篇 5小时前
下一篇 5小时前

相关推荐