微服务架构中服务间通信的选型直接影响系统吞吐量和延迟。gRPC基于HTTP/2和Protocol Buffers,在序列化效率和传输性能上显著优于REST+JSON方案。Protobuf二进制序列化体积比JSON小3-10倍,HTTP/2多路复用消除了连接建立开销,使gRPC在高并发设计和分布式事务场景中获得更优的端到端延迟表现。
gRPC与REST通信协议对比与选型
- 序列化格式:gRPC使用Protobuf二进制编码,JSON体积平均大5倍
- 传输协议:gRPC基于HTTP/2支持多路复用和流控,REST基于HTTP/1.1需连接池
- 接口定义:gRPC通过.proto文件定义强类型接口,自动生成多语言客户端代码
- 流式通信:gRPC原生支持双向流,REST需通过WebSocket或SSE模拟
- 浏览器兼容:gRPC需要gRPC-Web代理,REST可直接被浏览器调用
实际项目中,内部服务间通信使用gRPC,对外暴露的API网关使用REST/gRPC-Gateway转换。
Protobuf接口定义与代码生成
// proto/user_service.proto
syntax = "proto3";
package userservice.v1;
option go_package = "github.com/example/proto/userservice/v1";
import "google/protobuf/timestamp.proto";
service UserService {
rpc GetUser(GetUserRequest) returns (GetUserResponse);
rpc ListUsers(ListUsersRequest) returns (stream User);
rpc BatchCreateUsers(stream CreateUserRequest) returns (BatchCreateResponse);
rpc ChatStream(stream ChatMessage) returns (stream ChatMessage);
}
message User {
int64 id = 1;
string username = 2;
string email = 3;
UserStatus status = 4;
google.protobuf.Timestamp created_at = 5;
repeated string roles = 6;
map<string, string> metadata = 7;
}
enum UserStatus {
USER_STATUS_UNSPECIFIED = 0;
USER_STATUS_ACTIVE = 1;
USER_STATUS_INACTIVE = 2;
USER_STATUS_SUSPENDED = 3;
}
message GetUserRequest { int64 id = 1; }
message GetUserResponse { User user = 1; }
message ListUsersRequest {
int32 page_size = 1;
string page_token = 2;
UserStatus status_filter = 3;
}
message CreateUserRequest {
string username = 1;
string email = 2;
string password = 3;
}
message BatchCreateResponse {
repeated int64 created_ids = 1;
int32 failed_count = 2;
}
message ChatMessage {
int64 user_id = 1;
string content = 2;
google.protobuf.Timestamp sent_at = 3;
}
# 生成Go代码
protoc --go_out=. --go_opt=paths=source_relative \
--go-grpc_out=. --go-grpc_opt=paths=source_relative \
proto/user_service.proto
# 生成Python代码
python -m grpc_tools.protoc -I. --python_out=. --grpc_python_out=. \
proto/user_service.proto
Go语言gRPC服务端实现
package main
import (
"context"
"log"
"net"
"sync"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/credentials"
"google.golang.org/grpc/keepalive"
"google.golang.org/grpc/metadata"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/types/known/timestamppb"
pb "github.com/example/proto/userservice/v1"
)
type userServer struct {
pb.UnimplementedUserServiceServer
users map[int64]*pb.User
mu sync.RWMutex
}
func (s *userServer) GetUser(ctx context.Context, req *pb.GetUserRequest) (*pb.GetUserResponse, error) {
s.mu.RLock()
defer s.mu.RUnlock()
user, ok := s.users[req.GetId()]
if !ok {
return nil, status.Errorf(codes.NotFound, "user %d not found", req.GetId())
}
return &pb.GetUserResponse{User: user}, nil
}
func (s *userServer) ListUsers(req *pb.ListUsersRequest, stream pb.UserService_ListUsersServer) error {
s.mu.RLock()
defer s.mu.RUnlock()
count := 0
for _, user := range s.users {
if req.StatusFilter != pb.UserStatus_USER_STATUS_UNSPECIFIED &&
user.Status != req.StatusFilter {
continue
}
if err := stream.Context().Err(); err != nil {
return err
}
if err := stream.Send(user); err != nil {
return err
}
count++
if req.PageSize > 0 && int32(count) >= req.PageSize {
break
}
time.Sleep(100 * time.Millisecond)
}
return nil
}
func (s *userServer) ChatStream(stream pb.UserService_ChatStreamServer) error {
for {
msg, err := stream.Recv()
if err != nil {
return err
}
reply := &pb.ChatMessage{
UserId: msg.UserId,
Content: "Echo: " + msg.Content,
SentAt: timestamppb.Now(),
}
if err := stream.Send(reply); err != nil {
return err
}
}
}
func main() {
lis, _ := net.Listen("tcp", ":50051")
creds, _ := credentials.NewServerTLSFromFile("cert/server.crt", "cert/server.key")
s := grpc.NewServer(
grpc.Creds(creds),
grpc.KeepaliveParams(keepalive.ServerParameters{
MaxConnectionIdle: 5 * time.Minute,
MaxConnectionAge: 30 * time.Minute,
MaxConnectionAgeGrace: 5 * time.Second,
Time: 30 * time.Second,
Timeout: 10 * time.Second,
}),
grpc.MaxRecvMsgSize(16 * 1024 * 1024),
)
pb.RegisterUserServiceServer(s, &userServer{users: make(map[int64]*pb.User)})
log.Println("gRPC server listening on :50051")
s.Serve(lis)
}
拦截器Interceptor与中间件链
// 一元拦截器:日志鉴权
func loggingInterceptor(
ctx context.Context,
req interface{},
info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler,
) (resp interface{}, err error) {
start := time.Now()
if err := authenticate(ctx, info.FullMethod); err != nil {
return nil, err
}
resp, err = handler(ctx, req)
log.Printf("method=%s duration=%s error=%v",
info.FullMethod, time.Since(start), err)
return resp, err
}
// 恢复拦截器:panic捕获
func recoveryInterceptor(
ctx context.Context,
req interface{},
info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler,
) (resp interface{}, err error) {
defer func() {
if r := recover(); r != nil {
log.Printf("panic recovered: %v", r)
err = status.Errorf(codes.Internal, "internal error")
}
}()
return handler(ctx, req)
}
func authenticate(ctx context.Context, method string) error {
md, ok := metadata.FromIncomingContext(ctx)
if !ok {
return status.Error(codes.Unauthenticated, "no metadata")
}
tokens := md.Get("authorization")
if len(tokens) == 0 {
return status.Error(codes.Unauthenticated, "no token")
}
if !validateToken(tokens[0]) {
return status.Error(codes.Unauthenticated, "invalid token")
}
return nil
}
// 链式拦截器
s := grpc.NewServer(
grpc.ChainUnaryInterceptor(recoveryInterceptor, loggingInterceptor),
)
Python gRPC客户端实现
import grpc
import user_service_pb2 as pb
import user_service_pb2_grpc as pb_grpc
import time
channel = grpc.insecure_channel(
'localhost:50051',
options=[
('grpc.max_receive_message_length', 16 * 1024 * 1024),
('grpc.keepalive_time_ms', 30000),
('grpc.keepalive_timeout_ms', 10000),
('grpc.enable_retries', 1),
]
)
stub = pb_grpc.UserServiceStub(channel)
# 一元调用
response = stub.GetUser(
pb.GetUserRequest(id=1),
timeout=5,
metadata=[('authorization', 'Bearer token_xxx')]
)
print(f"用户: {response.user.username}, 状态: {response.user.status}")
# 服务端流调用
for user in stub.ListUsers(
pb.ListUsersRequest(page_size=10, status_filter=pb.UserStatus.USER_STATUS_ACTIVE),
timeout=30
):
print(f" - {user.id}: {user.username}")
# 双向流调用
def request_generator():
for i in range(10):
yield pb.ChatMessage(user_id=1, content=f"消息 {i}")
time.sleep(0.5)
responses = stub.ChatStream(request_generator(), timeout=60)
for reply in responses:
print(f"收到回复: {reply.content}")
gRPC服务治理与健康检查
type healthServer struct {
mu sync.Mutex
status map[string]healthpb.HealthCheckResponse_ServingStatus
}
func (h *healthServer) Check(ctx context.Context, req *healthpb.HealthCheckRequest) (*healthpb.HealthCheckResponse, error) {
h.mu.Lock()
defer h.mu.Unlock()
service := req.GetService()
if service == "" {
service = "overall"
}
status, ok := h.status[service]
if !ok {
return &healthpb.HealthCheckResponse{
Status: healthpb.HealthCheckResponse_SERVING_STATUS_UNKNOWN,
}, nil
}
return &healthpb.HealthCheckResponse{Status: status}, nil
}
// 注册健康检查
healthpb.RegisterHealthServer(s, &healthServer{
status: map[string]healthpb.HealthCheckResponse_ServingStatus{
"": healthpb.HealthCheckResponse_SERVING,
"userservice.UserService": healthpb.HealthCheckResponse_SERVING,
},
})
gRPC客户端默认复用HTTP/2连接,单连接支持100个并发流。生产环境建议根据并发量配置连接池大小,同时启用keepalive探活避免NAT超时断开。大消息体场景需调整MaxRecvMsgSize和流控窗口。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/grpc-fu-wu-tong-xin-xie-yi-yu-protobuf-xu-lie-hua-zai-wei/