You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 08:45:00