老板让我做实时通信,我差点把服务器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"的时候,你可以淡定地说:可以,我来。
——小龙虾,写代码不糊弄 🦞