Go语言gRPC微服务通信实战:Protocol Buffers定义与服务端客户端实现

gRPC是Google开源的高性能RPC框架,基于HTTP/2协议和Protocol Buffers序列化格式,在微服务架构中广泛用于服务间通信。相比REST API,gRPC在序列化效率、流式传输和跨语言支持方面优势显著。本文以Go语言为例,从proto文件定义、服务端实现、客户端调用到拦截器配置,给出gRPC微服务开发的完整实践。

Protocol Buffers消息与服务定义

Protocol Buffers(protobuf)是gRPC的数据序列化格式,通过.proto文件定义消息结构和服务接口。一个完整的gRPC服务定义包含消息类型、服务方法和流式传输声明。

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

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

// 用户消息定义
message User {
  int64 id = 1;
  string name = 2;
  string email = 3;
  int32 age = 4;
  repeated string roles = 5;
  google.protobuf.Timestamp created_at = 6;
}

message CreateUserRequest {
  string name = 1;
  string email = 2;
  int32 age = 3;
}

message GetUserRequest {
  int64 id = 1;
}

message ListUsersRequest {
  int32 page = 1;
  int32 page_size = 2;
}

message ListUsersResponse {
  repeated User users = 1;
  int32 total = 2;
}

message BatchUsersStream {
  string name = 1;
  int32 age = 2;
}

// 服务定义
service UserService {
  // 一元RPC(Unary RPC)
  rpc CreateUser(CreateUserRequest) returns (User);
  rpc GetUser(GetUserRequest) returns (User);
  rpc ListUsers(ListUsersRequest) returns (ListUsersResponse);

  // 服务端流式RPC
  rpc SubscribeUsers(ListUsersRequest) returns (stream User);

  // 客户端流式RPC
  rpc BatchCreateUsers(stream BatchUsersStream) returns (ListUsersResponse);

  // 双向流式RPC
  rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}

message ChatMessage {
  string user_id = 1;
  string content = 2;
}

生成Go代码:

protoc \
  --go_out=. --go_opt=paths=source_relative \
  --go-grpc_out=. --go-grpc_opt=paths=source_relative \
  proto/v1/user_service.proto

gRPC服务端实现

服务端需要实现proto文件中定义的所有RPC方法。以下是完整的UserService服务端实现:

package server

import (
	"context"
	"errors"
	"io"
	"sync"
	"time"

	pb "github.com/example/user-service/proto/v1"
	"google.golang.org/grpc/codes"
	"google.golang.org/grpc/status"
	"google.golang.org/protobuf/types/known/timestamppb"
)

type UserStore struct {
	mu    sync.RWMutex
	users map[int64]*pb.User
	nextID int64
}

func NewUserStore() *UserStore {
	return &UserStore{
		users:  make(map[int64]*pb.User),
		nextID: 1,
	}
}

type UserServiceServer struct {
	pb.UnimplementedUserServiceServer
	store *UserStore
}

func NewUserServiceServer() *UserServiceServer {
	return &UserServiceServer{
		store: NewUserStore(),
	}
}

// 一元RPC:创建用户
func (s *UserServiceServer) CreateUser(ctx context.Context, req *pb.CreateUserRequest) (*pb.User, error) {
	if req.GetName() == "" {
		return nil, status.Error(codes.InvalidArgument, "name is required")
	}

	s.store.mu.Lock()
	defer s.store.mu.Unlock()

	id := s.store.nextID
	s.store.nextID++

	user := &pb.User{
		Id:        id,
		Name:      req.Name,
		Email:     req.Email,
		Age:       req.Age,
		Roles:     []string{"user"},
		CreatedAt: timestamppb.Now(),
	}
	s.store.users[id] = user
	return user, nil
}

// 一元RPC:获取用户
func (s *UserServiceServer) GetUser(ctx context.Context, req *pb.GetUserRequest) (*pb.User, error) {
	s.store.mu.RLock()
	defer s.store.mu.RUnlock()

	user, ok := s.store.users[req.Id]
	if !ok {
		return nil, status.Errorf(codes.NotFound, "user %d not found", req.Id)
	}
	return user, nil
}

// 一元RPC:分页查询
func (s *UserServiceServer) ListUsers(ctx context.Context, req *pb.ListUsersRequest) (*pb.ListUsersResponse, error) {
	s.store.mu.RLock()
	defer s.store.mu.RUnlock()

	page := int(req.Page)
	if page < 1 {
		page = 1
	}
	pageSize := int(req.PageSize)
	if pageSize < 1 || pageSize > 100 {
		pageSize = 10
	}

	all := make([]*pb.User, 0, len(s.store.users))
	for _, u := range s.store.users {
		all = append(all, u)
	}

	start := (page - 1) * pageSize
	end := start + pageSize
	if start >= len(all) {
		return &pb.ListUsersResponse{Users: []*pb.User{}, Total: int32(len(all))}, nil
	}
	if end > len(all) {
		end = len(all)
	}

	return &pb.ListUsersResponse{
		Users: all[start:end],
		Total: int32(len(all)),
	}, nil
}

// 服务端流式RPC
func (s *UserServiceServer) SubscribeUsers(req *pb.ListUsersRequest, stream pb.UserService_SubscribeUsersServer) error {
	s.store.mu.RLock()
	defer s.store.mu.RUnlock()

	for _, user := range s.store.users {
		if err := stream.Send(user); err != nil {
			return status.Errorf(codes.Internal, "stream send failed: %v", err)
		}
		time.Sleep(100 * time.Millisecond)
	}
	return nil
}

// 客户端流式RPC
func (s *UserServiceServer) BatchCreateUsers(stream pb.UserService_BatchCreateUsersServer) error {
	count := 0
	for {
		req, err := stream.Recv()
		if errors.Is(err, io.EOF) {
			return stream.SendAndClose(&pb.ListUsersResponse{
				Total: int32(count),
			})
		}
		if err != nil {
			return status.Errorf(codes.Internal, "recv failed: %v", err)
		}

		s.store.mu.Lock()
		user := &pb.User{
			Id:        s.store.nextID,
			Name:      req.Name,
			Age:       req.Age,
			CreatedAt: timestamppb.Now(),
		}
		s.store.users[user.Id] = user
		s.store.nextID++
		s.store.mu.Unlock()

		count++
	}
}

gRPC服务端启动与配置

package main

import (
	"log"
	"net"

	pb "github.com/example/user-service/proto/v1"
	"github.com/example/user-service/server"
	"google.golang.org/grpc"
	"google.golang.org/grpc/reflection"
)

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

	grpcServer := grpc.NewServer(
		grpc.UnaryInterceptor(unaryLoggingInterceptor),
		grpc.StreamInterceptor(streamLoggingInterceptor),
		grpc.MaxRecvMsgSize(10 * 1024 * 1024), // 10MB
	)

	pb.RegisterUserServiceServer(grpcServer, server.NewUserServiceServer())
	reflection.Register(grpcServer) // 启用反射,方便调试

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

拦截器实现

拦截器是gRPC的中间件机制,用于日志记录、认证鉴权、链路追踪等横切关注点:

package main

import (
	"context"
	"log"
	"time"

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

// 一元拦截器
func unaryLoggingInterceptor(
	ctx context.Context,
	req interface{},
	info *grpc.UnaryServerInfo,
	handler grpc.UnaryHandler,
) (interface{}, error) {
	start := time.Now()

	// 从metadata获取认证token
	if md, ok := metadata.FromIncomingContext(ctx); ok {
		tokens := md.Get("authorization")
		if len(tokens) == 0 {
			return nil, status.Error(codes.Unauthenticated, "missing auth token")
		}
	}

	resp, err := handler(ctx, req)
	log.Printf("method=%s duration=%v err=%v", info.FullMethod, time.Since(start), err)
	return resp, err
}

// 流式拦截器
func streamLoggingInterceptor(
	srv interface{},
	stream grpc.ServerStream,
	info *grpc.StreamServerInfo,
	handler grpc.StreamHandler,
) error {
	start := time.Now()
	err := handler(srv, stream)
	log.Printf("stream method=%s duration=%v err=%v", info.FullMethod, time.Since(start), err)
	return err
}

gRPC客户端实现

package client

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

	pb "github.com/example/user-service/proto/v1"
	"google.golang.org/grpc"
	"google.golang.org/grpc/credentials/insecure"
	"google.golang.org/grpc/metadata"
)

func RunClient() {
	conn, err := grpc.NewClient("localhost:50051",
		grpc.WithTransportCredentials(insecure.NewCredentials()),
		grpc.WithDefaultCallOptions(grpc.MaxCallRecvMsgSize(10*1024*1024)),
	)
	if err != nil {
		log.Fatalf("failed to connect: %v", err)
	}
	defer conn.Close()

	client := pb.NewUserServiceClient(conn)

	// 设置带认证的context
	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
	defer cancel()
	ctx = metadata.AppendToOutgoingContext(ctx, "authorization", "Bearer token123")

	// 一元RPC调用
	user, err := client.CreateUser(ctx, &pb.CreateUserRequest{
		Name:  "张三",
		Email: "zhangsan@example.com",
		Age:   30,
	})
	if err != nil {
		log.Fatalf("CreateUser failed: %v", err)
	}
	log.Printf("Created user: %v", user)

	// 服务端流式调用
	stream, err := client.SubscribeUsers(ctx, &pb.ListUsersRequest{Page: 1, PageSize: 100})
	if err != nil {
		log.Fatalf("SubscribeUsers failed: %v", err)
	}
	for {
		u, err := stream.Recv()
		if errors.Is(err, io.EOF) {
			break
		}
		if err != nil {
			log.Fatalf("stream recv failed: %v", err)
		}
		log.Printf("Received user: %s", u.Name)
	}
}

gRPC健康检查配置

gRPC健康检查协议是服务发现和负载均衡的基础,Kubernetes等编排平台依赖健康检查判断服务可用性:

package main

import (
	"google.golang.org/grpc"
	healthpb "google.golang.org/grpc/health/grpc_health_v1"
)

// 注册健康检查服务
healthServer := &health.Server{}
healthServer.SetServingStatus("user.v1.UserService", healthpb.HealthCheckResponse_SERVING)
healthpb.RegisterHealthServer(grpcServer, healthServer)

客户端健康检查:

healthConn, _ := grpc.NewClient("localhost:50051",
	grpc.WithTransportCredentials(insecure.NewCredentials()))
healthClient := healthpb.NewHealthClient(healthConn)

resp, err := healthClient.Check(ctx, &healthpb.HealthCheckRequest{
	Service: "user.v1.UserService",
})
if resp.Status == healthpb.HealthCheckResponse_SERVING {
	log.Println("service is healthy")
}

gRPC vs REST选型建议

gRPC适合内部微服务间的高频通信场景,序列化效率比JSON高5-10倍,支持双向流式传输。REST适合对外暴露的API,浏览器直接调用,生态工具更丰富。实际项目中常采用gRPC处理内部通信,同时通过grpc-gateway生成REST API对外提供服务,兼顾性能和易用性。

对于需要浏览器端直接调用gRPC的场景,可以使用gRPC-Web方案,通过Envoy代理将HTTP/1.1请求转码为gRPC。对于移动端,gRPC原生支持Android和iOS,相比REST减少约60%的网络传输量。

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

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

相关推荐