gRPC基于HTTP/2和Protocol Buffers实现高性能RPC通信,支持四种调用模式:Unary(一元调用)、Server Streaming(服务端流)、Client Streaming(客户端流)、Bidirectional Streaming(双向流)。流式通信适合实时数据推送、大文件分块传输和双向交互场景。相比REST的JSON序列化,gRPC使用二进制Protobuf编码,传输体积小3-5倍,解析速度快10倍以上。掌握proto协议设计和流式通信模式,是微服务架构中构建高性能服务治理体系的关键。
Protocol Buffers协议设计规范
proto文件定义服务接口和消息结构,是gRPC通信的契约文件。良好的proto设计应遵循向后兼容原则,避免破坏性变更。
syntax = "proto3";
package chat.v1;
option go_package = "github.com/example/chat-service/api/chat/v1;chatv1";
import "google/protobuf/timestamp.proto";
message ChatMessage {
int64 id = 1;
string room_id = 2;
string user_id = 3;
string content = 4;
enum MessageType {
MESSAGE_TYPE_UNSPECIFIED = 0;
TEXT = 1;
IMAGE = 2;
FILE = 3;
}
MessageType type = 5;
google.protobuf.Timestamp created_at = 6;
// oneof实现互斥字段
oneof metadata {
ImageMeta image_meta = 7;
FileMeta file_meta = 8;
}
// map类型
map<string, string> extensions = 9;
}
message ImageMeta {
int32 width = 1;
int32 height = 2;
string url = 3;
}
// 服务定义: 四种调用模式
service ChatService {
rpc SendMessage(SendMessageRequest) returns (SendMessageResponse);
rpc SubscribeMessages(SubscribeRequest) returns (stream ChatMessage);
rpc UploadFile(stream FileChunk) returns (UploadResponse);
rpc ChatStream(stream ChatMessage) returns (stream ChatMessage);
}
proto3设计要点:字段编号1-15使用1字节编码,高频字段优先分配小编号。删除字段时保留编号(使用reserved),防止复用导致反序列化混乱。枚举第一个值必须以_UNSPECIFIED = 0结尾作为默认值。
Go实现gRPC服务端流式通信
服务端流式RPC适用于实时推送场景,如消息订阅、日志流、股票行情推送。客户端发送一次请求,服务端持续推送数据直到关闭连接。
package main
import (
"context"
"log"
"net"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
pb "github.com/example/chat-service/api/chat/v1"
)
type ChatServiceServer struct {
pb.UnimplementedChatServiceServer
subscribers map[string]chan *pb.ChatMessage
}
func (s *ChatServiceServer) SubscribeMessages(
req *pb.SubscribeRequest,
stream pb.ChatService_SubscribeMessagesServer,
) error {
msgChan := make(chan *pb.ChatMessage, 100)
s.subscribers[req.RoomId] = msgChan
defer func() {
delete(s.subscribers, req.RoomId)
close(msgChan)
}()
ctx := stream.Context()
for {
select {
case <-ctx.Done():
return ctx.Err()
case msg := <-msgChan:
if err := stream.Send(msg); err != nil {
return status.Errorf(codes.Internal, "send failed: %v", err)
}
case <-time.After(30 * time.Second):
// 心跳保持连接
if err := stream.Send(&pb.ChatMessage{Content: "ping"}); err != nil {
return err
}
}
}
}
func (s *ChatServiceServer) BroadcastMessage(roomID string, msg *pb.ChatMessage) {
if ch, ok := s.subscribers[roomID]; ok {
select {
case ch <- msg:
default:
log.Printf("room %s message queue full", roomID)
}
}
}
func main() {
lis, err := net.Listen("tcp", ":50051")
if err != nil {
log.Fatalf("listen failed: %v", err)
}
server := grpc.NewServer(
grpc.MaxRecvMsgSize(16*1024*1024),
grpc.MaxSendMsgSize(16*1024*1024),
)
pb.RegisterChatServiceServer(server, &ChatServiceServer{
subscribers: make(map[string]chan *pb.ChatMessage),
})
log.Println("gRPC server started :50051")
server.Serve(lis)
}
双向流式RPC实现实时聊天
双向流允许客户端和服务端同时发送消息流,适合聊天室、协同编辑等实时双向交互场景。API接口规范在此模式下需设计消息类型区分控制信令和数据消息。
func (s *ChatServiceServer) ChatStream(
stream pb.ChatService_ChatStreamServer,
) error {
// 为该连接分配用户ID
userID := generateID()
msgChan := make(chan *pb.ChatMessage, 256)
// goroutine接收客户端消息
go func() {
for {
msg, err := stream.Recv()
if err != nil {
close(msgChan)
return
}
msg.UserId = userID
msg.CreatedAt = timestamppb.Now()
// 广播给房间内所有用户
s.broadcastToRoom(msg.RoomId, msg)
}
}()
// 主循环发送消息给客户端
ctx := stream.Context()
for {
select {
case <-ctx.Done():
return ctx.Err()
case msg, ok := <-msgChan:
if !ok {
return nil
}
if err := stream.Send(msg); err != nil {
return err
}
}
}
}
客户端流式RPC实现大文件上传
客户端流式RPC适用于批量数据上传场景,客户端分批发送数据流,服务端统一处理并返回最终结果。消息中间件在此场景中可作为缓冲层。
// 客户端代码 (Go)
func uploadFile(client pb.ChatServiceClient, filePath string) error {
stream, err := client.UploadFile(context.Background())
if err != nil {
return err
}
file, err := os.Open(filePath)
if err != nil {
return err
}
defer file.Close()
buf := make([]byte, 64*1024) // 64KB chunks
for {
n, err := file.Read(buf)
if err == io.EOF {
break
}
if err != nil {
return err
}
chunk := &pb.FileChunk{
Data: buf[:n],
}
if err := stream.Send(chunk); err != nil {
return err
}
}
resp, err := stream.CloseAndRecv()
if err != nil {
return err
}
log.Printf("upload complete: %s, size: %d", resp.Url, resp.Size)
return nil
}
gRPC拦截器实现认证与链路追踪
拦截器(Interceptor)是gRPC中间件机制,分为一元拦截器和流拦截器。高并发设计中常用于认证、限流、日志和链路追踪。服务治理依赖拦截器实现横切关注点。
// 服务端一元拦截器: 认证 + 日志
func authUnaryInterceptor(
ctx context.Context,
req interface{},
info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler,
) (interface{}, error) {
// 从metadata提取token
md, ok := metadata.FromIncomingContext(ctx)
if !ok {
return nil, status.Errorf(codes.Unauthenticated, "no metadata")
}
tokens := md.Get("authorization")
if len(tokens) == 0 {
return nil, status.Errorf(codes.Unauthenticated, "no token")
}
// 验证token
userID, err := validateToken(tokens[0])
if err != nil {
return nil, status.Errorf(codes.Unauthenticated, "invalid token")
}
// 将userID注入context
ctx = context.WithValue(ctx, userIDKey{}, userID)
// 记录请求日志
start := time.Now()
resp, err := handler(ctx, req)
log.Printf("method=%s duration=%v err=%v", info.FullMethod, time.Since(start), err)
return resp, err
}
// 流拦截器: 认证
func authStreamInterceptor(
srv interface{},
ss grpc.ServerStream,
info *grpc.StreamServerInfo,
handler grpc.StreamHandler,
) error {
md, ok := metadata.FromIncomingContext(ss.Context())
if !ok {
return status.Errorf(codes.Unauthenticated, "no metadata")
}
tokens := md.Get("authorization")
if len(tokens) == 0 {
return status.Errorf(codes.Unauthenticated, "no token")
}
_, err := validateToken(tokens[0])
if err != nil {
return status.Errorf(codes.Unauthenticated, "invalid token")
}
return handler(srv, ss)
}
// 注册拦截器
server := grpc.NewServer(
grpc.UnaryInterceptor(authUnaryInterceptor),
grpc.StreamInterceptor(authStreamInterceptor),
)
gRPC连接池与负载均衡配置
gRPC客户端默认使用HTTP/2多路复用,单个TCP连接可并发多个请求。但在高并发场景下单连接可能成为瓶颈,需配置连接池和负载均衡策略。
// 客户端连接配置 (Go)
conn, err := grpc.Dial(
"dns:///chat-service:50051", // DNS服务发现
grpc.WithDefaultServiceConfig(`{
"loadBalancingConfig": [{"round_robin": {}}],
"methodConfig": [{
"name": [{"service": "chat.v1.ChatService"}],
"retryPolicy": {
"maxAttempts": 3,
"initialBackoff": "0.1s",
"maxBackoff": "1s",
"backoffMultiplier": 2,
"retryableStatusCodes": ["UNAVAILABLE", "DEADLINE_EXCEEDED"]
}
}]
}`),
grpc.WithTransportCredentials(credentials.NewTLS(&tls.Config{})),
)
// 健康检查
conn, err := grpc.Dial(
"chat-service:50051",
grpc.WithDefaultServiceConfig(`{"loadBalancingConfig": [{"round_robin":{}}]}`),
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultCallOptions(
grpc.MaxCallRecvMsgSize(16*1024*1024),
grpc.MaxCallSendMsgSize(16*1024*1024),
),
)
round_robin策略将请求均匀分配到多个后端实例。结合DNS服务发现(dns:///前缀),客户端自动发现Kubernetes Service背后所有Pod并建立连接。retryPolicy配置对UNAVAILABLE和DEADLINE_EXCEEDED错误自动重试,提升分布式事务的可靠性。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/grpc-liu-shi-tong-xin-yu-protocolbuffers-xie-yi-she-ji-gui/