gRPC服务通信实战:Protocol Buffers定义与流式RPC设计

gRPCProtocol Buffers基础

gRPC是Google开源的高性能RPC框架,基于HTTP/2协议传输,使用Protocol Buffers作为接口定义语言(IDL)和序列化格式。相比RESTful API的JSON序列化,Protobuf的二进制编码体积更小、解析速度更快,典型场景下序列化速度比JSON快3-10倍,数据体积小20-50%。

HTTP/2的多路复用特性允许在单个TCP连接上并发多个请求,避免了HTTP/1.1的队头阻塞问题。gRPC还支持流式通信,客户端和服务端可以在一个连接上持续发送和接收消息,适用于实时数据推送、大文件传输和日志流等场景。

安装protoc编译器和Go插件:

# 安装protoc
protoc --version

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

# 确保插件在PATH中
export PATH="$PATH:$(go env GOPATH)/bin"

Proto文件定义与代码生成

Protocol Buffers通过.proto文件定义服务接口和消息结构。以下是一个电商订单服务的完整proto定义。

syntax = "proto3";

package ecommerce.v1;
option go_package = "ecommerce/proto/v1;orderv1";

import "google/protobuf/timestamp.proto";
import "google/protobuf/empty.proto";

// 订单服务
service OrderService {
  // 一元调用:创建订单
  rpc CreateOrder(CreateOrderRequest) returns (Order);
  
  // 一元调用:查询订单
  rpc GetOrder(GetOrderRequest) returns (Order);
  
  // 服务端流:查询订单列表
  rpc ListOrders(ListOrdersRequest) returns (stream Order);
  
  // 客户端流:批量导入订单
  rpc BatchCreateOrders(stream CreateOrderRequest) returns (BatchResult);
  
  // 双向流:实时订单状态更新
  rpc OrderStream(stream OrderQuery) returns (stream OrderUpdate);
}

message Order {
  string id = 1;
  string user_id = 2;
  repeated OrderItem items = 3;
  double total_amount = 4;
  OrderStatus status = 5;
  google.protobuf.Timestamp created_at = 6;
}

message OrderItem {
  string product_id = 1;
  string product_name = 2;
  int32 quantity = 3;
  double price = 4;
}

enum OrderStatus {
  ORDER_STATUS_UNSPECIFIED = 0;
  ORDER_STATUS_PENDING = 1;
  ORDER_STATUS_PAID = 2;
  ORDER_STATUS_SHIPPED = 3;
  ORDER_STATUS_DELIVERED = 4;
  ORDER_STATUS_CANCELLED = 5;
}

message CreateOrderRequest {
  string user_id = 1;
  repeated OrderItem items = 2;
}

message GetOrderRequest {
  string id = 1;
}

message ListOrdersRequest {
  string user_id = 1;
  OrderStatus status = 2;
  int32 page_size = 3;
  string page_token = 4;
}

message OrderQuery {
  string order_id = 1;
  bool subscribe = 2;
}

message OrderUpdate {
  string order_id = 1;
  OrderStatus new_status = 2;
  string message = 3;
}

message BatchResult {
  int32 success_count = 1;
  int32 failure_count = 2;
  repeated string failed_ids = 3;
}

proto3语法中字段编号1-15使用1字节编码,16-2047使用2字节,高频字段应优先使用小编号。enum的第一个字段必须为0值作为默认值。repeated关键字表示数组类型,对应Go中的slice。import语句引入well-known types如Timestamp和Empty。

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

服务端实现与拦截器

gRPC服务端实现各RPC方法,并通过拦截器(Interceptor)实现统一的中间件逻辑如日志、认证、限流。

package main

import (
    "context"
    "log"
    "net"
    "google.golang.org/grpc"
    "google.golang.org/grpc/codes"
    "google.golang.org/grpc/status"
    "google.golang.org/grpc/credentials"
)

type OrderServer struct {
    orderv1.UnimplementedOrderServiceServer
    repo OrderRepository
}

// 一元调用
func (s *OrderServer) CreateOrder(ctx context.Context, 
    req *orderv1.CreateOrderRequest) (*orderv1.Order, error) {
    if req.UserId == "" {
        return nil, status.Error(codes.InvalidArgument, "user_id required")
    }
    order, err := s.repo.Create(ctx, req)
    if err != nil {
        return nil, status.Errorf(codes.Internal, "create failed: %v", err)
    }
    return order, nil
}

// 服务端流
func (s *OrderServer) ListOrders(req *orderv1.ListOrdersRequest,
    stream orderv1.OrderService_ListOrdersServer) error {
    orders, err := s.repo.ListByUser(stream.Context(), req)
    if err != nil {
        return status.Errorf(codes.Internal, "query failed: %v", err)
    }
    for _, order := range orders {
        if err := stream.Send(order); err != nil {
            return err
        }
    }
    return nil
}

// 双向流:实时订单状态推送
func (s *OrderServer) OrderStream(stream orderv1.OrderService_OrderStreamServer) error {
    var subscriptions = make(map[string]bool)
    
    for {
        query, err := stream.Recv()
        if err != nil {
            break
        }
        
        if query.Subscribe {
            subscriptions[query.OrderId] = true
        }
        
        // 推送订阅订单的状态更新
        for orderID := range subscriptions {
            update := s.repo.GetLatestUpdate(stream.Context(), orderID)
            if update != nil {
                if err := stream.Send(update); err != nil {
                    return err
                }
            }
        }
    }
    return nil
}

// 认证拦截器
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.Error(codes.Unauthenticated, "no metadata")
    }
    tokens := md.Get("authorization")
    if len(tokens) == 0 {
        return nil, status.Error(codes.Unauthenticated, "no token")
    }
    if err := validateToken(tokens[0]); err != nil {
        return nil, status.Error(codes.Unauthenticated, "invalid token")
    }
    return handler(ctx, req)
}

func main() {
    lis, err := net.Listen("tcp", ":50051")
    if err != nil {
        log.Fatal(err)
    }
    
    creds, _ := credentials.NewServerTLSFromFile(
        "cert/server.crt", "cert/server.key")
    
    server := grpc.NewServer(
        grpc.Creds(creds),
        grpc.UnaryInterceptor(authInterceptor),
        grpc.MaxRecvMsgSize(10*1024*1024),
    )
    orderv1.RegisterOrderServiceServer(server, &OrderServer{repo: NewOrderRepo()})
    
    log.Println("gRPC server on :50051")
    server.Serve(lis)
}

客户端调用与连接管理

func main() {
    creds := credentials.NewTLS(&tls.Config{
        InsecureSkipVerify: true, // 仅测试环境
    })
    
    conn, err := grpc.Dial("localhost:50051",
        grpc.WithTransportCredentials(creds),
        grpc.WithDefaultServiceConfig(`{
            "loadBalancingPolicy": "round_robin",
            "methodConfig": [{
                "name": [{"service": "ecommerce.v1.OrderService"}],
                "retryPolicy": {
                    "maxAttempts": 3,
                    "initialBackoff": "0.1s",
                    "maxBackoff": "1s",
                    "backoffMultiplier": 2.0,
                    "retryableStatusCodes": ["UNAVAILABLE"]
                }
            }]
        }`),
    )
    if err != nil {
        log.Fatal(err)
    }
    defer conn.Close()
    
    client := orderv1.NewOrderServiceClient(conn)
    
    // 设置超时和认证
    ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
    defer cancel()
    ctx = metadata.AppendToOutgoingContext(ctx, 
        "authorization", "Bearer "+token)
    
    // 一元调用
    order, err := client.CreateOrder(ctx, &orderv1.CreateOrderRequest{
        UserId: "user_123",
        Items: []*orderv1.OrderItem{
            {ProductId: "p1", Quantity: 2, Price: 99.9},
        },
    })
    
    // 服务端流调用
    stream, err := client.ListOrders(ctx, &orderv1.ListOrdersRequest{
        UserId: "user_123", PageSize: 20,
    })
    for {
        order, err := stream.Recv()
        if err == io.EOF { break }
        if err != nil { log.Fatal(err) }
        log.Printf("Order: %v", order.Id)
    }
}

gRPC客户端的连接管理需要关注连接池、重试策略和负载均衡。grpc.Dial配置了round_robin负载均衡策略和自动重试:当服务端返回UNAVAILABLE状态码时,客户端按指数退避策略重试,最多3次。连接复用gRPC底层的HTTP/2多路复用能力,单个连接支持并发100个以上的流式请求。

gRPC网关与REST兼容

并非所有客户端都支持gRPC协议,浏览器端尤其受限。grpc-gateway通过proto注解自动生成RESTful API代理,同时支持gRPC和HTTP调用。

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

service OrderService {
  rpc CreateOrder(CreateOrderRequest) returns (Order) {
    option (google.api.http) = {
      post: "/v1/orders"
      body: "*"
    };
  }
  rpc GetOrder(GetOrderRequest) returns (Order) {
    option (google.api.http) = {
      get: "/v1/orders/{id}"
    };
  }
}

# 生成网关代理代码
protoc --grpc-gateway_out=. --grpc-gateway_opt=paths=source_relative \
       proto/v1/order.proto

生成的gateway代理监听HTTP端口,将REST请求转换为gRPC调用转发给后端服务。前端通过标准fetch/axios调用REST API,后端服务间通信用gRPC,两者共用同一套proto定义,保持接口一致性。

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

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

相关推荐