Go与gRPC实现pub-sub时接收端关闭channel的问题如何解决?
Go gRPC Pub/Sub模式订阅端触发通道关闭的解决方案
核心改造思路
遵循Go channel发送端关闭的原则,通过通道包装+生命周期信号+并发安全的通道管理实现订阅端退出时的信号同步,避免发布端写入panic。
具体实现步骤
1. 改造通道存储结构,增加生命周期管理
用自定义结构体包装业务通道,增加退出信号、单次关闭保证,同时给全局通道列表加读写锁保证并发安全:
import "sync" type clientChan struct { ch chan *SendMessageRequest // 业务消息通道 done chan struct{} // 订阅端退出信号 closeOnce sync.Once // 保证资源仅关闭一次 } type GRPCServer struct { // 其他原有字段省略 clientChannels map[*clientChan]struct{} // 用map存储订阅端通道,方便快速删除 chMu sync.RWMutex // 保护clientChannels的并发读写 } // NewClientChannel 初始化订阅端通道并注册到全局列表 func (s *GRPCServer) NewClientChannel() *clientChan { cliCh := &clientChan{ ch: make(chan *SendMessageRequest, 10), // 可根据业务调整缓冲大小 done: make(chan struct{}), } s.chMu.Lock() s.clientChannels[cliCh] = struct{}{} s.chMu.Unlock() return cliCh }
2. 改造订阅端逻辑,退出时主动清理资源
在订阅端的流处理方法中增加defer逻辑,退出时主动注销通道、关闭资源,同时监听流上下文的取消信号:
func (s *GRPCServer) WatchMessageServer(req *WatchMessageRequest, stream ExampleService_WatchMessageServer) error { cliCh := s.NewClientChannel() // 退出前自动清理注册信息和通道资源 defer func() { // 从全局列表删除当前订阅端 s.chMu.Lock() delete(s.clientChannels, cliCh) s.chMu.Unlock() // 仅关闭一次资源 cliCh.closeOnce.Do(func() { close(cliCh.done) close(cliCh.ch) }) }() for { select { // 监听客户端断连信号 case <-stream.Context().Done(): return stream.Context().Err() case msg, ok := <-cliCh.ch: if !ok { return nil } // 发送失败会触发defer逻辑,自动清理资源 if err := stream.Send(msg); err != nil { return err } } } }
3. 改造发布端逻辑,发送前校验通道状态
发布端发送消息时通过select同时监听退出信号,避免往已关闭的通道写入:
func (s *GRPCServer) SendMessage(ctx context.Context, req *SendMessageRequest) (*emptypb.Empty, error) { s.chMu.RLock() defer s.chMu.RUnlock() for cliCh := range s.clientChannels { select { // 跳过已退出的订阅端 case <-cliCh.done: continue // 发送成功直接返回 case cliCh.ch <- req: // 可选:避免消费慢的订阅端阻塞发布流程,可根据业务决定是否保留default default: // 可在这里记录慢订阅日志、或做积压告警 } } return &emptypb.Empty{}, nil }
关键注意点
- 全程遵守Go channel关闭原则,仅通过sync.Once保证资源唯一关闭,不会出现重复关闭、或往关闭通道写入的panic
- 读写锁隔离全局通道列表的读写操作,避免并发读写map的panic
- 用map存储订阅端通道比slice删除效率高,适合高频订阅/取消订阅的聊天室场景
- 发送逻辑的default分支可根据业务取舍:要求不丢消息就去掉default、调大通道缓冲;要求发布端低延迟就保留default,丢弃积压消息或做异步补偿。
内容的提问来源于stack exchange,提问作者xakepp35
相关产品推荐
相关产品推荐

