gRPC Protocol Buffers接口定义与双向流式通信实现

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/

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

相关推荐