如何在Redigo阻塞接收消息时监听ctx.Done()终止进程?
当然可行!这是处理阻塞操作+上下文取消的标准方案
你的场景非常典型:每个WebSocket连接对应一个协程,用Redigo的PubSubConn.Receive()阻塞等待Redis消息,同时需要响应ctx.Done()信号来优雅终止协程。完全可以实现,这里给你两种常用的实现思路:
方法一:用Select同时监听上下文和Redis消息
把阻塞的psc.Receive()放到单独的goroutine里,然后用select在主协程中同时监听ctx.Done()、Redis消息和错误:
func handleUserConnection(ctx context.Context, ws *websocket.Conn, redisAddr, channel string) error { // 建立Redis连接 redisConn, err := redis.Dial("tcp", redisAddr) if err != nil { return fmt.Errorf("failed to connect to redis: %w", err) } defer redisConn.Close() psc := redis.PubSubConn{Conn: redisConn} defer psc.Close() // 订阅目标频道 if err := psc.Subscribe(channel); err != nil { return fmt.Errorf("failed to subscribe to channel %s: %w", channel, err) } // 启动goroutine接收Redis消息 msgChan := make(chan redis.Message, 1) errChan := make(chan error, 1) go func() { for { resp := psc.Receive() switch r := resp.(type) { case redis.Message: msgChan <- r case redis.Subscription: // 可选:处理订阅确认,比如打日志 log.Printf("successfully subscribed to %s, current subscribers: %d", r.Channel, r.Count) case error: errChan <- r return } } }() // 主循环:监听上下文取消和Redis消息 for { select { case <-ctx.Done(): log.Println("context canceled, shutting down connection") return ctx.Err() case msg := <-msgChan: // 把Redis消息转发到WebSocket if err := ws.WriteMessage(websocket.TextMessage, msg.Data); err != nil { return fmt.Errorf("failed to write to websocket: %w", err) } case err := <-errChan: return fmt.Errorf("redis receive error: %w", err) } } }
这种方式逻辑清晰,能同时处理多种事件(上下文取消、消息接收、错误),扩展性强。
方法二:利用上下文取消时关闭Redis连接
Redigo的PubSubConn.Close()会中断阻塞的Receive()调用,让它直接返回错误。我们可以监听ctx.Done(),触发时关闭PubSub连接,从而退出循环:
func handleUserConnection(ctx context.Context, ws *websocket.Conn, redisAddr, channel string) error { redisConn, err := redis.Dial("tcp", redisAddr) if err != nil { return fmt.Errorf("failed to connect to redis: %w", err) } defer redisConn.Close() psc := redis.PubSubConn{Conn: redisConn} defer psc.Close() if err := psc.Subscribe(channel); err != nil { return fmt.Errorf("failed to subscribe to channel %s: %w", channel, err) } // 启动goroutine监听上下文取消,关闭PubSub连接 go func() { <-ctx.Done() _ = psc.Close() // 关闭后,psc.Receive()会立即返回错误 }() // 消息接收循环 for { resp := psc.Receive() switch r := resp.(type) { case redis.Message: if err := ws.WriteMessage(websocket.TextMessage, r.Data); err != nil { return fmt.Errorf("failed to write to websocket: %w", err) } case redis.Subscription: log.Printf("subscribed to %s, subscribers count: %d", r.Channel, r.Count) case error: // 检查错误是否由上下文取消导致 select { case <-ctx.Done(): return ctx.Err() default: return fmt.Errorf("redis error: %w", r) } } } }
这种方式代码更简洁,利用了Redigo本身的连接特性,不需要额外的通道,适合逻辑相对简单的场景。
注意事项
- 一定要用
defer确保Redis连接和PubSubConn被正确关闭,避免资源泄漏。 - 处理
redis.Subscription类型的响应,这是Redis返回的订阅/取消订阅确认,虽然不是业务消息,但可以用来做状态校验或日志。 - 两种方式都能保证
ctx.Done()触发时,协程能快速退出,不会一直卡在psc.Receive()上。
内容的提问来源于stack exchange,提问作者Sam White
相关产品推荐
相关产品推荐

