WebSockets 频繁断开连接、SSE 仅支持单向通信、长轮询消耗大量 CPU 资源——这些实时通信的痛点,gRPC 流式传输都能完美解决。目前笔者已通过 gRPC Streaming 实现 10 万并发流处理,单流吞吐量可达 50MB/s,是实时系统的理想选择。

下面就来说说gRpc的流式使用。

核心:四种通信模式及适用场景

gRPC 提供四种优化后的通信模式,精准匹配不同业务需求,掌握其适用场景是构建高效实时系统的关键。

service StreamingService {
    rpc Subscribe(Topic) returns (stream Event);     // 服务端流式
    rpc Upload(stream Chunk) returns (Summary);      // 客户端流式  
    rpc Chat(stream Message) returns (stream Message); // 双向流式
    rpc GetStatus(Empty) returns (Status);           // 简单 RPC(Unary)
}

1. 服务端流式(Server Streaming)

  • 核心能力:服务端向客户端推送持续更新,服务端维护状态并在数据就绪时主动推送。
  • 典型场景:实时仪表盘、价格行情、通知推送、日志追踪。

2. 客户端流式(Client Streaming)

  • 核心能力:客户端向服务端持续发送数据,服务端接收并实时处理。
  • 典型场景:文件上传、遥测数据上报、批量数据导入(比 REST 上传快 10 倍)。

3. 双向流式(Bidirectional Streaming)

  • 核心能力:全双工通信,客户端与服务端可独立发送/接收数据,互不阻塞。
  • 典型场景:聊天系统、实时游戏、协同编辑工具。

4. 简单 RPC(Unary)

  • 核心能力:传统请求-响应模式,简单可靠。
  • 典型场景:标准 API 操作、CRUD 接口,无需流式传输的基础业务。

实战代码:三种流式场景实现

一、服务端流式:向客户端推送实时更新

服务端实现(Go)
type server struct {
    pb.UnimplementedStreamingServiceServer
    subscribers map[string][]chan *pb.Event // 按主题存储订阅者通道
    mu          sync.RWMutex                // 并发安全锁
}

// 订阅主题,持续接收事件
func (s *server) Subscribe(req *pb.Topic, stream pb.StreamingService_SubscribeServer) error {
    // 创建带缓冲通道,避免慢消费者阻塞
    ch := make(chan *pb.Event, 100)
    
    // 注册订阅者
    s.mu.Lock()
    s.subscribers[req.Name] = append(s.subscribers[req.Name], ch)
    s.mu.Unlock()
    
    // 退出时清理资源
    defer func() {
        s.mu.Lock()
        subs := s.subscribers[req.Name]
        for i, sub := range subs {
            if sub == ch {
                s.subscribers[req.Name] = append(subs[:i], subs[i+1:]...)
                break
            }
        }
        s.mu.Unlock()
        close(ch)
    }()
    
    // 持续向客户端发送事件
    for {
        select {
        case event := <-ch:
            if err := stream.Send(event); err != nil {
                return err // 客户端断开连接
            }
        case <-stream.Context().Done():
            return nil // 客户端主动取消
        }
    }
}

// 向指定主题的所有订阅者发布事件
func (s *server) PublishEvent(topic string, event *pb.Event) {
    s.mu.RLock()
    subscribers := s.subscribers[topic]
    s.mu.RUnlock()
    
    for _, ch := range subscribers {
        select {
        case ch <- event:
        default:
            // 通道满时丢弃消息(可优化为背压处理)
            log.Printf("Dropping message for slow consumer")
        }
    }
}
客户端实现(带自动重连)
// 带自动重连的订阅方法
func SubscribeWithReconnect(client pb.StreamingServiceClient, topic string) {
    backoff := time.Second
    maxBackoff := time.Minute
    
    for {
        err := subscribe(client, topic)
        if err != nil {
            log.Printf("Stream error: %v, reconnecting in %v", err, backoff)
            time.Sleep(backoff)
            backoff *= 2 // 指数退避
            if backoff > maxBackoff {
                backoff = maxBackoff
            }
        } else {
            backoff = time.Second // 连接成功重置退避时间
        }
    }
}

// 基础订阅逻辑
func subscribe(client pb.StreamingServiceClient, topic string) error {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()
    
    stream, err := client.Subscribe(ctx, &pb.Topic{Name: topic})
    if err != nil {
        return err
    }
    
    // 持续接收服务端事件
    for {
        event, err := stream.Recv()
        if err == io.EOF {
            return nil // 流正常结束
        }
        if err != nil {
            return err // 连接异常
        }
        
        // 处理事件(业务逻辑自定义)
        handleEvent(event)
    }
}

二、客户端流式:高效大文件上传

服务端实现(Go)
// 处理文件上传流
func (s *server) Upload(stream pb.StreamingService_UploadServer) error {
    var (
        fileID   = uuid.New().String()       // 生成唯一文件ID
        file     *os.File                    // 临时文件
        size     int64                       // 文件总大小
        checksum = sha256.New()              // 校验和计算器
    )
    
    // 创建临时文件
    file, err := os.CreateTemp("", "upload-*.tmp")
    if err != nil {
        return status.Errorf(codes.Internal, "failed to create file: %v", err)
    }
    defer os.Remove(file.Name()) // 异常时清理临时文件
    
    // 持续接收客户端上传的文件分片
    for {
        chunk, err := stream.Recv()
        if err == io.EOF {
            break // 上传完成
        }
        if err != nil {
            return status.Errorf(codes.Internal, "failed to receive chunk: %v", err)
        }
        
        // 写入分片数据
        n, err := file.Write(chunk.Data)
        if err != nil {
            return status.Errorf(codes.Internal, "failed to write: %v", err)
        }
        
        size += int64(n)
        checksum.Write(chunk.Data)
        
        // 限制文件最大1GB
        if size > 1<<30 {
            return status.Errorf(codes.InvalidArgument, "file too large")
        }
    }
    
    // 移动临时文件到最终存储路径
    finalPath := fmt.Sprintf("/uploads/%s", fileID)
    if err := os.Rename(file.Name(), finalPath); err != nil {
        return status.Errorf(codes.Internal, "failed to save: %v", err)
    }
    
    // 返回上传结果摘要
    return stream.SendAndClose(&pb.Summary{
        FileId:   fileID,
        Size:     size,
        Checksum: fmt.Sprintf("%x", checksum.Sum(nil)),
    })
}
客户端实现(带进度跟踪)
// 上传文件并显示进度
func UploadFile(client pb.StreamingServiceClient, filepath string) error {
    file, err := os.Open(filepath)
    if err != nil {
        return err
    }
    defer file.Close()
    
    // 获取文件大小(用于计算进度)
    stat, _ := file.Stat()
    totalSize := stat.Size()
    
    // 设置5分钟超时上下文
    ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
    defer cancel()
    
    // 发起上传流请求
    stream, err := client.Upload(ctx)
    if err != nil {
        return err
    }
    
    // 1MB分片上传
    buf := make([]byte, 1<<20)
    var sent int64
    
    for {
        n, err := file.Read(buf)
        if err == io.EOF {
            break
        }
        if err != nil {
            return err
        }
        
        // 发送分片
        if err := stream.Send(&pb.Chunk{Data: buf[:n]}); err != nil {
            return err
        }
        
        // 更新并显示进度
        sent += int64(n)
        progress := float64(sent) / float64(totalSize) * 100
        fmt.Printf("\rUploading: %.1f%%", progress)
    }
    
    // 接收上传结果
    summary, err := stream.CloseAndRecv()
    if err != nil {
        return err
    }
    
    fmt.Printf("\nUpload complete: %s (%d bytes)\n", summary.FileId, summary.Size)
    return nil
}

三、双向流式:实时聊天系统

服务端实现(Go)
// 聊天服务端结构体
type chatServer struct {
    rooms map[string]*Room // 房间集合
    mu    sync.RWMutex     // 房间操作锁
}

// 聊天房间结构体
type Room struct {
    clients map[string]pb.StreamingService_ChatServer // 房间内客户端
    mu      sync.RWMutex                              // 客户端操作锁
}

// 双向流式聊天实现
func (s *chatServer) Chat(stream pb.StreamingService_ChatServer) error {
    // 第一个消息必须是加入房间请求
    msg, err := stream.Recv()
    if err != nil {
        return err
    }
    
    if msg.Type != pb.MessageType_JOIN {
        return status.Errorf(codes.InvalidArgument, "first message must be JOIN")
    }
    
    roomID := msg.RoomId
    userID := msg.UserId
    
    // 获取或创建房间
    s.mu.Lock()
    room, exists := s.rooms[roomID]
    if !exists {
        room = &Room{clients: make(map[string]pb.StreamingService_ChatServer)}
        s.rooms[roomID] = room
    }
    s.mu.Unlock()
    
    // 将客户端加入房间
    room.mu.Lock()
    room.clients[userID] = stream
    room.mu.Unlock()
    
    // 广播用户加入通知
    s.broadcast(room, &pb.Message{
        Type:      pb.MessageType_USER_JOINED,
        UserId:    userID,
        Timestamp: time.Now().Unix(),
    }, userID)
    
    // 退出时清理客户端并广播离开通知
    defer func() {
        room.mu.Lock()
        delete(room.clients, userID)
        room.mu.Unlock()
        
        s.broadcast(room, &pb.Message{
            Type:      pb.MessageType_USER_LEFT,
            UserId:    userID,
            Timestamp: time.Now().Unix(),
        }, userID)
    }()
    
    // 持续处理客户端消息
    for {
        msg, err := stream.Recv()
        if err == io.EOF {
            return nil
        }
        if err != nil {
            return err
        }
        
        // 丰富消息元数据
        msg.Timestamp = time.Now().Unix()
        msg.UserId = userID
        
        // 广播消息到房间内所有客户端
        s.broadcast(room, msg, "")
    }
}

// 广播消息(支持排除特定用户)
func (s *chatServer) broadcast(room *Room, msg *pb.Message, exclude string) {
    room.mu.RLock()
    defer room.mu.RUnlock()
    
    for userID, client := range room.clients {
        if userID == exclude {
            continue
        }
        
        // 协程发送,避免阻塞
        go func(c pb.StreamingService_ChatServer, m *pb.Message) {
            if err := c.Send(m); err != nil {
                log.Printf("Failed to send to %s: %v", userID, err)
            }
        }(client, msg)
    }
}
客户端实现(Go)
// 启动聊天客户端
func StartChat(client pb.StreamingServiceClient, roomID, userID string) error {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()
    
    // 发起聊天流请求
    stream, err := client.Chat(ctx)
    if err != nil {
        return err
    }
    
    // 发送加入房间消息
    if err := stream.Send(&pb.Message{
        Type:   pb.MessageType_JOIN,
        RoomId: roomID,
        UserId: userID,
    }); err != nil {
        return err
    }
    
    // 启动接收消息协程
    go func() {
        for {
            msg, err := stream.Recv()
            if err != nil {
                log.Printf("Receive error: %v", err)
                cancel()
                return
            }
            displayMessage(msg) // 显示消息(业务逻辑自定义)
        }
    }()
    
    // 从标准输入读取消息并发送
    scanner := bufio.NewScanner(os.Stdin)
    for scanner.Scan() {
        text := scanner.Text()
        if text == "/quit" {
            return nil // 退出聊天
        }
        
        if err := stream.Send(&pb.Message{
            Type:    pb.MessageType_CHAT,
            Content: text,
        }); err != nil {
            return err
        }
    }
    
    return scanner.Err()
}

关键技术:流控、元数据与连接管理

1. 流控与背压处理

gRPC 内置流控机制,但需手动适配业务场景:

// 服务端流控:发送前检查客户端就绪状态
func (s *server) StreamWithFlowControl(req *pb.Request, stream pb.Service_StreamServer) error {
    for i := 0; i < 1000000; i++ {
        // 客户端接收缓冲区满时会阻塞
        if err := stream.Send(&pb.Response{Data: generateData(i)}); err != nil {
            if status.Code(err) == codes.Unavailable {
                log.Printf("Client overwhelmed at message %d", i)
            }
            return err
        }
        
        // 可选:每100条消息限流10ms
        if i%100 == 0 {
            time.Sleep(10 * time.Millisecond)
        }
    }
    return nil
}

// 客户端背压处理:限制并发处理数
func ConsumeWithBackpressure(client pb.ServiceClient) error {
    stream, err := client.Stream(context.Background(), &pb.Request{})
    if err != nil {
        return err
    }
    
    sem := make(chan struct{}, 10) // 最大10个并发处理
    
    for {
        msg, err := stream.Recv()
        if err == io.EOF {
            break
        }
        if err != nil {
            return err
        }
        
        sem <- struct{}{} // 获取并发令牌
        go func(m *pb.Response) {
            defer func() { <-sem }() // 释放令牌
            
            // 模拟慢处理
            time.Sleep(100 * time.Millisecond)
            process(m)
        }(msg)
    }
    
    // 等待所有任务完成
    for i := 0; i < cap(sem); i++ {
        sem <- struct{}{}
    }
    
    return nil
}

2. 元数据与头部传递

支持在流中传递元数据(如认证信息、版本号):

// 服务端发送元数据
func (s *server) AuthenticatedStream(req *pb.Request, stream pb.Service_StreamServer) error {
    // 发送初始元数据
    header := metadata.Pairs(
        "stream-id", uuid.New().String(),
        "server-version", "1.0.0",
    )
    stream.SendHeader(header)
    
    // 发送流数据
    for i := 0; i < 100; i++ {
        stream.Send(&pb.Response{Data: fmt.Sprintf("Message %d", i)})
    }
    
    // 发送尾随元数据
    trailer := metadata.Pairs(
        "message-count", "100",
        "checksum", "abc123",
    )
    stream.SetTrailer(trailer)
    
    return nil
}

// 客户端读取元数据
func ReadWithMetadata(client pb.ServiceClient) error {
    stream, err := client.Stream(context.Background(), &pb.Request{})
    if err != nil {
        return err
    }
    
    // 获取初始元数据
    header, err := stream.Header()
    if err != nil {
        return err
    }
    streamID := header.Get("stream-id")[0]
    log.Printf("Stream ID: %s", streamID)
    
    // 读取流消息
    for {
        _, err := stream.Recv()
        if err == io.EOF {
            break
        }
        if err != nil {
            return err
        }
    }
    
    // 获取尾随元数据
    trailer := stream.Trailer()
    count := trailer.Get("message-count")[0]
    log.Printf("Received %s messages", count)
    
    return nil
}

3. 可靠连接管理(保活+重连)

// 创建高可用客户端(带保活和重连)
func NewRobustClient(addr string) (pb.ServiceClient, error) {
    conn, err := grpc.Dial(addr,
        grpc.WithTransportCredentials(insecure.NewCredentials()),
        // 保活配置:每10秒发ping,3秒超时,无流时也发送
        grpc.WithKeepaliveParams(keepalive.ClientParameters{
            Time:                10 * time.Second,
            Timeout:             3 * time.Second,
            PermitWithoutStream: true,
        }),
        // 消息大小限制:50MB
        grpc.WithDefaultCallOptions(
            grpc.MaxCallRecvMsgSize(50*1024*1024),
            grpc.MaxCallSendMsgSize(50*1024*1024),
        ),
        // 流重连拦截器
        grpc.WithStreamInterceptor(streamRetryInterceptor()),
    )
    if err != nil {
        return nil, err
    }
    
    return pb.NewServiceClient(conn), nil
}

// 流重连拦截器
func streamRetryInterceptor() grpc.StreamClientInterceptor {
    return func(ctx context.Context, desc *grpc.StreamDesc, cc *grpc.ClientConn,
        method string, streamer grpc.Streamer, opts ...grpc.CallOption) (grpc.ClientStream, error) {
        
        var stream grpc.ClientStream
        var err error
        
        // 最多重试3次
        for i := 0; i < 3; i++ {
            stream, err = streamer(ctx, desc, cc, method, opts...)
            if err == nil {
                return stream, nil
            }
            
            // 仅重试可用状态和资源耗尽错误
            code := status.Code(err)
            if code == codes.Unavailable || code == codes.ResourceExhausted {
                time.Sleep(time.Duration(i+1) * time.Second)
                continue
            }
            break // 非重试错误直接返回
        }
        
        return nil, err
    }
}

4. 流负载均衡

// 客户端负载均衡配置(轮询+健康检查)
conn, err := grpc.Dial("dns:///myservice.local:50051",
    grpc.WithDefaultServiceConfig(`{
        "loadBalancingPolicy": "round_robin",
        "healthCheckConfig": {
            "serviceName": "StreamingService"
        }
    }`),
    grpc.WithTransportCredentials(insecure.NewCredentials()),
)

// 服务端注册到Consul服务发现
func RegisterWithConsul(consulAddr, serviceID, grpcAddr string) error {
    client, err := consul.NewClient(&consul.Config{
        Address: consulAddr,
    })
    if err != nil {
        return err
    }
    
    return client.Agent().ServiceRegister(&consul.AgentServiceRegistration{
        ID:      serviceID,
        Name:    "streaming-service",
        Port:    50051,
        Address: grpcAddr,
        Check: &consul.AgentServiceCheck{
            GRPC:                           grpcAddr,
            Interval:                       "10s", // 每10秒健康检查
            DeregisterCriticalServiceAfter: "1m",  // 异常1分钟后注销
        },
    })
}

性能基准测试

通信方式吞吐量延迟(p99)CPU使用率
REST (HTTP/1.1)10K req/s50ms80%
WebSocket50K msg/s10ms60%
gRPC Unary80K req/s5ms40%
gRPC Streaming500K msg/s1ms35%

常见陷阱与避坑指南

  1. 未处理重连:流连接必然会断开,必须实现自动重连机制。
  2. 忽略流控:生产者速度远快于消费者时,会导致内存溢出(OOM)。
  3. 未设置超时:流可能无限运行,需通过上下文设置合理 deadlines。
  4. 处理器阻塞:耗时操作需用协程异步处理,避免阻塞流通道。
  5. 资源未清理:流处理器中必须用 defer 清理资源(如连接、通道)。

模式选择建议

  • 服务端流式:实时馈送、日志传输、监控数据
  • 客户端流式:文件上传、批量处理、遥测数据
  • 双向流式:聊天、游戏、协同编辑
  • 简单 RPC:简单请求响应、CRUD 操作

实用技巧:优先使用简单 RPC,仅在流式传输能带来明显性能或体验提升时,再引入流式模式——避免为不必要的复杂度买单。

Logo

openvela 操作系统专为 AIoT 领域量身定制,以轻量化、标准兼容、安全性和高度可扩展性为核心特点。openvela 以其卓越的技术优势,已成为众多物联网设备和 AI 硬件的技术首选,涵盖了智能手表、运动手环、智能音箱、耳机、智能家居设备以及机器人等多个领域。

更多推荐