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/