gRPC通信模型与Protocol Buffers接口定义
gRPC基于HTTP/2构建,使用Protocol Buffers作为默认接口定义语言(IDL)和序列化格式。相比JSON+REST,gRPC在微服务间通信中具有更低的序列化开销和更高的吞吐量,单次连接可复用多个请求/响应(HTTP/2多路复用),支持四种通信模式:Unary RPC(一元调用)、Server Streaming(服务端流)、Client Streaming(客户端流)和Bidirectional Streaming(双向流)。
Protocol Buffers通过.proto文件定义服务接口和消息结构,由protoc编译器生成目标语言的桩代码。二进制编码在传输效率上比JSON快3-10倍,且字段编号机制保证了前后向兼容性。
// chat.proto - 双向流式聊天服务定义
syntax = "proto3";
package chat.v1;
option go_package = "github.com/example/chat/api/v1;chatv1";
// 聊天消息
message ChatMessage {
string user_id = 1;
string content = 2;
int64 timestamp = 3;
MessageType type = 4;
}
enum MessageType {
MESSAGE_TYPE_UNSPECIFIED = 0;
MESSAGE_TYPE_TEXT = 1;
MESSAGE_TYPE_IMAGE = 2;
MESSAGE_TYPE_SYSTEM = 3;
}
// 聊房间信息
message RoomInfo {
string room_id = 1;
string name = 2;
repeated string members = 3;
}
// 双向流式聊天服务
service ChatService {
// Unary: 创建聊天室
rpc CreateRoom(CreateRoomRequest) returns (RoomInfo);
// Server Streaming: 获取房间历史消息(服务端推送)
rpc GetHistory(GetHistoryRequest) returns (stream ChatMessage);
// Client Streaming: 客户端批量上传消息
rpc UploadMessages(stream ChatMessage) returns (UploadResponse);
// Bidirectional Streaming: 实时双向聊天
rpc ChatStream(stream ChatMessage) returns (stream ChatMessage);
}
Go语言gRPC服务端双向流实现
使用protoc生成Go桩代码后,实现ChatService接口。双向流的核心是通过goroutine并发处理读写两个方向的流:
package main
import (
"context"
"io"
"log"
"net"
"sync"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
pb "github.com/example/chat/api/v1"
)
type ChatServer struct {
pb.UnimplementedChatServiceServer
mu sync.RWMutex
rooms map[string]*Room
}
type Room struct {
id string
members map[string]pb.ChatService_ChatStreamServer
mu sync.RWMutex
}
func NewChatServer() *ChatServer {
return &ChatServer{
rooms: make(map[string]*Room),
}
}
// ChatStream 双向流式聊天
func (s *ChatServer) ChatStream(stream pb.ChatService_ChatStreamServer) error {
// 第一条消息用于注册用户信息
firstMsg, err := stream.Recv()
if err != nil {
return status.Error(codes.Internal, "failed to receive first message")
}
roomID := firstMsg.UserId // 简化:用user_id关联房间
userID := firstMsg.UserId
room := s.getOrCreateRoom(roomID)
room.mu.Lock()
room.members[userID] = stream
room.mu.Unlock()
// 通知房间内其他成员新用户加入
s.broadcast(room, &pb.ChatMessage{
UserId: "system",
Content: userID + " joined room",
Timestamp: time.Now().Unix(),
Type: pb.MessageType_MESSAGE_TYPE_SYSTEM,
}, userID)
// 接收客户端消息并广播
for {
msg, err := stream.Recv()
if err == io.EOF {
// 客户端断开连接
room.mu.Lock()
delete(room.members, userID)
room.mu.Unlock()
s.broadcast(room, &pb.ChatMessage{
UserId: "system",
Content: userID + " left room",
Timestamp: time.Now().Unix(),
Type: pb.MessageType_MESSAGE_TYPE_SYSTEM,
}, userID)
return nil
}
if err != nil {
log.Printf("recv error from %s: %v", userID, err)
return err
}
msg.Timestamp = time.Now().Unix()
s.broadcast(room, msg, "")
}
}
func (s *ChatServer) broadcast(room *Room, msg *pb.ChatMessage, excludeUser string) {
room.mu.RLock()
defer room.mu.RUnlock()
for uid, s := range room.members {
if uid == excludeUser {
continue
}
if err := s.Send(msg); err != nil {
log.Printf("send to %s failed: %v", uid, err)
}
}
}
gRPC拦截器链与中间件机制
gRPC拦截器(Interceptor)类似HTTP中间件,在请求处理前后插入横切逻辑。Unary拦截器处理一元RPC,Stream拦截器处理流式RPC。常见用途包括认证鉴权、日志记录、指标采集、链路追踪和限流熔断。
// Unary拦截器:认证与日志
func AuthUnaryInterceptor(
ctx context.Context,
req interface{},
info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler,
) (resp interface{}, err error) {
// 从metadata提取token
md, ok := metadata.FromIncomingContext(ctx)
if !ok {
return nil, status.Error(codes.Unauthenticated, "missing metadata")
}
tokens := md.Get("authorization")
if len(tokens) == 0 {
return nil, status.Error(codes.Unauthenticated, "missing auth token")
}
// 验证token(简化)
userID, err := validateToken(tokens[0])
if err != nil {
return nil, status.Error(codes.Unauthenticated, "invalid token")
}
// 注入用户ID到context
ctx = context.WithValue(ctx, userIDKey{}, userID)
start := time.Now()
resp, err = handler(ctx, req)
log.Printf("method=%s user=%s duration=%v err=%v",
info.FullMethod, userID, time.Since(start), err)
return resp, err
}
// Stream拦截器:适用于双向流
func LogStreamInterceptor(
srv interface{},
ss grpc.ServerStream,
info *grpc.StreamServerInfo,
handler grpc.StreamHandler,
) error {
start := time.Now()
log.Printf("stream started: %s client=%v", info.FullMethod, ss.Context())
err := handler(srv, ss)
log.Printf("stream ended: %s duration=%v err=%v",
info.FullMethod, time.Since(start), err)
return err
}
// 组装拦截器链
func main() {
lis, _ := net.Listen("tcp", ":50051")
grpcServer := grpc.NewServer(
grpc.UnaryInterceptor(AuthUnaryInterceptor),
grpc.StreamInterceptor(LogStreamInterceptor),
)
pb.RegisterChatServiceServer(grpcServer, NewChatServer())
log.Println("gRPC server listening on :50051")
grpcServer.Serve(lis)
}
客户端双向流连接与重连机制
gRPC客户端通过grpc.Dial建立连接,Channel默认复用HTTP/2连接。双向流客户端需要处理流的中断和自动重连:
func startChatClient(ctx context.Context, addr string) error {
conn, err := grpc.Dial(addr,
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultCallOptions(
grpc.MaxCallRecvMsgSize(4*1024*1024),
),
)
if err != nil {
return err
}
defer conn.Close()
client := pb.NewChatServiceClient(conn)
// 发起双向流
stream, err := client.ChatStream(ctx)
if err != nil {
return err
}
// 注册用户
stream.Send(&pb.ChatMessage{
UserId: "user-001",
Content: "user-001",
Type: pb.MessageType_MESSAGE_TYPE_TEXT,
})
// goroutine接收服务端推送的消息
go func() {
for {
msg, err := stream.Recv()
if err == io.EOF {
log.Println("stream closed by server")
return
}
if err != nil {
log.Printf("recv error: %v", err)
return
}
log.Printf("[%s] %s: %s",
time.Unix(msg.Timestamp, 0).Format("15:04:05"),
msg.UserId, msg.Content)
}
}()
// 主循环从stdin读取并发送
scanner := bufio.NewScanner(os.Stdin)
for scanner.Scan() {
text := scanner.Text()
if text == "/quit" {
stream.CloseSend()
break
}
stream.Send(&pb.ChatMessage{
UserId: "user-001",
Content: text,
Type: pb.MessageType_MESSAGE_TYPE_TEXT,
})
}
return nil
}
gRPC负载均衡与健康检查
gRPC客户端内置负载均衡,支持round_robin和pick_first策略。通过解析器获取后端地址列表,客户端在HTTP/2连接层轮询发送请求到不同实例:
// 客户端负载均衡配置
conn, err := grpc.Dial(
"dns:///chat-service:50051", // DNS解析器自动获取所有A记录
grpc.WithDefaultServiceConfig(`{
"loadBalancingConfig": {
"round_robin": {}
}
}`),
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
// gRPC健康检查协议(服务端需注册HealthServer)
// health.proto:
// service Health { rpc Check(HealthCheckRequest) returns (HealthCheckResponse); }
// 服务端注册健康检查
health := health.NewServer()
health.SetServingStatus("chat.v1.ChatService", healthpb.HealthCheckResponse_SERVING)
healthpb.RegisterHealthServer(grpcServer, health)
// 客户端配置健康检查(自动排除不健康节点)
grpc.WithDefaultServiceConfig(`{
"loadBalancingConfig": {"round_robin": {}},
"healthCheckConfig": {
"serviceName": "chat.v1.ChatService"
}
}`)
round_robin策略下每个RPC请求轮询到不同后端实例,适合无状态服务。有状态场景使用pick_first(默认),客户端固定到一个后端直到该后端不可用。配合健康检查协议,客户端自动剔除返回NOT_SERVING状态的实例,无需依赖外部负载均衡器即可实现服务发现和容错切换。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/grpcprotocolbuffers-jie-kou-ding-yi-yu-shuang-xiang-liu-shi/