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/