gRPC微服务通信实战:Protocol Buffers定义与拦截器中间件设计

gRPC基于HTTP/2和Protocol Buffers构建,提供高性能、强类型的微服务通信方案。相比REST的JSON文本序列化,gRPC使用二进制Protobuf编码,消息体积更小、解析速度更快。HTTP/2的多路复用和头部压缩机制使gRPC在内部服务间通信场景下性能优势明显,适合高吞吐、低延迟的微服务架构

Protocol Buffers服务定义与代码生成

使用.proto文件定义服务接口和消息格式。一个典型的用户服务定义:

syntax = "proto3";

package user.v1;

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

service UserService {
  rpc GetUser(GetUserRequest) returns (GetUserResponse);
  rpc ListUsers(ListUsersRequest) returns (stream User);
  rpc CreateUserBatch(stream CreateUserRequest) returns (CreateUserBatchResponse);
  rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}

message GetUserRequest {
  int64 id = 1;
}

message GetUserResponse {
  User user = 1;
}

message User {
  int64 id = 1;
  string name = 2;
  string email = 3;
  UserStatus status = 4;
  int64 created_at = 5;
}

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

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

message CreateUserRequest {
  string name = 1;
  string email = 2;
}

message CreateUserBatchResponse {
  repeated int64 ids = 1;
}

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

gRPC支持四种调用模式:一元调用(Unary)、服务端流(Server Streaming)、客户端流(Client Streaming)和双向流(Bidirectional Streaming)。使用protoc工具生成各语言客户端和服务端代码:

protoc --go_out=. --go-grpc_out=.   --proto_path=proto   proto/user/v1/user.proto

Go服务端实现与注册

生成代码后实现服务接口,以Go为例:

package server

import (
    "context"
    "database/sql"
    userv1 "github.com/example/proto/user/v1"
)

type UserServer struct {
    userv1.UnimplementedUserServiceServer
    db *sql.DB
}

func (s *UserServer) GetUser(ctx context.Context, req *userv1.GetUserRequest) (*userv1.GetUserResponse, error) {
    var user userv1.User
    err := s.db.QueryRowContext(ctx,
        "SELECT id, name, email, status, created_at FROM users WHERE id = ?",
        req.Id,
    ).Scan(&user.Id, &user.Name, &user.Email, &user.Status, &user.CreatedAt)
    if err != nil {
        return nil, status.Errorf(codes.NotFound, "user %d not found", req.Id)
    }
    return &userv1.GetUserResponse{User: &user}, nil
}

func (s *UserServer) ListUsers(req *userv1.ListUsersRequest, stream userv1.UserService_ListUsersServer) error {
    rows, err := s.db.Query(
        "SELECT id, name, email, status FROM users WHERE status = ? LIMIT ?",
        req.StatusFilter, req.PageSize,
    )
    if err != nil {
        return status.Errorf(codes.Internal, "query failed: %v", err)
    }
    defer rows.Close()

    for rows.Next() {
        var user userv1.User
        if err := rows.Scan(&user.Id, &user.Name, &user.Email, &user.Status); err != nil {
            return status.Errorf(codes.Internal, "scan failed: %v", err)
        }
        if err := stream.Send(&user); err != nil {
            return err
        }
    }
    return nil
}

启动gRPC服务器并注册服务:

func main() {
    lis, err := net.Listen("tcp", ":50051")
    if err != nil {
        log.Fatalf("listen failed: %v", err)
    }

    grpcServer := grpc.NewServer(
        grpc.UnaryInterceptor(unaryInterceptor),
        grpc.StreamInterceptor(streamInterceptor),
    )

    userv1.RegisterUserServiceServer(grpcServer, &UserServer{db: db})

    log.Println("gRPC server listening on :50051")
    grpcServer.Serve(lis)
}

拦截器中间件与请求链设计

gRPC拦截器类似HTTP中间件,在请求处理前后插入横切逻辑。一元拦截器用于日志、认证、限流等:

func unaryInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
    start := time.Now()

    // 认证检查
    if err := checkAuth(ctx); err != nil {
        return nil, err
    }

    // 请求日志
    log.Printf("gRPC call: %s, request: %+v", info.FullMethod, req)

    // 执行handler
    resp, err := handler(ctx, req)

    // 响应日志与耗时
    duration := time.Since(start)
    if err != nil {
        log.Printf("gRPC error: %s, duration: %v, error: %v", info.FullMethod, duration, err)
        metrics.GRPCErrors.WithLabelValues(info.FullMethod, status.Code(err).String()).Inc()
    } else {
        log.Printf("gRPC done: %s, duration: %v", info.FullMethod, duration)
        metrics.GRPCLatency.WithLabelValues(info.FullMethod).Observe(duration.Seconds())
    }

    return resp, err
}

流式拦截器处理流请求:

func streamInterceptor(srv interface{}, ss grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error {
    ctx := ss.Context()

    if err := checkAuth(ctx); err != nil {
        return err
    }

    wrappedStream := &serverStreamWrapper{
        ServerStream: ss,
        method:       info.FullMethod,
    }

    return handler(srv, wrappedStream)
}

type serverStreamWrapper struct {
    grpc.ServerStream
    method string
}

func (w *serverStreamWrapper) SendMsg(m interface{}) error {
    log.Printf("stream send: %s, msg: %+v", w.method, m)
    return w.ServerStream.SendMsg(m)
}

func (w *serverStreamWrapper) RecvMsg(m interface{}) error {
    err := w.ServerStream.RecvMsg(m)
    if err == nil {
        log.Printf("stream recv: %s, msg: %+v", w.method, m)
    }
    return err
}

多拦截器链式组合

实际项目需要多个拦截器协同工作,如认证-日志-限流-恢复。Go gRPC支持传入多个拦截器但需要手动链式包装,或使用go-grpc-middleware库:

import grpcMiddleware "github.com/grpc-ecosystem/go-grpc-middleware"

grpcServer := grpc.NewServer(
    grpcMiddleware.WithUnaryServerChain(
        recovery.UnaryServerInterceptor(),   // panic恢复
        logging.UnaryServerInterceptor(),     // 请求日志
        auth.UnaryServerInterceptor(),        // JWT认证
        ratelimit.UnaryServerInterceptor(),   // 限流
    ),
    grpcMiddleware.WithStreamServerChain(
        recovery.StreamServerInterceptor(),
        logging.StreamServerInterceptor(),
        auth.StreamServerInterceptor(),
    ),
)

拦截器按声明顺序执行,每个拦截器调用handler前是请求预处理阶段,调用后是响应后处理阶段。执行顺序类似洋葱模型:外层拦截器先执行预处理,内层先执行后处理。

客户端调用与连接池管理

gRPC客户端使用grpc.Dial建立连接,底层复用HTTP/2连接,支持多路复用:

conn, err := grpc.Dial(
    "localhost:50051",
    grpc.WithTransportCredentials(insecure.NewCredentials()),
    grpc.WithDefaultServiceConfig(`{
        "loadBalancingPolicy": "round_robin",
        "methodConfig": [{
            "name": [{"service": "user.v1.UserService"}],
            "retryPolicy": {
                "maxAttempts": 3,
                "initialBackoff": "0.1s",
                "maxBackoff": "1s",
                "retryableStatusCodes": ["UNAVAILABLE", "DEADLINE_EXCEEDED"]
            }
        }]
    }`),
    grpc.WithUnaryInterceptor(clientLogInterceptor),
)
if err != nil {
    log.Fatalf("dial failed: %v", err)
}
defer conn.Close()

client := userv1.NewUserServiceClient(conn)
resp, err := client.GetUser(ctx, &userv1.GetUserRequest{Id: 1})

Service Config配置了负载均衡策略和重试策略。DNS解析多个地址时round_robin策略自动轮询。重试策略仅对幂等操作安全,需谨慎配置retryableStatusCodes

连接池在gRPC中由底层HTTP/2连接自动管理。单个TCP连接支持100个并发流(HTTP/2默认),可通过grpc.WithMaxConcurrentStreams调整。高并发场景下建议保持连接复用,避免频繁建连。

错误处理与状态码规范

gRPC定义了标准状态码,服务端通过status.Errorf返回带状态码的错误:

switch {
case err == sql.ErrNoRows:
    return nil, status.Error(codes.NotFound, "user not found")
case strings.Contains(err.Error(), "duplicate"):
    return nil, status.Error(codes.AlreadyExists, "email already registered")
case strings.Contains(err.Error(), "connection refused"):
    return nil, status.Error(codes.Unavailable, "database unavailable")
default:
    return nil, status.Error(codes.Internal, err.Error())
}

客户端根据状态码执行重试或向上传递错误:status.Code(err)获取状态码,status.Convert(err).Message()获取错误消息。建议在API网关层将gRPC状态码映射为HTTP状态码,如NotFound对应404、AlreadyExists对应409、PermissionDenied对应403。

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

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

相关推荐