gRPC微服务通信实战:Protobuf序列化与拦截器中间件

gRPC协议与HTTP/2传输层

gRPC是Google开源的高性能RPC框架,基于HTTP/2协议和Protocol Buffers序列化格式。与RESTful JSON相比,gRPC的优势在于:HTTP/2多路复用减少连接开销,Protobuf二进制序列化体积比JSON小3-10倍,强类型接口定义文件(.proto)自动生成多语言客户端代码,原生支持流式通信和双向流式RPC。微服务内部通信使用gRPC能显著降低延迟和带宽消耗。

Protobuf接口定义与代码生成

Protobuf通过.proto文件定义服务接口和消息格式。以下是一个电商订单服务的完整定义:

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

package order.v1;
option go_package = "github.com/example/order-service/gen/order/v1;orderv1";

// 订单服务定义
service OrderService {
  // 创建订单
  rpc CreateOrder(CreateOrderRequest) returns (CreateOrderResponse);
  // 查询订单
  rpc GetOrder(GetOrderRequest) returns (GetOrderResponse);
  // 订单流式查询(服务端流)
  rpc StreamOrders(StreamOrdersRequest) returns (stream Order);
  // 双向流式通信
  rpc OrderChat(stream OrderMessage) returns (stream OrderMessage);
}

// 创建订单请求
message CreateOrderRequest {
  string user_id = 1;
  repeated OrderItem items = 2;
  Address shipping_address = 3;
  PaymentMethod payment_method = 4;
}

message OrderItem {
  string product_id = 1;
  int32 quantity = 2;
  double unit_price = 3;
}

message Address {
  string province = 1;
  string city = 2;
  string detail = 3;
  string phone = 4;
}

enum PaymentMethod {
  PAYMENT_UNSPECIFIED = 0;
  PAYMENT_WECHAT = 1;
  PAYMENT_ALIPAY = 2;
  PAYMENT_CARD = 3;
}

message CreateOrderResponse {
  string order_id = 1;
  double total_amount = 2;
  OrderStatus status = 3;
  string created_at = 4;
}

enum OrderStatus {
  STATUS_UNSPECIFIED = 0;
  STATUS_PENDING = 1;
  STATUS_PAID = 2;
  STATUS_SHIPPED = 3;
  STATUS_COMPLETED = 4;
  STATUS_CANCELLED = 5;
}

proto3语法的关键规则:所有字段编号必须唯一且不可变更(兼容性保证),枚举值0必须为UNSPECIFIED作为默认值,repeated字段对应列表类型,message嵌套无深度限制。字段编号1-15使用1字节编码,16-2047使用2字节,高频字段应使用小编号。

生成Go代码:

# 安装protoc和插件
go install google.golang.org/protobuf/cmd/protoc-gen-go@latest
go install google.golang.org/grpc/cmd/protoc-gen-go-grpc@latest

# 生成代码
protoc --go_out=. --go_opt=paths=source_relative   --go-grpc_out=. --go-grpc_opt=paths=source_relative   proto/order.proto

# 或使用buf工具管理proto
# buf.gen.yaml
version: v1
plugins:
  - plugin: buf.build/protocolbuffers/go
    out: gen
    opt: paths=source_relative
  - plugin: buf.build/grpc/go
    out: gen
    opt: paths=source_relative

服务端实现与拦截器中间件

gRPC服务端实现接口定义中声明的RPC方法,拦截器(Interceptor)用于在RPC调用前后插入横切逻辑,如日志、认证、限流、链路追踪。

package main

import (
    "context"
    "log"
    "net"
    "time"

    "google.golang.org/grpc"
    "google.golang.org/grpc/codes"
    "google.golang.org/grpc/status"
    "google.golang.org/grpc/middleware"
)

// 一元拦截器(Unary Interceptor)
func loggingInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (resp interface{}, err error) {
    start := time.Now()
    log.Printf("RPC开始: %s, 请求: %+v", info.FullMethod, req)

    resp, err = handler(ctx, req)

    log.Printf("RPC结束: %s, 耗时: %v, 错误: %v", info.FullMethod, time.Since(start), err)
    return resp, err
}

// 认证拦截器
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, "缺少metadata")
    }

    tokens := md.Get("authorization")
    if len(tokens) == 0 {
        return nil, status.Error(codes.Unauthenticated, "缺少认证token")
    }

    // 验证token
    userID, err := validateToken(tokens[0])
    if err != nil {
        return nil, status.Error(codes.Unauthenticated, "token无效")
    }

    // 将userID放入context传递给下游
    ctx = context.WithValue(ctx, "userID", userID)
    return handler(ctx, req)
}

// 限流拦截器
func rateLimitInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
    if !rateLimiter.Allow() {
        return nil, status.Error(codes.ResourceExhausted, "请求频率超限")
    }
    return handler(ctx, req)
}

// 流式拦截器
func streamLoggingInterceptor(srv interface{}, ss grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error {
    log.Printf("流式RPC开始: %s", info.FullMethod)
    err := handler(srv, ss)
    log.Printf("流式RPC结束: %s, 错误: %v", info.FullMethod, err)
    return err
}

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

    // 注册多个拦截器,按顺序执行
    s := grpc.NewServer(
        grpc.UnaryInterceptor(middleware.ChainUnaryServer(
            loggingInterceptor,
            authInterceptor,
            rateLimitInterceptor,
        )),
        grpc.StreamInterceptor(streamLoggingInterceptor),
    )

    // 注册服务实现
    pb.RegisterOrderServiceServer(s, &orderServiceImpl{})

    log.Println("gRPC服务启动,监听 :50051")
    if err := s.Serve(lis); err != nil {
        log.Fatalf("服务启动失败: %v", err)
    }
}

拦截器通过grpc.UnaryInterceptor注册,middleware.ChainUnaryServer将多个拦截器串联成链。执行顺序是从前到后依次进入,从后到前依次返回,类似洋葱模型。

客户端调用与连接管理

gRPC客户端使用连接池和负载均衡管理到服务端的连接。Go客户端示例:

package main

import (
    "context"
    "log"
    "time"

    "google.golang.org/grpc"
    "google.golang.org/grpc/credentials/insecure"
    "google.golang.org/grpc/balancer/roundrobin"
    "google.golang.org/grpc/metadata"
)

func main() {
    // 创建连接(支持多后端负载均衡)
    conn, err := grpc.NewClient(
        "dns:///order-service:50051",  // DNS解析多实例
        grpc.WithTransportCredentials(insecure.NewCredentials()),
        grpc.WithDefaultServiceConfig(`{"loadBalancingConfig":[{"round_robin":{}}]}`),
    )
    if err != nil {
        log.Fatalf("连接失败: %v", err)
    }
    defer conn.Close()

    client := pb.NewOrderServiceClient(conn)

    // 添加认证metadata
    ctx := metadata.AppendToOutgoingContext(context.Background(),
        "authorization", "Bearer eyJhbGc...",
    )

    // 设置超时
    ctx, cancel := context.WithTimeout(ctx, 5*time.Second)
    defer cancel()

    // 一元RPC调用
    resp, err := client.CreateOrder(ctx, &pb.CreateOrderRequest{
        UserId: "user-123",
        Items: []*pb.OrderItem{
            {ProductId: "prod-001", Quantity: 2, UnitPrice: 99.9},
        },
        PaymentMethod: pb.PaymentMethod_PAYMENT_WECHAT,
    })
    if err != nil {
        log.Fatalf("创建订单失败: %v", err)
    }
    log.Printf("订单创建成功: %s, 金额: %.2f", resp.OrderId, resp.TotalAmount)

    // 服务端流式调用
    stream, err := client.StreamOrders(ctx, &pb.StreamOrdersRequest{UserId: "user-123"})
    if err != nil {
        log.Fatalf("流式查询失败: %v", err)
    }

    for {
        order, err := stream.Recv()
        if err == io.EOF {
            break
        }
        if err != nil {
            log.Fatalf("接收流数据失败: %v", err)
        }
        log.Printf("订单: %s, 状态: %v", order.OrderId, order.Status)
    }
}

DNS解析模式下,gRPC客户端通过DNS A记录获取多个后端地址,使用round_robin负载均衡策略轮询发送请求。生产环境建议使用xDS(Envoy的发现服务)或Consul做服务发现。

错误处理与状态码规范

gRPC使用标准状态码传递错误信息,应用层错误应映射到合适的gRPC状态码:

var (
    // 使用status.Error返回带详细信息的错误
    ErrOrderNotFound    = status.Error(codes.NotFound, "订单不存在")
    ErrInvalidQuantity  = status.Error(codes.InvalidArgument, "商品数量必须大于0")
    ErrPermissionDenied = status.Error(codes.PermissionDenied, "无权操作此订单")
    ErrInternal         = status.Error(codes.Internal, "服务内部错误")
)

// 带详细信息的错误(使用ErrorInfo)
func createOrder(ctx context.Context, req *pb.CreateOrderRequest) (*pb.CreateOrderResponse, error) {
    if len(req.Items) == 0 {
        st := status.New(codes.InvalidArgument, "订单商品列表为空")
        // 添加详细错误信息
        details := &errdetails.BadRequest_FieldViolation{
            Field:       "items",
            Description: "至少需要一个商品",
        }
        st, _ = st.WithDetails(details)
        return nil, st.Err()
    }

    // 业务逻辑错误转换为gRPC状态码
    order, err := orderService.Create(ctx, req)
    if errors.Is(err, ErrOrderNotFound) {
        return nil, status.Error(codes.NotFound, err.Error())
    }
    if errors.Is(err, ErrPermissionDenied) {
        return nil, status.Error(codes.PermissionDenied, err.Error())
    }
    if err != nil {
        return nil, status.Error(codes.Internal, err.Error())
    }

    return &pb.CreateOrderResponse{
        OrderId:     order.ID,
        TotalAmount: order.Total,
        Status:      order.Status,
    }, nil
}

客户端通过status.Code(err)提取错误码,根据不同码值执行重试或降级逻辑。codes.Unavailable表示服务暂时不可用可重试,codes.InvalidArgument表示参数错误不应重试。

健康检查与优雅停机

gRPC健康检查协议(grpc.health.v1)是标准化的健康检查接口,负载均衡器通过它判断服务实例是否可接收流量:

// 实现健康检查服务
type healthServer struct {
    pb.UnimplementedHealthServer
    mu     sync.RWMutex
    status map[string]pb.HealthCheckResponse_ServingStatus
}

func (h *healthServer) Check(ctx context.Context, req *pb.HealthCheckRequest) (*pb.HealthCheckResponse, error) {
    h.mu.RLock()
    defer h.mu.RUnlock()
    status, ok := h.status[req.Service]
    if !ok {
        return &pb.HealthCheckResponse{Status: pb.HealthCheckResponse_UNKNOWN}, nil
    }
    return &pb.HealthCheckResponse{Status: status}, nil
}

// 优雅停机:先将状态设为NOT_SERVING,等待已有请求处理完成
func (s *Server) GracefulStop() {
    // 1. 健康检查状态设为NOT_SERVING
    s.health.SetServingStatus("", pb.HealthCheckResponse_NOT_SERVING)

    // 2. 等待负载均衡器将流量摘除
    time.Sleep(5 * time.Second)

    // 3. 停止接收新请求,等待已有请求完成
    s.grpcServer.GracefulStop()
}

优雅停机流程:设置健康检查为NOT_SERVING,负载均衡器下次探测时摘除该实例,等待已有RPC请求完成,最后关闭连接。配合Kubernetes的preStop hook和terminationGracePeriodSeconds使用。

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

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

相关推荐