gRPC是Google开源的高性能RPC框架,基于HTTP/2协议传输,使用Protocol Buffers作为接口定义语言和序列化格式。相比RESTful JSON API,gRPC在吞吐量、延迟和类型安全方面有显著优势,适合微服务间的高频内部通信。Go语言作为gRPC的原生支持语言之一,拥有完善的工具链和生态。
Protobuf接口定义与代码生成
gRPC服务的接口通过.proto文件定义,使用protoc编译器生成多语言客户端和服务端代码。先定义一个用户服务的proto文件:
// proto/user_service.proto
syntax = "proto3";
package user.v1;
option go_package = "github.com/example/proto/user/v1;userv1";
// 用户服务定义
service UserService {
rpc GetUser(GetUserRequest) returns (GetUserResponse);
rpc ListUsers(ListUsersRequest) returns (ListUsersResponse);
rpc CreateUser(CreateUserRequest) returns (CreateUserResponse);
rpc UpdateUser(UpdateUserRequest) returns (UpdateUserResponse);
rpc DeleteUser(DeleteUserRequest) returns (DeleteUserResponse);
}
message GetUserRequest {
int64 id = 1;
}
message GetUserResponse {
int64 id = 1;
string name = 2;
string email = 3;
string phone = 4;
int32 status = 5;
int64 created_at = 6;
}
message ListUsersRequest {
int32 page = 1;
int32 page_size = 2;
string keyword = 3;
}
message ListUsersResponse {
repeated GetUserResponse users = 1;
int32 total = 2;
}
message CreateUserRequest {
string name = 1;
string email = 2;
string phone = 3;
}
message CreateUserResponse {
int64 id = 1;
}
message UpdateUserRequest {
int64 id = 1;
string name = 2;
string email = 3;
string phone = 4;
}
message UpdateUserResponse {
bool success = 1;
}
message DeleteUserRequest {
int64 id = 1;
}
message DeleteUserResponse {
bool success = 1;
}
使用protoc生成Go代码:
# 安装protoc和Go插件
protoc --go_out=. --go-grpc_out=. --go_opt=paths=source_relative --go-grpc_opt=paths=source_relative proto/user_service.proto
# 或使用buf工具(推荐)
# buf generate
Go服务端实现与拦截器中间件
gRPC服务端实现proto定义的接口,通过拦截器(Interceptor)实现日志、认证、限流等横切关注点。以下是完整的服务端实现:
package main
import (
"context"
"log"
"net"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
pb "github.com/example/proto/user/v1"
)
type UserServer struct {
pb.UnimplementedUserServiceServer
// 注入数据库、缓存等依赖
db UserRepo
}
// GetUser 实现
func (s *UserServer) GetUser(
ctx context.Context, req *pb.GetUserRequest,
) (*pb.GetUserResponse, error) {
user, err := s.db.FindByID(ctx, req.Id)
if err != nil {
return nil, status.Errorf(codes.NotFound, "user not found: %v", err)
}
return &pb.GetUserResponse{
Id: user.ID,
Name: user.Name,
Email: user.Email,
Phone: user.Phone,
Status: user.Status,
CreatedAt: user.CreatedAt.Unix(),
}, nil
}
// ListUsers 实现
func (s *UserServer) ListUsers(
ctx context.Context, req *pb.ListUsersRequest,
) (*pb.ListUsersResponse, error) {
users, total, err := s.db.List(ctx, req.Page, req.PageSize, req.Keyword)
if err != nil {
return nil, status.Errorf(codes.Internal, "list failed: %v", err)
}
resp := &pb.ListUsersResponse{
Total: int32(total),
}
for _, u := range users {
resp.Users = append(resp.Users, &pb.GetUserResponse{
Id: u.ID,
Name: u.Name,
Email: u.Email,
})
}
return resp, nil
}
// 日志拦截器
func loggingInterceptor(
ctx context.Context, req interface{},
info *grpc.UnaryServerInfo, handler grpc.UnaryHandler,
) (resp interface{}, err error) {
start := time.Now()
resp, err = handler(ctx, req)
log.Printf(
"method=%s duration=%v err=%v",
info.FullMethod, time.Since(start), err,
)
return resp, err
}
// 认证拦截器
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.Errorf(codes.Unauthenticated, "missing metadata")
}
tokens := md.Get("authorization")
if len(tokens) == 0 {
return nil, status.Errorf(codes.Unauthenticated, "missing token")
}
// 验证token逻辑
if !validateToken(tokens[0]) {
return nil, status.Errorf(codes.Unauthenticated, "invalid token")
}
return handler(ctx, req)
}
func main() {
lis, err := net.Listen("tcp", ":50051")
if err != nil {
log.Fatalf("failed to listen: %v", err)
}
server := grpc.NewServer(
grpc.ChainUnaryInterceptor(
loggingInterceptor,
authInterceptor,
),
)
pb.RegisterUserServiceServer(server, &UserServer{
db: NewUserRepo(),
})
log.Println("gRPC server listening on :50051")
if err := server.Serve(lis); err != nil {
log.Fatalf("failed to serve: %v", err)
}
}
Go客户端调用与连接池管理
gRPC客户端使用grpc.Dial创建连接,支持连接池、重试、负载均衡等能力:
package main
import (
"context"
"log"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/balancer/roundrobin"
pb "github.com/example/proto/user/v1"
)
type UserClient struct {
conn *grpc.ClientConn
client pb.UserServiceClient
}
func NewUserClient(addr string) (*UserClient, error) {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
conn, err := grpc.DialContext(ctx, addr,
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultServiceConfig(`{
"loadBalancingConfig": [{"round_robin": {}}]
}`),
grpc.WithDefaultCallOptions(
grpc.MaxCallRecvMsgSize(4*1024*1024),
),
)
if err != nil {
return nil, err
}
return &UserClient{
conn: conn,
client: pb.NewUserServiceClient(conn),
}, nil
}
func (c *UserClient) GetUser(ctx context.Context, id int64) (*pb.GetUserResponse, error) {
ctx, cancel := context.WithTimeout(ctx, 3*time.Second)
defer cancel()
return c.client.GetUser(ctx, &pb.GetUserRequest{Id: id})
}
func (c *UserClient) Close() error {
return c.conn.Close()
}
流式通信与双向流处理
gRPC支持三种流式通信模式:服务端流、客户端流和双向流。流式接口适合实时数据推送、大文件分块传输、聊天等场景:
// proto扩展:流式接口
service ChatService {
// 服务端流:实时消息推送
rpc Subscribe(SubscribeRequest) returns (stream Message);
// 客户端流:批量上传
rpc UploadBatch(stream UploadChunk) returns (UploadResponse);
// 双向流:实时聊天
rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}
// Go双向流实现
func (s *ChatServer) Chat(
stream pb.ChatService_ChatServer,
) error {
// 为每个连接创建一个通道
userID := generateID()
ch := make(chan *pb.ChatMessage, 100)
// 注册到聊天室
s.chatroom.Register(userID, ch)
defer s.chatroom.Unregister(userID)
// 启动goroutine发送消息
go func() {
for msg := range ch {
stream.Send(msg)
}
}()
// 接收客户端消息
for {
msg, err := stream.Recv()
if err != nil {
break
}
// 广播到聊天室所有成员
s.chatroom.Broadcast(msg)
}
return nil
}
gRPC与REST网关集成
微服务架构中,内部通信使用gRPC,对外API使用REST。grpc-gateway可以自动从proto生成RESTful代理,避免手动维护两套接口定义:
// 在proto中添加HTTP注解
import "google/api/annotations.proto";
service UserService {
rpc GetUser(GetUserRequest) returns (GetUserResponse) {
option (google.api.http) = {
get: "/api/v1/users/{id}"
};
}
rpc CreateUser(CreateUserRequest) returns (CreateUserResponse) {
option (google.api.http) = {
post: "/api/v1/users"
body: "*"
};
}
rpc ListUsers(ListUsersRequest) returns (ListUsersResponse) {
option (google.api.http) = {
get: "/api/v1/users"
};
}
}
// 生成网关代码后启动反向代理
func main() {
ctx := context.Background()
ctx, cancel := context.WithCancel(ctx)
defer cancel()
mux := runtime.NewServeMux()
opts := []grpc.DialOption{grpc.WithTransportCredentials(insecure.NewCredentials())}
pb.RegisterUserServiceHandlerFromEndpoint(ctx, mux, ":50051", opts)
// HTTP请求自动转发到gRPC服务
http.ListenAndServe(":8080", mux)
}
gRPC健康检查与服务发现集成
Kubernetes等编排平台对gRPC服务有原生健康检查支持。使用grpc-health-probe实现标准健康检查接口:
// 注册健康检查服务
import (
healthpb "google.golang.org/grpc/health/proto"
health "google.golang.org/grpc/health"
)
func main() {
server := grpc.NewServer(...)
// 注册业务服务
pb.RegisterUserServiceServer(server, &UserServer{})
// 注册健康检查服务
healthServer := health.NewServer()
healthServer.SetServingStatus("user.v1.UserService", healthpb.HealthCheckResponse_SERVING)
healthpb.RegisterHealthServer(server, healthServer)
server.Serve(lis)
}
// Kubernetes Pod配置
// livenessProbe:
// exec:
// command: ["/bin/grpc_health_probe", "-addr=:50051"]
// readinessProbe:
// exec:
// command: ["/bin/grpc_health_probe", "-addr=:50051", "-service=user.v1.UserService"]
gRPC over HTTP/2天然支持多路复用,一个TCP连接可以并发处理多个请求,避免了HTTP/1.1的连接限制问题。在微服务高并发场景下,gRPC的吞吐量通常比REST JSON高出3-10倍,延迟降低50%以上。配合protobuf的二进制序列化,消息体积比JSON减少30%-60%,显著降低网络带宽消耗。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-grpc-wei-fu-wu-tong-xin-yu-protobuf-jie-kou-ding/