gRPC服务定义与Protobuf通信实战:高性能微服务RPC方案

gRPC是Google开源的高性能RPC框架,基于HTTP/2协议传输、Protobuf序列化数据。相比REST+JSON方案,gRPC在吞吐量上提升5-10倍,序列化体积减少3-5倍,延迟降低50%以上。gRPC原生支持流式通信、双向流控和连接复用,特别适合内部微服务间高频调用场景。本文以Go语言为例演示gRPC服务定义、实现、流式通信和拦截器机制。

Protobuf消息与服务定义

Protocol Buffers(Protobuf)是Google的序列化格式,通过.proto文件定义数据结构和服务接口。Protobuf采用二进制编码,字段通过编号而非名称标识,序列化后的体积远小于JSON:

// proto/user_service.proto
syntax = "proto3";

package user.v1;
option go_package = "github.com/example/proto/user/v1;userv1";

// 用户消息定义
message User {
  int64 id = 1;
  string name = 2;
  string email = 3;
  string phone = 4;
  UserStatus status = 5;
  int64 created_at = 6;
  repeated Address addresses = 7;
  
  enum UserStatus {
    UNKNOWN = 0;
    ACTIVE = 1;
    INACTIVE = 2;
    SUSPENDED = 3;
  }
}

message Address {
  string province = 1;
  string city = 2;
  string district = 3;
  string detail = 4;
  bool is_default = 5;
}

// 请求与响应消息
message GetUserRequest {
  int64 user_id = 1;
}

message GetUserResponse {
  User user = 1;
}

message ListUsersRequest {
  int32 page_size = 1;
  string page_token = 2;
  User.UserStatus status_filter = 3;
}

message ListUsersResponse {
  repeated User users = 1;
  string next_page_token = 2;
  int32 total_count = 3;
}

message BatchGetUsersRequest {
  repeated int64 user_ids = 1;
}

message BatchGetUsersResponse {
  map<int64, User> users = 1;
  repeated int64 not_found_ids = 2;
}

// 服务定义(支持Unary和Stream两种调用模式)
service UserService {
  // Unary RPC:一元调用(请求-响应)
  rpc GetUser(GetUserRequest) returns (GetUserResponse);
  rpc ListUsers(ListUsersRequest) returns (ListUsersResponse);
  rpc BatchGetUsers(BatchGetUsersRequest) returns (BatchGetUsersResponse);
  
  // Server Streaming:服务端流式返回
  rpc StreamUsers(ListUsersRequest) returns (stream User);
  
  // Client Streaming:客户端流式发送
  rpc CreateUsers(stream User) returns (BatchGetUsersResponse);
  
  // Bidirectional Streaming:双向流式通信
  rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}

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

使用protoc编译器生成Go代码。字段编号1-15占用1字节,16-2047占用2字节,频繁使用的字段应使用小编号。repeated关键字表示数组,map表示键值对。optional字段在proto3中默认所有字段均为optional,零值不会被序列化。

# 安装protoc和Go插件
protoc --go_out=. --go-grpc_out=. \
  proto/user_service.proto

# 生成的文件结构:
# proto/user/v1/user_service.pb.go      # 消息定义
# proto/user/v1/user_service_grpc.pb.go # gRPC服务接口

gRPC服务端实现

服务端实现UserServiceServer接口,注册到gRPC Server。每个RPC方法对应一个Go函数,ctx参数携带请求元数据(deadline、metadata等):

package server

import (
    "context"
    "errors"
    "time"
    
    userv1 "github.com/example/proto/user/v1"
    "google.golang.org/grpc/codes"
    "google.golang.org/grpc/status"
)

type UserServiceServer struct {
    userv1.UnimplementedUserServiceServer
    db    *sql.DB
    cache *redis.Client
}

// Unary RPC实现
func (s *UserServiceServer) GetUser(ctx context.Context, req *userv1.GetUserRequest) (*userv1.GetUserResponse, error) {
    if req.UserId <= 0 {
        return nil, status.Error(codes.InvalidArgument, "user_id must be positive")
    }
    
    // 先查缓存
    cached, err := s.cache.Get(ctx, fmt.Sprintf("user:%d", req.UserId)).Result()
    if err == nil {
        var user userv1.User
        if proto.Unmarshal([]byte(cached), &user) == nil {
            return &userv1.GetUserResponse{User: &user}, nil
        }
    }
    
    // 查数据库
    user, err := s.queryUserFromDB(ctx, req.UserId)
    if err != nil {
        if errors.Is(err, sql.ErrNoRows) {
            return nil, status.Error(codes.NotFound, "user not found")
        }
        return nil, status.Error(codes.Internal, err.Error())
    }
    
    // 写入缓存
    if data, err := proto.Marshal(user); err == nil {
        s.cache.Set(ctx, fmt.Sprintf("user:%d", req.UserId), data, 5*time.Minute)
    }
    
    return &userv1.GetUserResponse{User: user}, nil
}

// Server Streaming RPC实现
func (s *UserServiceServer) StreamUsers(req *userv1.ListUsersRequest, stream userv1.UserService_StreamUsersServer) error {
    // 分批查询并流式推送
    offset := 0
    batchSize := int32(100)
    
    for {
        users, err := s.queryUsersBatch(stream.Context(), batchSize, offset, req.StatusFilter)
        if err != nil {
            return status.Error(codes.Internal, err.Error())
        }
        if len(users) == 0 {
            return nil  // 数据发送完毕
        }
        
        for _, user := range users {
            if err := stream.Send(user); err != nil {
                return err  // 客户端断开或取消
            }
        }
        
        offset += len(users)
    }
}

// Bidirectional Streaming RPC实现
func (s *UserServiceServer) Chat(stream userv1.UserService_ChatServer) error {
    // 每个连接启动一个goroutine处理消息
    for {
        msg, err := stream.Recv()
        if err != nil {
            return err
        }
        
        // 业务处理...
        reply := &userv1.ChatMessage{
            UserId:    msg.UserId,
            Content:   "processed: " + msg.Content,
            Timestamp: time.Now().Unix(),
        }
        
        if err := stream.Send(reply); err != nil {
            return err
        }
    }
}

// 注册服务
func StartGRPCServer(addr string) error {
    lis, err := net.Listen("tcp", addr)
    if err != nil {
        return err
    }
    
    server := grpc.NewServer(
        grpc.MaxRecvMsgSize(16 * 1024 * 1024),  // 16MB
        grpc.MaxSendMsgSize(16 * 1024 * 1024),
        grpc.KeepaliveParams(keepalive.ServerParameters{
            MaxConnectionIdle:     5 * time.Minute,
            MaxConnectionAge:      30 * time.Minute,
            MaxConnectionAgeGrace: 5 * time.Second,
            Time:                  30 * time.Second,
            Timeout:              10 * time.Second,
        }),
    )
    
    userv1.RegisterUserServiceServer(server, &UserServiceServer{...})
    return server.Serve(lis)
}

gRPC客户端调用与连接管理

客户端使用grpc.Dial建立连接,支持连接池、负载均衡和超时控制。通过grpc.WithDefaultServiceConfig配置客户端负载均衡策略:

package client

import (
    "context"
    "time"
    
    "google.golang.org/grpc"
    "google.golang.org/grpc/credentials/insecure"
    "google.golang.org/grpc/balancer/roundrobin"
)

func NewUserClient(target string) (userv1.UserServiceClient, *grpc.ClientConn, error) {
    conn, err := grpc.Dial(target,
        grpc.WithTransportCredentials(insecure.NewCredentials()),
        // 客户端负载均衡:DNS解析多地址后轮询
        grpc.WithDefaultServiceConfig(`{
            "loadBalancingConfig": [{"round_robin": {}}]
        }`),
        grpc.WithDefaultCallOptions(
            grpc.MaxCallRecvMsgSize(16 * 1024 * 1024),
        ),
    )
    if err != nil {
        return nil, nil, err
    }
    
    return userv1.NewUserServiceClient(conn), conn, nil
}

// Unary调用示例
func GetUserByID(client userv1.UserServiceClient, userID int64) (*userv1.User, error) {
    ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
    defer cancel()
    
    resp, err := client.GetUser(ctx, &userv1.GetUserRequest{UserId: userID})
    if err != nil {
        if st, ok := status.FromError(err); ok {
            switch st.Code() {
            case codes.NotFound:
                return nil, fmt.Errorf("用户不存在")
            case codes.DeadlineExceeded:
                return nil, fmt.Errorf("请求超时")
            case codes.Unavailable:
                return nil, fmt.Errorf("服务不可用")
            }
        }
        return nil, err
    }
    return resp.User, nil
}

// Server Streaming调用示例
func ListAllUsers(client userv1.UserServiceClient) error {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()
    
    stream, err := client.StreamUsers(ctx, &userv1.ListUsersRequest{
        PageSize: 100,
    })
    if err != nil {
        return err
    }
    
    for {
        user, err := stream.Recv()
        if err == io.EOF {
            break
        }
        if err != nil {
            return err
        }
        fmt.Printf("收到用户: ID=%d, Name=%s\n", user.Id, user.Name)
    }
    return nil
}

拦截器与中间件机制

gRPC拦截器(Interceptor)类似HTTP中间件,在请求处理前后执行通用逻辑。Unary拦截器用于认证、日志、指标采集、链路追踪等横切关注点:

// 服务端Unary拦截器
func LoggingInterceptor(ctx context.Context, req interface{}, 
    info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
    
    start := time.Now()
    
    // 请求前:记录请求信息
    fmt.Printf("[gRPC] %s request: %+v\n", info.FullMethod, req)
    
    // 调用实际处理函数
    resp, err := handler(ctx, req)
    
    // 请求后:记录响应和耗时
    duration := time.Since(start)
    status := "OK"
    if err != nil {
        status = err.Error()
    }
    fmt.Printf("[gRPC] %s response: status=%s duration=%v\n", 
        info.FullMethod, status, duration)
    
    // 上报指标到Prometheus
    grpcRequestDuration.WithLabelValues(info.FullMethod, status).
        Observe(duration.Seconds())
    
    return resp, err
}

// 认证拦截器(JWT验证)
func AuthInterceptor(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.Error(codes.Unauthenticated, "no metadata")
    }
    
    tokens := md.Get("authorization")
    if len(tokens) == 0 {
        return nil, status.Error(codes.Unauthenticated, "no token")
    }
    
    // 验证JWT
    claims, err := validateJWT(tokens[0])
    if err != nil {
        return nil, status.Error(codes.Unauthenticated, "invalid token")
    }
    
    // 将用户信息注入context
    ctx = context.WithValue(ctx, "userID", claims.UserID)
    
    return handler(ctx, req)
}

// 链式拦截器
func ChainUnaryInterceptors(interceptors ...grpc.UnaryServerInterceptor) grpc.UnaryServerInterceptor {
    return func(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
        chain := handler
        for i := len(interceptors) - 1; i >= 0; i-- {
            chain = func(current grpc.UnaryHandler, interceptor grpc.UnaryServerInterceptor) grpc.UnaryHandler {
                return func(ctx context.Context, req interface{}) (interface{}, error) {
                    return interceptor(ctx, req, info, current)
                }
            }(chain, interceptors[i])
        }
        return chain(ctx, req)
    }
}

// 注册拦截器
server := grpc.NewServer(
    grpc.ChainUnaryInterceptor(
        RecoveryInterceptor,     // Panic恢复(最先执行)
        LoggingInterceptor,      // 日志记录
        AuthInterceptor,         // 认证
        MetricsInterceptor,      // 指标采集
    ),
)

gRPC Gateway RESTful API代理

gRPC适合内部服务通信,但浏览器和移动端原生不支持gRPC。grpc-gateway插件通过protoc生成反向代理,将HTTP/JSON请求转换为gRPC调用,同时对外提供RESTful API:

// 在proto文件中添加HTTP映射注解
import "google/api/annotations.proto";

service UserService {
  rpc GetUser(GetUserRequest) returns (GetUserResponse) {
    option (google.api.http) = {
      get: "/v1/users/{user_id}"
    };
  }
  
  rpc ListUsers(ListUsersRequest) returns (ListUsersResponse) {
    option (google.api.http) = {
      get: "/v1/users"
    };
  }
  
  rpc CreateUser(CreateUserRequest) returns (User) {
    option (google.api.http) = {
      post: "/v1/users"
      body: "user"
    };
  }
}

# 生成gateway代理代码
protoc --grpc-gateway_out=. proto/user_service.proto

# 启动HTTP网关,转发到gRPC服务
mux := runtime.NewServeMux()
userv1.RegisterUserServiceHandlerFromEndpoint(ctx, mux, "localhost:9090")
http.ListenAndServe(":8080", mux)

grpc-gateway生成的REST API自动处理JSON到Protobuf的转换,支持query参数映射、path参数提取和body序列化。通过一个.proto文件同时定义gRPC接口和REST API,保持两种协议的接口一致性。生产环境通常采用gRPC做内部通信、grpc-gateway对外暴露HTTP接口的混合架构。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/grpc-fu-wu-ding-yi-yu-protobuf-tong-xin-shi-zhan-gao-xing/

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

相关推荐