一、从一个半夜报警说起
凌晨两点,手机连着震了七八下。群里值班同事发来一张监控图:某个微服务的内存占用像爬坡一样,从 1G 慢慢涨到 4G,然后被 OOM Killer 干掉,重启,再涨,再被杀。大家第一反应是缓存没设过期,第二反应是数据库连接没关,查了一圈都没问题。最后我把堆栈导出来一看,发现 goroutine 的数量异常庞大,有几十万个,全部卡在同一个地方:等待从流里读数据。
那一刻我意识到,这不是普通的内存泄漏,而是 gRPC 流式 RPC 用错了姿势导致的 goroutine 泄漏。goroutine 虽然很轻,每个栈初始只有几 KB,但几十万个堆在一起,内存照样扛不住,而且它们引用的对象、buffer、连接资源也会跟着一起驻留,内存自然就爆了。
今天咱们就掰开揉碎聊聊,gRPC 流式调用里,那个“未正确关闭流”和“done channel 处理不当”到底是怎么把 goroutine 困死的,以及我们该怎么排查、怎么修。
二、先搞懂 gRPC 的三种流式模式
gRPC 有四种调用方式:一元调用(就是普通请求-响应)、服务端流式、客户端流式、双向流式。流式调用意味着连接上会持续有数据流动,而 gRPC 底层是用 HTTP/2 的多路复用实现的,每个流都对应一个独立的 stream。
这里重点说双向流式,因为最容易出问题。服务端和客户端各持有一个流对象,服务端可以不断往流里写消息,客户端也可以不断往流里写消息,读写是独立的。比如一个实时聊天服务,客户端发一条消息,服务端回一条,这就是典型的双向流。
要使用流式 RPC,我们得先定义 proto 文件,比如:
// 这是 proto 文件,定义聊天服务的消息和接口
syntax = "proto3";
package chat;
// 定义聊天消息结构
message ChatMessage {
string user = 1; // 用户名
string content = 2; // 消息内容
}
// 定义聊天服务
service ChatService {
// 双向流式 RPC:客户端上传消息流,服务端返回响应流
rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}
然后通过 protoc 生成 Go 代码。生成后的客户端长这样:
// 客户端调用方式
stream, err := client.Chat(ctx) // 拿到一个双向流
if err != nil {
log.Fatalf("创建流失败: %v", err)
}
// 发送消息
stream.Send(&ChatMessage{User: "张三", Content: "你好"})
// 接收响应
resp, err := stream.Recv()
看起来很简单对吧?但坑就在这个 stream 上。如果流没有正确关闭,或者提前退出了,goroutine 就会一直阻塞在 Recv() 或 Send() 上,永远等不到数据,也等不到错误。
三、真凶一号:忘记关闭流
很多人在用流式 RPC 的时候,会这样写客户端代码:
// 这段代码有严重问题,会导致 goroutine 泄漏
func chatWithServer(ctx context.Context, client chat.ChatServiceClient) {
// 创建双向流
stream, err := client.Chat(ctx)
if err != nil {
log.Printf("创建流失败: %v", err)
return
}
// 启动一个 goroutine 专门接收服务端消息
done := make(chan struct{})
go func() {
for {
resp, err := stream.Recv()
if err != nil {
// 一旦出错就退出循环
log.Printf("接收消息失败: %v", err)
close(done) // 通知主流程结束
return
}
log.Printf("收到回复: %s", resp.Content)
}
}()
// 主流程发送一批消息
for i := 0; i < 10; i++ {
if err := stream.Send(&ChatMessage{User: "张三", Content: fmt.Sprintf("第%d条", i)}); err != nil {
log.Printf("发送失败: %v", err)
break
}
}
// 注意:这里没有关闭流!没有调用 CloseSend()!
// 也没有等待 done channel,函数直接返回了
select {
case <-done:
case <-time.After(5 * time.Second):
log.Println("等待超时,直接退出")
}
// 函数返回后,stream 对象失去引用,但底层 goroutine 还在阻塞
}
这段代码漏洞百出。首先,发送完消息后,我们没有调用 CloseSend() 来告诉服务端“我说完了”。服务端可能还在等后续消息,导致服务端那边的 goroutine 也一直挂着。更要命的是,接收消息那个 goroutine 一直在 Recv() 里阻塞。就算我们用了 select 等待 done,超时后函数直接返回了,但那个 goroutine 并没有被杀死,它还在那个 stream 上等着,底层的 HTTP/2 连接也被它占着。
为什么 Recv() 会一直阻塞?因为 gRPC 流的 Recv() 在没有消息时是阻塞的,只有三种情况会返回:收到消息、收到错误、流被关闭。如果服务端一直不发送消息,也不关闭流,客户端这边就永远卡住。
正确的做法是:当我们不需要再发送数据时,主动调用 CloseSend()。这会向服务端发送一个半关闭信号,告诉对方“我的发送结束了”,但接收还能继续。同时,我们还需要确保接收 goroutine 能够正常退出,比如服务端收到 CloseSend() 后,也会关闭它的响应流,这样客户端的 Recv() 就会返回 io.EOF,从而退出循环。
改进后的代码:
// 修复后的客户端:正确关闭发送流,并等待接收协程退出
func chatWithServer(ctx context.Context, client chat.ChatServiceClient) {
// 创建双向流
stream, err := client.Chat(ctx)
if err != nil {
log.Printf("创建流失败: %v", err)
return
}
// 保证在函数退出前关闭发送方向
defer stream.CloseSend()
done := make(chan struct{})
// 接收协程:专门处理服务端推送
go func() {
defer close(done) // 无论什么路径退出,都通知主流程
for {
resp, err := stream.Recv()
if err != nil {
if err == io.EOF {
// 服务端正常关闭流,这是预期内的结束
log.Println("服务端关闭了流")
} else {
// 其他错误,比如网络断、context 取消
log.Printf("接收消息错误: %v", err)
}
return
}
log.Printf("收到回复: %s", resp.Content)
}
}()
// 发送消息
for i := 0; i < 10; i++ {
if err := stream.Send(&ChatMessage{User: "张三", Content: fmt.Sprintf("第%d条", i)}); err != nil {
log.Printf("发送失败: %v", err)
return
}
}
// 发送完毕,主动关闭发送流
// 注意:这里不能提前 defer 掉,因为我们需要先发送完再关闭
// 如果在上面的循环里遇到错误,defer 也会触发 CloseSend()
// 所以这里手动调用一次也可以,但 defer 会重复调用,CloseSend 重复调用是无害的
stream.CloseSend()
// 等待接收协程退出,并设置超时保护
select {
case <-done:
log.Println("接收协程已正常退出")
case <-time.After(10 * time.Second):
// 超时说明服务端迟迟不关流,可能是业务异常
// 此时我们可以取消 context 来强制关闭底层连接
log.Println("等待接收协程超时,可能流未关闭")
// 使用 cancel 上下文可以强制释放
// 但这里上下文是外部传入的,更好的做法是调用 CloseSend 或取消 context
}
}
看,这样就清晰多了。但仍有隐患:如果服务端一直不关流,我们即使调用了 CloseSend(),Recv() 依然会阻塞。这时候需要一个超时控制,或者用 context 取消。
四、真凶二号:done channel 乱用
所谓的 done channel,一般指用于通知协程结束的 channel。很多开发者会自己造一个 channel 来管理协程生命周期,但经常用错。最常见的错误是:往 channel 发信号时没人接收,或者关闭 channel 时重复 close 导致 panic,或者根本没人 close channel 导致接收方永远阻塞。
比如这段代码:
// 错误示范:done channel 被错误地关闭
func badDemo(ctx context.Context, client chat.ChatServiceClient) {
stream, err := client.Chat(ctx)
if err != nil {
return
}
done := make(chan bool)
// 协程 A:负责接收
go func() {
for {
_, err := stream.Recv()
if err != nil {
// 这里尝试关闭 done,但可能已经被协程 B 关闭了
close(done) // 重复关闭会引起 panic
return
}
}
}()
// 协程 B:负责监控超时
go func() {
time.Sleep(5 * time.Second)
close(done) // 也是关闭 done
}()
<-done // 主流程在这里等待
}
两个协程都可能在 close(done),一旦其中一个先执行了 close,另一个再 close 就是 close of closed channel panic,整个程序直接崩溃。另外,如果协程 A 因为某些原因一直不退出,协程 B 关闭了 done,主流程走开了,但协程 A 依然卡在 Recv() 里,泄漏依旧。
正确做法是:使用 sync.Once 来保证只关闭一次,或者干脆用 context 来取消。gRPC 本身支持 context 取消,当我们调用 cancel() 时,底层的流会被中断,Recv() 会返回错误,协程就能退出。所以与其自己造 done channel,不如直接用 context:
// 推荐做法:使用 context 控制协程生命周期
func goodDemo(ctx context.Context, client chat.ChatServiceClient) {
// 创建可取消的上下文
ctx, cancel := context.WithCancel(ctx)
defer cancel() // 确保退出时取消
stream, err := client.Chat(ctx)
if err != nil {
log.Printf("创建流失败: %v", err)
return
}
// 接收协程
go func() {
for {
resp, err := stream.Recv()
if err != nil {
// context 取消会导致这里得到错误,从而退出
log.Printf("接收结束: %v", err)
return
}
log.Printf("收到: %s", resp.Content)
}
}()
// 发送完数据
for i := 0; i < 5; i++ {
if err := stream.Send(&ChatMessage{User: "李四", Content: "hello"}); err != nil {
log.Printf("发送错误: %v", err)
return
}
}
// 发送完毕,关闭发送流
stream.CloseSend()
// 这里不需要自己造 done channel,用 context 控制即可
// 如果业务逻辑处理完了,调用 cancel(),所有协程都会收到取消信号
}
原理是:gRPC 的流在创建时绑定了 context,一旦 context 被取消,底层连接会关闭,所有阻塞的 Recv() 和 Send() 都会立即返回 error,协程自然就退出了。这样我们就不需要小心翼翼地管理 done channel 的关闭了。
但是,在某些场景下我们仍然需要主动等待协程退出,比如要确认接收协程把数据都处理完了。这时候可以使用 sync.WaitGroup 或用一个专门的 channel 来通知,但一定要保证不会重复关闭或永远等待。
五、服务端同样有雷
别以为只有客户端会泄漏,服务端写不好一样炸。服务端流式处理时,如果客户端意外断开,服务端的循环里可能还在往流里写消息,或者还在等待读取客户端消息,导致 goroutine 卡死。
看一个典型的服务端业务方法:
// 服务端实现双向流 RPC
func (s *chatServer) Chat(stream chat.ChatService_ChatServer) error {
// 这个 goroutine 负责接收客户端消息
go func() {
for {
msg, err := stream.Recv()
if err != nil {
// 客户端断开了,这里收到错误
// 但如果外层函数已经返回,这个 goroutine 还怎么退出?
log.Printf("接收客户端消息失败: %v", err)
return
}
log.Printf("收到客户端消息: %s", msg.Content)
}
}()
// 这里循环发送消息给客户端
for i := 0; i < 100; i++ {
if err := stream.Send(&ChatMessage{User: "服务器", Content: "响应"}); err != nil {
return err
}
}
// 函数返回了,但上面那个接收 goroutine 还在跑吗?
// 这个问题很严重
return nil
}
服务端的方法返回后,gRPC 框架会认为这个流已经结束了,底层的流对象会被清理。但是我们在方法内部自己启动的 goroutine 并不受框架管理,如果它还在 stream.Recv() 上阻塞,它依然持有这个 stream 对象的引用,导致资源无法释放。更糟糕的是,如果这个 goroutine 永远不退出,它就是泄漏的。
为什么它会卡住?因为客户端可能没有关闭发送流,它会一直等客户端发消息。而我们已经返回了,不会再有人去读这个流了。
正确的服务端做法是:不要在方法内部随意启动 goroutine 去长期阻塞。如果必须用,一定要确保它能在方法返回之前退出,或者用 context 控制。其实 gRPC 服务端的 stream 对象自带 context(通过 stream.Context() 获取),我们可以监听这个 context 的取消来退出 goroutine:
// 正确的服务端实现:使用 stream.Context() 来控制接收协程退出
func (s *chatServer) Chat(stream chat.ChatService_ChatServer) error {
ctx := stream.Context() // 获取流的上下文
done := make(chan struct{})
// 接收协程:处理客户端消息
go func() {
defer close(done)
for {
msg, err := stream.Recv()
if err != nil {
return
}
// 处理消息...
log.Printf("收到客户端消息: %s", msg.Content)
}
}()
// 等待接收协程退出,或者上下文取消
select {
case <-done:
// 接收协程退出,可能是客户端关闭了发送流或出错
case <-ctx.Done():
// 客户端断开连接,context 被取消
log.Println("客户端断开连接")
}
// 此时函数可以安全返回,因为协程已经退出
return nil
}
这样写,接收协程的生命周期就和流的生命周期绑定在一起了。客户端断开、方法返回、context 取消,这些事件都能触发协程退出。
六、排查泄漏的实战套路
遇到 goroutine 泄漏,别急着瞎猜。Go 有现成的 pprof 工具,可以快速定位问题。
6.1 先看 goroutine 数量
在代码里引入 net/http/pprof,然后在运行中访问 /debug/pprof/goroutine?debug=1 就能看到所有 goroutine 的栈信息。也可以直接在代码里打印:
// 在关键位置打印 goroutine 数量
import "runtime"
fmt.Println("当前 goroutine 数量:", runtime.NumGoroutine())
如果发现数量持续增长,那基本就是泄漏了。
6.2 抓取 goroutine 栈
用 pprof 拿到阻塞的 goroutine 栈,栈里会显示它们卡在哪个函数上。比如:
goroutine 12345 [chan receive]:
google.golang.org/grpc/internal/transport.(*http2Client).recvMsg(0xc0000b2000, 0xc0000f0000)
.../grpc/internal/transport/http2_client.go:...
看到 chan receive 和 recvMsg,就说明它正在阻塞等待接收消息。再往上找,能看到是哪个业务函数调用 stream.Recv() 的。
6.3 检查 stream 是否关闭
一个简单的方法:在 goroutine 栈里找到 (*controlBuffer).get 或者 Recv 阻塞的调用,大概率是流没关。也可以打印流的 Context().Err(),如果为 nil 说明流还活着,如果返回 context.Canceled 或 DeadlineExceeded 则说明被取消了。
6.4 用单位测试复现
我们还可以写一个小测试,通过不断创建流但不关闭,来观察 goroutine 数量变化:
// 测试代码:验证泄漏的复现
func TestLeak(t *testing.T) {
// 先记录初始 goroutine 数量
base := runtime.NumGoroutine()
for i := 0; i < 1000; i++ {
// 每次创建流,但不做任何清理
stream, _ := client.Chat(context.Background())
_ = stream // 故意不关闭
}
time.Sleep(time.Second)
after := runtime.NumGoroutine()
t.Logf("初始 goroutine: %d, 之后: %d", base, after)
}
如果 after 远大于 base,说明确实泄漏了。然后我们尝试修复代码,再跑测试,看数量是否回归正常。
七、这类问题的应用场景
流式 RPC 在实时通信、数据推送、大文件传输等场景下特别常见。比如:
- 聊天系统:客户端和服务端互发消息,双向流最合适。
- 实时监控:服务端每隔几秒推送一次 CPU 或内存数据,客户端持续接收。
- 事件驱动:比如订单状态变更,服务端将事件流推给客户端。
这些场景都需要长时间维持流的连接,也正因为时间长,才更容易发生泄漏。一次两次没事,但连接一多,内存就吃紧了。
八、流式 RPC 的技术优缺点
任何技术都有两面性。流式 RPC 的优点很明显:
- 节省连接:多个消息复用同一个 HTTP/2 连接,减少握手开销。
- 实时性好:服务端可以在任意时刻推送数据,不需要客户端轮询。
- 内存友好:消息可以逐个处理,不需要一次性把所有数据加载到内存。
缺点也很明显:
- 生命周期管理复杂:流的打开、关闭、异常处理都比一元 RPC 复杂,稍有不慎就会泄漏。
- 排查难度大:泄漏发生在底层连接和 goroutine 层面,肉眼很难发现。
- 服务端压力大:长时间占用的连接也会占用服务端资源,如果客户端乱来,服务端也可能被拖垮。
九、注意事项总结
结合上面案例,我总结几个保命经验:
- 凡是创建了流,必须确保有一条路径能够关闭它。不管是
defer stream.CloseSend()还是cancel(),一定要有。 - 不要在函数里随手启动一个 goroutine 去
Recv(),你必须有办法让它退出。优先使用 context 而不是手写的 done channel。 - 如果必须要用 done channel,请用
sync.Once来防止重复关闭,并且要做超时处理。 - 服务端要监听
stream.Context().Done(),随时感知客户端断开。 - 线上服务一定要开 pprof,并且监控 goroutine 数量,设定告警阈值。
- 每一次
Recv()返回错误时,要区分io.EOF和真实错误。io.EOF是流正常关闭,不是错误。 - context.WithCancel 的 cancel 函数要及时调用,否则 context 本身也会泄漏(它是上下文取消的订阅者)。
- 别忘了一个容易被忽略的细节:
CloseSend()只能关闭发送方向,接收方向还需要依赖对方关闭流或我们自己取消 context。
十、总结:别让“小流”变成“大灾”
这次排查让我深刻体会到,gRPC 流式调用不是简单的 API 调用,它的底层有状态、有连接、有协程。任何一个没关闭的流,都可能变成一只藏在内存里的吞噬兽。平时开发时,我们只顾着业务逻辑是否正确,却忘了这些“隐形的资源”。
修复这个问题并不难,难的是在写代码的第一时间就意识到:流是有生命周期的,协程是需要被管理的。如果你正被类似的 goroutine 泄漏折磨,不妨顺着这个思路去检查:每个流有没有关闭?每个协程有没有退出路径?每个 done channel 是否安全?排查完一圈,大概率能找到元凶。
技术这条路,坑坑洼洼,但每踩一个坑,把它填平,后面的路就好走了。希望这篇文章能帮你在排查内存泄漏时,少走一些弯路。
Comments