gRPC微服务通信实战:Protocol Buffers与流式RPC配置

gRPC是Google开源的高性能RPC框架,基于HTTP/2协议和Protocol Buffers序列化机制,在微服务架构中承担服务间通信职责。相比RESTful API的JSON文本序列化,gRPC使用二进制编码,传输体积更小、解析速度更快。gRPC原生支持流式通信、超时控制、拦截器链等特性,适合高并发微服务通信场景。

Protocol Buffers消息定义与代码生成

Protocol Buffers(proto3)是gRPC的接口定义语言(IDL),用于描述服务接口和消息结构。proto文件定义了服务端和客户端的通信契约,通过protoc编译器生成多语言代码。以下为订单服务的proto定义:

// order.proto
syntax = "proto3";

package order;
option go_package = "./proto/order";
option java_package = "com.example.order.grpc";
option java_multiple_files = true;

service OrderService {
  // 一元RPC(Unary RPC)
  rpc GetOrder(GetOrderRequest) returns (GetOrderResponse);
  rpc CreateOrder(CreateOrderRequest) returns (CreateOrderResponse);

  // 服务端流式RPC
  rpc StreamOrders(StreamOrdersRequest) returns (stream Order);

  // 客户端流式RPC
  rpc BatchCreateOrders(stream CreateOrderRequest) returns (BatchCreateResponse);

  // 双向流式RPC
  rpc OrderChat(stream OrderMessage) returns (stream OrderMessage);
}

message GetOrderRequest {
  string order_id = 1;
}

message GetOrderResponse {
  string order_id = 1;
  string user_id = 2;
  string product_name = 3;
  int32 quantity = 4;
  double total_price = 5;
  string status = 6;
  int64 created_at = 7;
}

message CreateOrderRequest {
  string user_id = 1;
  string product_id = 2;
  int32 quantity = 3;
}

message CreateOrderResponse {
  string order_id = 1;
  string status = 2;
}

message StreamOrdersRequest {
  string user_id = 1;
  string status = 2;
}

message Order {
  string order_id = 1;
  string status = 2;
  double amount = 3;
}

message BatchCreateResponse {
  int32 success_count = 1;
  int32 fail_count = 2;
  repeated string order_ids = 3;
}

message OrderMessage {
  string order_id = 1;
  string message = 2;
  string type = 3;
}

使用protoc生成Go语言代码:

# 安装protoc和插件
protoc --go_out=. --go-grpc_out=.   --go_opt=paths=source_relative   --go-grpc_opt=paths=source_relative   order.proto

gRPC服务端与客户端实现

Go语言中gRPC服务端实现需要注册服务、配置监听端口和拦截器。以下为服务端核心代码:

package main

import (
    "context"
    "log"
    "net"
    "google.golang.org/grpc"
    "google.golang.org/grpc/codes"
    "google.golang.org/grpc/status"
    pb "example.com/order/proto/order"
)

type OrderServer struct {
    pb.UnimplementedOrderServiceServer
    orders map[string]*pb.GetOrderResponse
}

func (s *OrderServer) GetOrder(ctx context.Context, req *pb.GetOrderRequest) (*pb.GetOrderResponse, error) {
    order, ok := s.orders[req.GetOrderId()]
    if !ok {
        return nil, status.Errorf(codes.NotFound, "order not found: %s", req.GetOrderId())
    }
    return order, nil
}

func (s *OrderServer) CreateOrder(ctx context.Context, req *pb.CreateOrderRequest) (*pb.CreateOrderResponse, error) {
    orderID := generateOrderID()
    s.orders[orderID] = &pb.GetOrderResponse{
        OrderId:     orderID,
        UserId:      req.GetUserId(),
        ProductName: "示例商品",
        Quantity:    req.GetQuantity(),
        Status:      "CREATED",
    }
    return &pb.CreateOrderResponse{
        OrderId: orderID,
        Status:  "CREATED",
    }, nil
}

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

    server := grpc.NewServer(
        grpc.UnaryInterceptor(loggingInterceptor),
        grpc.MaxRecvMsgSize(10*1024*1024),
    )
    pb.RegisterOrderServiceServer(server, &OrderServer{
        orders: make(map[string]*pb.GetOrderResponse),
    })

    log.Println("gRPC server listening on :50051")
    if err := server.Serve(lis); err != nil {
        log.Fatalf("failed to serve: %v", err)
    }
}

客户端调用:

func main() {
    conn, err := grpc.Dial("localhost:50051",
        grpc.WithTransportCredentials(insecure.NewCredentials()),
        grpc.WithDefaultCallOptions(
            grpc.MaxCallRecvMsgSize(10*1024*1024),
        ),
    )
    if err != nil {
        log.Fatal(err)
    }
    defer conn.Close()

    client := pb.NewOrderServiceClient(conn)

    // 一元RPC调用
    resp, err := client.CreateOrder(context.Background(), &pb.CreateOrderRequest{
        UserId:    "user_001",
        ProductId: "prod_001",
        Quantity:  2,
    })
    if err != nil {
        log.Fatal(err)
    }
    log.Printf("Created order: %s, status: %s", resp.GetOrderId(), resp.GetStatus())
}

流式RPC与双向通信配置

服务端流式RPC适用于批量查询场景,客户端发送一次请求,服务端返回多条数据流。客户端流式RPC适用于批量写入场景,客户端发送多条数据,服务端返回一次汇总结果。双向流式RPC支持实时双向通信,可用于聊天、协同编辑等场景。

// 服务端流式RPC实现
func (s *OrderServer) StreamOrders(req *pb.StreamOrdersRequest, stream pb.OrderService_StreamOrdersServer) error {
    for _, order := range s.orders {
        if order.GetStatus() == req.GetStatus() {
            if err := stream.Send(&pb.Order{
                OrderId: order.GetOrderId(),
                Status:  order.GetStatus(),
                Amount:  order.GetTotalPrice(),
            }); err != nil {
                return err
            }
        }
    }
    return nil
}

// 客户端接收流式数据
func streamOrders(client pb.OrderServiceClient, userID, status string) {
    stream, err := client.StreamOrders(context.Background(), &pb.StreamOrdersRequest{
        UserId: userID,
        Status: status,
    })
    if err != nil {
        log.Fatal(err)
    }
    for {
        order, err := stream.Recv()
        if err == io.EOF {
            break
        }
        if err != nil {
            log.Fatal(err)
        }
        log.Printf("Order: %s, Status: %s, Amount: %.2f",
            order.GetOrderId(), order.GetStatus(), order.GetAmount())
    }
}

拦截器链与超时重试机制

gRPC拦截器分为一元拦截器(UnaryInterceptor)和流拦截器(StreamInterceptor),用于实现日志记录、认证鉴权、指标采集等横切关注点。生产环境建议配置超时和重试策略:

// 一元拦截器:日志记录
func loggingInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
    start := time.Now()
    resp, err := handler(ctx, req)
    log.Printf("%s  duration=%v  error=%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) {
    md, ok := metadata.FromIncomingContext(ctx)
    if !ok {
        return nil, status.Errorf(codes.Unauthenticated, "missing metadata")
    }
    tokens := md.Get("authorization")
    if len(tokens) == 0 || !validateToken(tokens[0]) {
        return nil, status.Errorf(codes.Unauthenticated, "invalid token")
    }
    return handler(ctx, req)
}

// 客户端超时与重试
func main() {
    conn, _ := grpc.Dial("localhost:50051",
        grpc.WithTransportCredentials(insecure.NewCredentials()),
        grpc.WithDefaultCallOptions(
            grpc.MaxCallRecvMsgSize(10*1024*1024),
        ),
        grpc.WithUnaryInterceptor(retryInterceptor),
    )
    defer conn.Close()
}

func retryInterceptor(ctx context.Context, method string, req, reply interface{}, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error {
    ctx, cancel := context.WithTimeout(ctx, 5*time.Second)
    defer cancel()

    maxRetries := 3
    var lastErr error
    for i := 0; i < maxRetries; i++ {
        lastErr = invoker(ctx, method, req, reply, cc, opts...)
        if lastErr == nil {
            return nil
        }
        if status.Code(lastErr) == codes.Unavailable {
            time.Sleep(time.Duration(i+1) * time.Second)
            continue
        }
        return lastErr
    }
    return lastErr
}

gRPC服务治理还需要配合健康检查(grpc.health.v1)、负载均衡(客户端侧DNS或服务发现)和链路追踪(OpenTelemetry gRPC interceptor)。生产部署建议启用TLS加密通信,使用envoy或grpc-go内置的balancer实现客户端负载均衡。Protocol Buffers的版本兼容性需要遵循字段编号不复用、新增字段使用optional标记的规则,确保服务端和客户端的proto定义向前兼容。

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

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

相关推荐