gRPC服务通信协议与Protobuf序列化在微服务中的应用实战

微服务架构中服务间通信的选型直接影响系统吞吐量和延迟。gRPC基于HTTP/2和Protocol Buffers,在序列化效率和传输性能上显著优于REST+JSON方案。Protobuf二进制序列化体积比JSON小3-10倍,HTTP/2多路复用消除了连接建立开销,使gRPC在高并发设计和分布式事务场景中获得更优的端到端延迟表现。

gRPC与REST通信协议对比与选型

  • 序列化格式:gRPC使用Protobuf二进制编码,JSON体积平均大5倍
  • 传输协议:gRPC基于HTTP/2支持多路复用和流控,REST基于HTTP/1.1需连接池
  • 接口定义:gRPC通过.proto文件定义强类型接口,自动生成多语言客户端代码
  • 流式通信:gRPC原生支持双向流,REST需通过WebSocket或SSE模拟
  • 浏览器兼容:gRPC需要gRPC-Web代理,REST可直接被浏览器调用

实际项目中,内部服务间通信使用gRPC,对外暴露的API网关使用REST/gRPC-Gateway转换。

Protobuf接口定义与代码生成

// proto/user_service.proto
syntax = "proto3";
package userservice.v1;
option go_package = "github.com/example/proto/userservice/v1";
import "google/protobuf/timestamp.proto";

service UserService {
  rpc GetUser(GetUserRequest) returns (GetUserResponse);
  rpc ListUsers(ListUsersRequest) returns (stream User);
  rpc BatchCreateUsers(stream CreateUserRequest) returns (BatchCreateResponse);
  rpc ChatStream(stream ChatMessage) returns (stream ChatMessage);
}

message User {
  int64 id = 1;
  string username = 2;
  string email = 3;
  UserStatus status = 4;
  google.protobuf.Timestamp created_at = 5;
  repeated string roles = 6;
  map<string, string> metadata = 7;
}

enum UserStatus {
  USER_STATUS_UNSPECIFIED = 0;
  USER_STATUS_ACTIVE = 1;
  USER_STATUS_INACTIVE = 2;
  USER_STATUS_SUSPENDED = 3;
}

message GetUserRequest { int64 id = 1; }
message GetUserResponse { User user = 1; }
message ListUsersRequest {
  int32 page_size = 1;
  string page_token = 2;
  UserStatus status_filter = 3;
}
message CreateUserRequest {
  string username = 1;
  string email = 2;
  string password = 3;
}
message BatchCreateResponse {
  repeated int64 created_ids = 1;
  int32 failed_count = 2;
}
message ChatMessage {
  int64 user_id = 1;
  string content = 2;
  google.protobuf.Timestamp sent_at = 3;
}
# 生成Go代码
protoc --go_out=. --go_opt=paths=source_relative \
    --go-grpc_out=. --go-grpc_opt=paths=source_relative \
    proto/user_service.proto

# 生成Python代码
python -m grpc_tools.protoc -I. --python_out=. --grpc_python_out=. \
    proto/user_service.proto

Go语言gRPC服务端实现

package main

import (
    "context"
    "log"
    "net"
    "sync"
    "time"
    "google.golang.org/grpc"
    "google.golang.org/grpc/codes"
    "google.golang.org/grpc/credentials"
    "google.golang.org/grpc/keepalive"
    "google.golang.org/grpc/metadata"
    "google.golang.org/grpc/status"
    "google.golang.org/protobuf/types/known/timestamppb"
    pb "github.com/example/proto/userservice/v1"
)

type userServer struct {
    pb.UnimplementedUserServiceServer
    users map[int64]*pb.User
    mu    sync.RWMutex
}

func (s *userServer) GetUser(ctx context.Context, req *pb.GetUserRequest) (*pb.GetUserResponse, error) {
    s.mu.RLock()
    defer s.mu.RUnlock()
    user, ok := s.users[req.GetId()]
    if !ok {
        return nil, status.Errorf(codes.NotFound, "user %d not found", req.GetId())
    }
    return &pb.GetUserResponse{User: user}, nil
}

func (s *userServer) ListUsers(req *pb.ListUsersRequest, stream pb.UserService_ListUsersServer) error {
    s.mu.RLock()
    defer s.mu.RUnlock()
    count := 0
    for _, user := range s.users {
        if req.StatusFilter != pb.UserStatus_USER_STATUS_UNSPECIFIED && 
           user.Status != req.StatusFilter {
            continue
        }
        if err := stream.Context().Err(); err != nil {
            return err
        }
        if err := stream.Send(user); err != nil {
            return err
        }
        count++
        if req.PageSize > 0 && int32(count) >= req.PageSize {
            break
        }
        time.Sleep(100 * time.Millisecond)
    }
    return nil
}

func (s *userServer) ChatStream(stream pb.UserService_ChatStreamServer) error {
    for {
        msg, err := stream.Recv()
        if err != nil {
            return err
        }
        reply := &pb.ChatMessage{
            UserId: msg.UserId,
            Content: "Echo: " + msg.Content,
            SentAt: timestamppb.Now(),
        }
        if err := stream.Send(reply); err != nil {
            return err
        }
    }
}

func main() {
    lis, _ := net.Listen("tcp", ":50051")
    creds, _ := credentials.NewServerTLSFromFile("cert/server.crt", "cert/server.key")
    s := grpc.NewServer(
        grpc.Creds(creds),
        grpc.KeepaliveParams(keepalive.ServerParameters{
            MaxConnectionIdle:     5 * time.Minute,
            MaxConnectionAge:      30 * time.Minute,
            MaxConnectionAgeGrace: 5 * time.Second,
            Time:                  30 * time.Second,
            Timeout:               10 * time.Second,
        }),
        grpc.MaxRecvMsgSize(16 * 1024 * 1024),
    )
    pb.RegisterUserServiceServer(s, &userServer{users: make(map[int64]*pb.User)})
    log.Println("gRPC server listening on :50051")
    s.Serve(lis)
}

拦截器Interceptor与中间件链

// 一元拦截器:日志鉴权
func loggingInterceptor(
    ctx context.Context,
    req interface{},
    info *grpc.UnaryServerInfo,
    handler grpc.UnaryHandler,
) (resp interface{}, err error) {
    start := time.Now()
    if err := authenticate(ctx, info.FullMethod); err != nil {
        return nil, err
    }
    resp, err = handler(ctx, req)
    log.Printf("method=%s duration=%s error=%v",
        info.FullMethod, time.Since(start), err)
    return resp, err
}

// 恢复拦截器:panic捕获
func recoveryInterceptor(
    ctx context.Context,
    req interface{},
    info *grpc.UnaryServerInfo,
    handler grpc.UnaryHandler,
) (resp interface{}, err error) {
    defer func() {
        if r := recover(); r != nil {
            log.Printf("panic recovered: %v", r)
            err = status.Errorf(codes.Internal, "internal error")
        }
    }()
    return handler(ctx, req)
}

func authenticate(ctx context.Context, method string) error {
    md, ok := metadata.FromIncomingContext(ctx)
    if !ok {
        return status.Error(codes.Unauthenticated, "no metadata")
    }
    tokens := md.Get("authorization")
    if len(tokens) == 0 {
        return status.Error(codes.Unauthenticated, "no token")
    }
    if !validateToken(tokens[0]) {
        return status.Error(codes.Unauthenticated, "invalid token")
    }
    return nil
}

// 链式拦截器
s := grpc.NewServer(
    grpc.ChainUnaryInterceptor(recoveryInterceptor, loggingInterceptor),
)

Python gRPC客户端实现

import grpc
import user_service_pb2 as pb
import user_service_pb2_grpc as pb_grpc
import time

channel = grpc.insecure_channel(
    'localhost:50051',
    options=[
        ('grpc.max_receive_message_length', 16 * 1024 * 1024),
        ('grpc.keepalive_time_ms', 30000),
        ('grpc.keepalive_timeout_ms', 10000),
        ('grpc.enable_retries', 1),
    ]
)

stub = pb_grpc.UserServiceStub(channel)

# 一元调用
response = stub.GetUser(
    pb.GetUserRequest(id=1),
    timeout=5,
    metadata=[('authorization', 'Bearer token_xxx')]
)
print(f"用户: {response.user.username}, 状态: {response.user.status}")

# 服务端流调用
for user in stub.ListUsers(
    pb.ListUsersRequest(page_size=10, status_filter=pb.UserStatus.USER_STATUS_ACTIVE),
    timeout=30
):
    print(f"  - {user.id}: {user.username}")

# 双向流调用
def request_generator():
    for i in range(10):
        yield pb.ChatMessage(user_id=1, content=f"消息 {i}")
        time.sleep(0.5)

responses = stub.ChatStream(request_generator(), timeout=60)
for reply in responses:
    print(f"收到回复: {reply.content}")

gRPC服务治理与健康检查

type healthServer struct {
    mu sync.Mutex
    status map[string]healthpb.HealthCheckResponse_ServingStatus
}

func (h *healthServer) Check(ctx context.Context, req *healthpb.HealthCheckRequest) (*healthpb.HealthCheckResponse, error) {
    h.mu.Lock()
    defer h.mu.Unlock()
    service := req.GetService()
    if service == "" {
        service = "overall"
    }
    status, ok := h.status[service]
    if !ok {
        return &healthpb.HealthCheckResponse{
            Status: healthpb.HealthCheckResponse_SERVING_STATUS_UNKNOWN,
        }, nil
    }
    return &healthpb.HealthCheckResponse{Status: status}, nil
}

// 注册健康检查
healthpb.RegisterHealthServer(s, &healthServer{
    status: map[string]healthpb.HealthCheckResponse_ServingStatus{
        "": healthpb.HealthCheckResponse_SERVING,
        "userservice.UserService": healthpb.HealthCheckResponse_SERVING,
    },
})

gRPC客户端默认复用HTTP/2连接,单连接支持100个并发流。生产环境建议根据并发量配置连接池大小,同时启用keepalive探活避免NAT超时断开。大消息体场景需调整MaxRecvMsgSize和流控窗口。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/grpc-fu-wu-tong-xin-xie-yi-yu-protobuf-xu-lie-hua-zai-wei/

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

相关推荐