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

基于Redigo的聊天服务器:如何在订阅阻塞时监听ctx.Done()终止进程?

当然可行!这是Go中处理阻塞操作与上下文取消的典型场景

你完全可以通过goroutine + select语句来实现同时监听ctx.Done()信号和Redis订阅的阻塞接收操作,既不会让psc.Receive()阻塞主协程,又能在上下文取消时及时终止整个进程。

具体实现思路

核心是把阻塞的psc.Receive()放到单独的goroutine中,将收到的消息/错误转发到通道,然后在主协程用select同时监听ctx.Done()和这些通道,这样就能响应上下文取消信号了。

示例代码

func handleWebSocketClient(ctx context.Context, wsConn *websocket.Conn, redisPool *redis.Pool) {
    // 从连接池获取Redis连接
    redisConn, err := redisPool.GetContext(ctx)
    if err != nil {
        log.Printf("Failed to get Redis connection: %v", err)
        wsConn.Close()
        return
    }
    defer redisConn.Close()

    // 初始化Redis订阅客户端
    psc := redis.PubSubConn{Conn: redisConn}
    defer psc.Close()

    // 订阅目标频道
    if err := psc.Subscribe("chat:global"); err != nil {
        log.Printf("Failed to subscribe to channel: %v", err)
        wsConn.Close()
        return
    }

    // 创建用于传递消息和错误的通道
    msgChan := make(chan interface{}, 1)
    errChan := make(chan error, 1)

    // 启动goroutine处理Redis的阻塞接收
    go func() {
        for {
            resp := psc.Receive()
            switch v := resp.(type) {
            case redis.Message, redis.Subscription:
                msgChan <- v
            case error:
                errChan <- v
                return // 遇到错误时退出goroutine
            }
        }
    }()

    // 主循环:同时监听上下文取消、Redis消息和错误
    for {
        select {
        case <-ctx.Done():
            // 上下文取消,执行清理逻辑
            log.Println("Context canceled, closing client connections")
            psc.Unsubscribe() // 先取消订阅,避免Redis端残留订阅
            wsConn.Close()
            return
        case msg := <-msgChan:
            // 将Redis消息转发到WebSocket
            switch m := msg.(type) {
            case redis.Message:
                if err := wsConn.WriteMessage(websocket.TextMessage, m.Data); err != nil {
                    log.Printf("Failed to send message to WebSocket: %v", err)
                    return
                }
            case redis.Subscription:
                // 可以处理订阅确认、取消等通知(可选)
                log.Printf("Subscription status: %s channel %s", v.Kind, v.Channel)
            }
        case err := <-errChan:
            log.Printf("Redis subscription error: %v", err)
            wsConn.Close()
            return
        }
    }
}

关键注意点

  • 避免goroutine泄漏:当ctx.Done()触发时,一定要调用psc.Close()或psc.Unsubscribe(),这样后台的Receive() goroutine会收到一个关闭错误,从而退出循环,不会一直阻塞。
  • 资源及时释放:用defer确保Redis连接和WebSocket连接在函数退出时被关闭,避免资源泄漏。
  • 上下文传递:确保创建WebSocket协程时传入的上下文是可取消的(比如用context.WithCancel或context.WithTimeout),这样当用户断开WebSocket连接时,可以主动取消上下文,触发整个清理流程。

内容的提问来源于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:42:36