gRPC服务定义与Go语言实现:Protocol Buffers与流式通信实战

RESTful API 基于 HTTP/1.1 文本协议,JSON 序列化开销大且缺乏强类型约束。gRPC 基于 HTTP/2 和 Protocol Buffers 二进制编码,在微服务间通信场景中提供更高吞吐、更低延迟和更强的类型安全。本文从 Protobuf 定义到 Go 语言 gRPC 服务实现,覆盖 unary 调用和流式通信的完整实践。

Protocol Buffers消息定义与代码生成

gRPC 接口使用 Protocol Buffers IDL(Interface Definition Language)定义。以下是一个订单服务的完整 proto 定义:

syntax = "proto3";

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

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

// 订单服务定义
service OrderService {
  // Unary RPC:创建订单
  rpc CreateOrder(CreateOrderRequest) returns (CreateOrderResponse);
  
  // Unary RPC:查询订单详情
  rpc GetOrder(GetOrderRequest) returns (GetOrderResponse);
  
  // Server Streaming:订阅订单状态变更
  rpc StreamOrderUpdates(StreamOrderUpdatesRequest) returns (stream OrderUpdate);
  
  // Client Streaming:批量上传订单
  rpc BatchCreateOrders(stream CreateOrderRequest) returns (BatchCreateOrdersResponse);
  
  // Bidirectional Streaming:实时订单对话
  rpc OrderChat(stream OrderMessage) returns (stream OrderMessage);
}

message CreateOrderRequest {
  string user_id = 1;
  repeated OrderItem items = 2;
  string shipping_address = 3;
  PaymentMethod payment_method = 4;
}

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

enum PaymentMethod {
  PAYMENT_METHOD_UNSPECIFIED = 0;
  PAYMENT_METHOD_ALIPAY = 1;
  PAYMENT_METHOD_WECHAT = 2;
  PAYMENT_METHOD_CARD = 3;
}

message CreateOrderResponse {
  string order_id = 1;
  google.protobuf.Timestamp created_at = 2;
  double total_amount = 3;
}

message GetOrderRequest {
  string order_id = 1;
  google.protobuf.FieldMask field_mask = 2;
}

message GetOrderResponse {
  Order order = 1;
}

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

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 StreamOrderUpdatesRequest {
  string user_id = 1;
  repeated string order_ids = 2;
}

message OrderUpdate {
  string order_id = 1;
  OrderStatus new_status = 2;
  google.protobuf.Timestamp updated_at = 3;
}

message BatchCreateOrdersResponse {
  int32 success_count = 1;
  int32 failure_count = 2;
  repeated string failed_order_indices = 3;
}

message OrderMessage {
  string session_id = 1;
  string content = 2;
  string user_id = 3;
}

使用 protoc 编译器生成 Go 代码:

# 安装 protoc 和 Go 插件
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 \
  api/v1/order.proto

生成的代码包含消息结构体(OrderCreateOrderRequest 等)和 gRPC 服务接口(OrderServiceServerOrderServiceClient)。枚举值 0 始终是 _UNSPECIFIED,proto3 中枚举零值代表未指定,用于区分”未设置”和”零值”。

Go语言gRPC服务端实现

实现 OrderServiceServer 接口:

package server

import (
	"context"
	"log"
	"sync"
	"time"

	"google.golang.org/grpc"
	"google.golang.org/grpc/codes"
	"google.golang.org/grpc/status"
	"google.golang.org/protobuf/types/known/timestamppb"

	pb "github.com/example/order-service/api/v1"
)

type OrderServer struct {
	pb.UnimplementedOrderServiceServer
	mu     sync.RWMutex
	orders map[string]*pb.Order
}

func NewOrderServer() *OrderServer {
	return &OrderServer{
		orders: make(map[string]*pb.Order),
	}
}

// Unary RPC: 创建订单
func (s *OrderServer) CreateOrder(ctx context.Context, req *pb.CreateOrderRequest) (*pb.CreateOrderResponse, error) {
	if req.UserId == "" {
		return nil, status.Error(codes.InvalidArgument, "user_id is required")
	}
	if len(req.Items) == 0 {
		return nil, status.Error(codes.InvalidArgument, "items cannot be empty")
	}

	orderID := generateOrderID()
	var totalAmount float64
	for _, item := range req.Items {
		totalAmount += float64(item.Quantity) * item.UnitPrice
	}

	order := &pb.Order{
		Id:        orderID,
		UserId:    req.UserId,
		Items:     req.Items,
		Status:    pb.OrderStatus_ORDER_STATUS_PENDING,
		CreatedAt: timestamppb.Now(),
		UpdatedAt: timestamppb.Now(),
	}

	s.mu.Lock()
	s.orders[orderID] = order
	s.mu.Unlock()

	return &pb.CreateOrderResponse{
		OrderId:     orderID,
		CreatedAt:   timestamppb.Now(),
		TotalAmount: totalAmount,
	}, nil
}

// Server Streaming: 订阅订单状态变更
func (s *OrderServer) StreamOrderUpdates(req *pb.StreamOrderUpdatesRequest, stream pb.OrderService_StreamOrderUpdatesServer) error {
	// 模拟订单状态变更推送
	ticker := time.NewTicker(5 * time.Second)
	defer ticker.Stop()

	statuses := []pb.OrderStatus{
		pb.OrderStatus_ORDER_STATUS_PAID,
		pb.OrderStatus_ORDER_STATUS_SHIPPED,
		pb.OrderStatus_ORDER_STATUS_DELIVERED,
	}

	for i := 0; i < len(statuses); i++ {
		select {
		case <-stream.Context().Done():
			return stream.Context().Err()
		case <-ticker.C:
			for _, orderID := range req.OrderIds {
				err := stream.Send(&pb.OrderUpdate{
					OrderId:    orderID,
					NewStatus:  statuses[i],
					UpdatedAt: timestamppb.Now(),
				})
				if err != nil {
					return status.Errorf(codes.Internal, "failed to send update: %v", err)
				}
			}
		}
	}
	return nil
}

// Client Streaming: 批量上传订单
func (s *OrderServer) BatchCreateOrders(stream pb.OrderService_BatchCreateOrdersServer) error {
	successCount := 0
	failureCount := 0
	var failedIndices []string

	for {
		req, err := stream.Recv()
		if err != nil {
			break // 客户端结束发送
		}

		// 处理每个订单创建请求
		_, createErr := s.CreateOrder(stream.Context(), req)
		if createErr != nil {
			failureCount++
			failedIndices = append(failedIndices, req.UserId)
		} else {
			successCount++
		}
	}

	return stream.SendAndClose(&pb.BatchCreateOrdersResponse{
		SuccessCount:      int32(successCount),
		FailureCount:      int32(failureCount),
		FailedOrderIndices: failedIndices,
	})
}

// Bidirectional Streaming: 实时订单对话
func (s *OrderServer) OrderChat(stream pb.OrderService_OrderChatServer) error {
	for {
		msg, err := stream.Recv()
		if err != nil {
			return err
		}

		log.Printf("收到消息: session=%s user=%s content=%s",
			msg.SessionId, msg.UserId, msg.Content)

		// 处理并回复
		reply := &pb.OrderMessage{
			SessionId: msg.SessionId,
			UserId:    "system",
			Content:   "已收到您的订单咨询:" + msg.Content,
		}

		if err := stream.Send(reply); err != nil {
			return err
		}
	}
}

func generateOrderID() string {
	return time.Now().Format("20060102150405") + "-" + randomHex(4)
}

gRPC服务启动与拦截器配置

package main

import (
	"fmt"
	"log"
	"net"
	"time"

	"google.golang.org/grpc"
	"google.golang.org/grpc/keepalive"
	"google.golang.org/grpc/reflection"

	pb "github.com/example/order-service/api/v1"
	"github.com/example/order-service/internal/server"
)

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

	// 创建带拦截器的 gRPC 服务
	srv := grpc.NewServer(
		grpc.UnaryInterceptor(unaryLoggingInterceptor),
		grpc.StreamInterceptor(streamLoggingInterceptor),
		grpc.KeepaliveParams(keepalive.ServerParameters{
			MaxConnectionIdle:     5 * time.Minute,
			MaxConnectionAge:      30 * time.Minute,
			MaxConnectionAgeGrace: 5 * time.Second,
			Time:                  30 * time.Second,
			Timeout:               10 * time.Second,
		}),
	)

	pb.RegisterOrderServiceServer(srv, server.NewOrderServer())

	// 启用 gRPC 反射(调试用,生产环境可关闭)
	reflection.Register(srv)

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

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

	log.Printf("unary: %s  duration=%s  error=%v",
		info.FullMethod, duration, err)
	return resp, err
}

func streamLoggingInterceptor(srv interface{}, ss grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error {
	start := time.Now()
	err := handler(srv, ss)
	duration := time.Since(start)

	log.Printf("stream: %s  duration=%s  error=%v",
		info.FullMethod, duration, err)
	return err
}

Keepalive 参数配置连接保活策略:MaxConnectionIdle 空闲连接最长存活时间,MaxConnectionAge 连接最大生命周期,超过后服务端主动关闭连接触发客户端重连,防止长连接负载不均衡。

gRPC客户端调用与连接管理

package main

import (
	"context"
	"io"
	"log"
	"time"

	"google.golang.org/grpc"
	"google.golang.org/grpc/credentials/insecure"

	pb "github.com/example/order-service/api/v1"
)

func main() {
	// 建立连接(带连接池和重试)
	conn, err := grpc.Dial("localhost:50051",
		grpc.WithTransportCredentials(insecure.NewCredentials()),
		grpc.WithDefaultServiceConfig(`{"loadBalancingPolicy":"round_robin"}`),
	)
	if err != nil {
		log.Fatalf("连接失败: %v", err)
	}
	defer conn.Close()

	client := pb.NewOrderServiceClient(conn)

	// Unary 调用
	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
	defer cancel()

	resp, err := client.CreateOrder(ctx, &pb.CreateOrderRequest{
		UserId: "user-001",
		Items: []*pb.OrderItem{
			{ProductId: "prod-1", Quantity: 2, UnitPrice: 99.9},
		},
		ShippingAddress: "北京市朝阳区",
		PaymentMethod:   pb.PaymentMethod_PAYMENT_METHOD_ALIPAY,
	})
	if err != nil {
		log.Fatalf("创建订单失败: %v", err)
	}
	log.Printf("订单创建成功: ID=%s, 金额=%.2f", resp.OrderId, resp.TotalAmount)

	// Server Streaming 调用
	stream, err := client.StreamOrderUpdates(context.Background(), &pb.StreamOrderUpdatesRequest{
		UserId:   "user-001",
		OrderIds: []string{resp.OrderId},
	})
	if err != nil {
		log.Fatalf("订阅失败: %v", err)
	}

	for {
		update, err := stream.Recv()
		if err == io.EOF {
			break
		}
		if err != nil {
			log.Fatalf("接收更新失败: %v", err)
		}
		log.Printf("订单 %s 状态变更: %s", update.OrderId, update.NewStatus)
	}
}

grpc.Dial 使用 round_robin 负载均衡策略,适用于多实例部署场景。客户端连接为长连接复用,无需每次请求建立新连接。对于短生命周期的 CLI 工具,使用 defer conn.Close() 确保连接正确关闭。

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

(0)
小编小编
上一篇 1天前
下一篇 1天前

相关推荐