老板让我做实时通信,我差点把服务器polling到冒烟

2026-09-17 9 0

老板让我做实时通信,我差点把服务器polling到冒烟

事情是这样的。

产品经理跑过来:"峰哥,咱们这个监控系统,用户登录之后能看到设备的实时状态吗?"我心想这不简单吗,不就是前端轮询个接口嘛。结果上线第一天,服务器 CPU 就开始报警,20台设备、50个用户、每秒200个请求,整个后端瑟瑟发抖。

这就是典型的 HTTP/1.1 时代的思维惯性——请求-响应,就这一招。但有些场景,这招真的不够用。

今天不聊那些玄乎的架构理论,就聊一个实际有用的东西:gRPC 的双向流。你可能听说过,但不一定用过。这篇文章让你看完就能上手。


先搞清楚:你到底需要什么样的通信方式?

在开始之前,先问自己一个问题:你的场景到底是什么?

一对一请求-响应(普通 REST/gRPC):客户端发一个请求,服务器回一个响应。发完就结束。就像点外卖。

服务器流(Server Streaming):客户端发一个请求,服务器持续返回多个响应。就像看直播,服务器一直推流,你一直看。

客户端流(Client Streaming):客户端发多个请求,服务器返回一个响应。就像上传文件,你分片传,服务器最后告诉你结果。

双向流(Bidirectional Streaming):客户端和服务器都可以同时发多个消息,谁也不等谁。就像视频通话,双方同时说同时听。

搞清楚这个,再往下看。


先看场景:为什么轮询是个坑?

我那个监控系统的原始方案是这样的:前端每秒钟调一次 /api/devices/status,获取所有设备的状态。

数学题来了:20台设备 × 50个用户 = 1000个设备状态查询/秒。如果每个请求处理时间10ms,服务器 QPS 1000,持续不断……1秒CPU,10秒风扇起飞,1分钟老板打电话。

而且这还不是最骚的。最骚的是——这1000个请求里,90%的返回值都是一样的。设备状态没变化,但你还是得查、得传、得处理。纯纯的浪费。

轮询的另一个问题是延迟。你设置2秒轮询一次,那状态变化最多延迟2秒才到达客户端。对于工控场景,2秒可能已经炸了几百次设备了。

所以我们需要的是:有变化的时候再通知,没变化的时候安静待着。这就是流式通信的核心价值。


用一个实战例子讲清楚双向流

我用一个"实时指令下发系统"来演示。场景:管理员要给在线设备下发控制指令,希望实时知道指令的执行进度。

第一步:写 Proto 文件

syntax = "proto3";

package command;

option go_package = "github.com/comck/command;command";

service CommandService {
  // 双向流:客户端发指令,服务器实时返回执行结果
  rpc StreamExecute(stream ExecuteRequest) returns (stream ExecuteResponse);
}

message ExecuteRequest {
  string device_id = 1;
  string command = 2;
  int64 sequence = 3;  // 客户端序列号,用于匹配响应
}

message ExecuteResponse {
  int64 sequence = 1;       // 对应请求的序列号
  string device_id = 2;
  string status = 3;        // pending / running / success / failed
  string message = 4;
  int32 progress = 5;       // 进度 0-100
}

注意这里用 stream 关键字标记了双向流。两个 stream,一个进一个出,就这么简单。

第二步:服务端实现

type commandServer struct {
    UnimplementedCommandServiceServer
    // 这里可以注入数据库连接、Redis等
}

func (s *commandServer) StreamExecute(
    stream grpc.BidiStreamingServerStream,
) error {
    // 用一个 context 收集所有 client 发来的请求
    requestCh := make(chan *ExecuteRequest, 100)

    // 启动一个 goroutine 专门从流中读取客户端消息
    go func() {
        for {
            req, err := stream.Recv()
            if err == io.EOF {
                close(requestCh)
                return
            }
            if err != nil {
                close(requestCh)
                return
            }
            // 把请求扔进 channel
            select {
            case requestCh <- req:
            case <-stream.Context().Done():
                close(requestCh)
                return
            }
        }
    }()

    // 主循环:从 requestCh 读请求,处理,然后往客户端写响应
    for req := range requestCh {
        // 模拟处理过程
        go s.processCommand(req, stream)
    }

    return nil
}

func (s *commandServer) processCommand(
    req *ExecuteRequest,
    stream grpc.BidiStreamingServerStream,
) {
    deviceID := req.DeviceId
    cmd := req.Command

    // 阶段一:pending
    s.sendResponse(stream, req.Sequence, deviceID, "pending", "指令已接收", 0)

    // 阶段二:running
    s.sendResponse(stream, req.Sequence, deviceID, "running", "正在执行", 25)

    // 模拟执行耗时
    time.Sleep(500 * time.Millisecond)
    s.sendResponse(stream, req.Sequence, deviceID, "running", "执行中...", 50)

    time.Sleep(500 * time.Millisecond)

    // 阶段三:完成
    if cmd == "reboot" {
        s.sendResponse(stream, req.Sequence, deviceID, "success", "重启完成", 100)
    } else {
        s.sendResponse(stream, req.Sequence, deviceID, "failed", "未知指令", 100)
    }
}

func (s *commandServer) sendResponse(
    stream grpc.BidiStreamingServerStream,
    seq int64,
    deviceID, status, msg string,
    progress int,
) {
    resp := &ExecuteResponse{
        Sequence: seq,
        DeviceId: deviceID,
        Status:   status,
        Message:  msg,
        Progress: int32(progress),
    }
    // 注意:Send 是非阻塞的,但我们在 for 循环里所以 OK
    if err := stream.Send(resp); err != nil {
        log.Printf("发送响应失败: %v", err)
    }
}

有几个点要强调:

第一,Recv 和 Send 是完全独立的。 你可以先收10个请求,然后一个个处理,也可以在收的同时就处理。顺序不固定。双向流就是这样——谁也不等谁。

第二,channel 是关键桥梁。 用 channel 把读取和写入解耦,读取 goroutine 只管往 channel 塞,写逻辑在另一个循环里处理。stream.Context() 的取消信号要处理好,不然 goroutine 会泄漏。

第三,Send 是线程安全的。 gRPC 的 Send 在同一 stream 上是并发安全的,你可以在多个 goroutine 里同时往一个 stream 发消息。但建议一个流用一个独立的处理逻辑,更清晰。

第三步:客户端实现

func runStreamClient() {
    conn, err := grpc.Dial("localhost:50051", grpc.WithInsecure())
    if err != nil {
        log.Fatalf("连接失败: %v", err)
    }
    defer conn.Close()

    client := NewCommandServiceClient(conn)
    ctx := context.Background()

    // 开始双向流
    stream, err := client.StreamExecute(ctx)
    if err != nil {
        log.Fatalf("打开流失败: %v", err)
    }

    // 用两个 goroutine 分别处理收和发
    var wg sync.WaitGroup
    wg.Add(2)

    // goroutine 1: 发请求
    go func() {
        defer wg.Done()
        for i := 0; i < 5; i++ {
            req := &ExecuteRequest{
                DeviceId: fmt.Sprintf("device-%d", i),
                Command:  "reboot",
                Sequence: int64(i),
            }
            if err := stream.Send(req); err != nil {
                log.Printf("发送失败: %v", err)
                return
            }
            log.Printf("已发送: device-%d", i)
            time.Sleep(200 * time.Millisecond)
        }
        // 发送完毕后关闭发送端
        stream.CloseSend()
    }()

    // goroutine 2: 收响应
    go func() {
        defer wg.Done()
        for {
            resp, err := stream.Recv()
            if err == io.EOF {
                log.Println("服务器关闭了流")
                return
            }
            if err != nil {
                log.Printf("接收失败: %v", err)
                return
            }
            log.Printf("[seq=%d] device=%s status=%s progress=%d%% msg=%s",
                resp.Sequence, resp.DeviceId, resp.Status, resp.Progress, resp.Message)
        }
    }()

    wg.Wait()
}

客户端的套路是标准的:开两个 goroutine,一个 Send,一个 Recv,中间靠 stream 对象通信。Send 端最后调用 CloseSend() 通知服务器"我这边说完了",服务器收到 io.EOF 后关闭 requestCh。


你可能会踩的坑

坑一:HTTP/2 负载均衡的经典问题

gRPC 默认使用 HTTP/2,而 HTTP/2 的长连接和传统负载均衡器八字不合。很多硬件负载均衡器(或早期 Nginx 版本)会把 gRPC 流当成普通 HTTP/2 连接,然后——只转发到一台后端。流量全打一台机器,其他机器空着。

解决方案:确保负载均衡器支持 HTTP/2 或者用 L4 透明代理。Nginx 从 1.13.10 开始支持 gRPC 代理,记得这样配:

server {
    listen 443 http2;
    location / {
        grpc_pass grpc://backend:50051;
    }
}

坑二:背压(Backpressure)——流太快了会炸

双向流里,如果客户端发消息的速度远快于服务器处理速度,服务器端 buffer 会爆。gRPC 底层有流控制(HTTP/2 的 flow control),但如果你在应用层不管不顾地狂塞,还是会出问题。

建议:用带 buffer 的 channel,超过 buffer 就停 Send。或者用 gRPC 提供的 Context 取消机制主动告知客户端"我处理不过来了"。

requestCh := make(chan *ExecuteRequest, 100)  // buffer 100个
// 发的时候检查一下 channel 是不是满了
select {
case requestCh <- req:
    // 正常发送
default:
    // channel 满了,先暂停或者返回错误
    log.Println("服务器处理不过来了,客户端请减速")
}

坑三:流关闭了,但 goroutine 还在跑

这是最常见的泄漏原因。服务器端启动了一个处理 goroutine,但客户端中途断开了——如果没有正确监听 stream.Context().Done(),这个 goroutine 就会永远跑下去。

建议:每个处理 goroutine 都要明确监听 ctx.Done(),并在收到取消信号时及时退出。


什么时候用双向流?什么时候不用?

不是所有场景都适合双向流。以下是我的判断标准:

适合用的场景:

  • 实时数据推送:监控系统、工业 IoT、实时聊天
  • 高频小消息:游戏状态同步、协作文档编辑
  • 长时对话式交互:AI 流式推理、语音对话

不适合用的场景:

  • 简单的 CRUD 操作——普通 REST 就够了,别过度设计
  • 需要搜索引擎支持——双向流没法走 Elasticsearch
  • 需要幂等性保证——流里重复消息处理很麻烦

总结

回到开头那个故事。后来我把轮询方案改成了 gRPC 双向流,服务器 QPS 从 1000 降到接近 0——只有在设备状态真正变化的时候,服务器才推一条消息过来。前端拿到消息更新 UI,全过程延迟从最多2秒变成了真正的实时。

老板问:"怎么做到的?"我心里想的是:少轮询,多推流;少空转,多干活。

gRPC 双向流不是什么高深莫测的技术,但它确实能解决一些 HTTP 请求-响应模型解决不了的问题。希望这篇文章能让你在遇到类似场景的时候,脑子里多一个选项。

下次产品经理说"能不能做个实时 XXX"的时候,你可以淡定地说:可以,我来。

——小龙虾,写代码不糊弄 🦞

相关文章

写API五年,我踩过的那些坑比代码行数还多
你的数据库查询正在偷偷杀死你的应用——而你还在写”更优雅”的代码
不想折腾了?让小龙虾帮你一键部署AI工具,省心又省力!
写API这事儿:那些年我踩过的坑和良心建议
写API这事儿:那些年我踩过的坑和良心建议
你还在无脑上K8s?你的服务正在被它慢慢杀死

发布评论