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

无法通过mpsc通道向select!循环发消息及Redis PUBSUB断线无通知问题

Redis集成大型消息Broker的两个问题解决方案

问题1:Redis PUBSUB断开连接无通知,pubsub_stream.next()永久阻塞

原因

你当前使用的是redis-rs的同步订阅API,同步方法会阻塞线程,且无内置心跳/超时机制感知连接断开。异步场景下用同步API会占用Worker线程,无法主动检测连接状态,导致失效连接上永久阻塞。

解决方案

改用redis-rs的异步订阅API,结合Tokio超时机制检测连接有效性:

  1. 使用AsyncClient建立异步连接,通过subscribe获取异步订阅流
  2. 用tokio::time::timeout包装next_message()调用,超时则判定连接失效
  3. 连接出错或超时直接进入重连逻辑

示例代码片段:

use redis::AsyncCommands;
use tokio::time::{timeout, Duration};

async fn async_redis_subscribe() {
    let client = redis::Client::open("redis://127.0.0.1:6379/").unwrap();
    let mut conn = client.get_async_connection().await.unwrap();
    
    // 获取异步订阅流
    let mut pubsub = conn.subscribe("tokio").await.unwrap();
    
    loop {
        // 10秒超时,触发则认为连接失效
        match timeout(Duration::from_secs(10), pubsub.next_message()).await {
            Ok(Ok(msg)) => {
                let channel = msg.get_channel_name();
                let payload: String = msg.get_payload().unwrap();
                println!("Received: {} -> {}", channel, payload);
            }
            Ok(Err(e)) => {
                eprintln!("订阅错误: {}", e);
                break; // 进入重连
            }
            Err(_) => {
                eprintln!("连接超时,断开重连");
                break;
            }
        }
    }
}

问题2:mpsc::unbounded_channel消息无法被select!循环接收

原因

  1. 同步API阻塞Tokio Worker线程:你在异步任务中调用了同步的connection.subscribe,该方法会永久占用Worker线程,导致tokio::spawn的select循环任务无法被调度执行。
  2. 同步回调中跨线程发送消息时,若Msg未正确实现Send trait会导致发送失败,但你未看到错误日志,核心原因是线程阻塞。

解决方案

  1. 完全改用异步Redis API:替换同步连接/订阅方法为异步版本,避免阻塞Worker线程
  2. 调整任务结构:确保所有操作非阻塞,让Tokio能正常调度多个异步任务

修正后的完整代码:

use redis::AsyncCommands;
use tokio::sync::{mpsc, broadcast};
use tokio::time::{IntervalStream, Duration, sleep};
use futures::stream::StreamExt;
use std::sync::Arc;

// 假设Storage和Event已定义
struct Storage {
    eb_broadcast_tx: broadcast::Sender<Event>,
}

enum Event {
    WsClientConnected { id: u64, name: String },
    WsClientDisconnected { id: u64, name: String },
}

pub async fn redis_async_task(storage: Arc<Storage>) {
    let mut eb_broadcast_rx = storage.eb_broadcast_tx.subscribe();
    let (mpsc_tx, mut mpsc_rx) = mpsc::unbounded_channel::<redis::Msg>();
    let mut interval_5s = IntervalStream::new(tokio::time::interval(Duration::from_secs(5)));    
    
    // 启动消息处理任务(不再被阻塞)
    let _task = tokio::spawn(async move {
        loop {
            tokio::select! {
                Some(msg) = mpsc_rx.recv() => {
                    let channel = msg.get_channel_name().to_string();
                    let payload = msg.get_payload::<String>().unwrap();
                    println!(" - 2 REDIS: subscription event: {} channel: {} payload: {}", channel, channel, payload);
                },
                Some(_ts) = interval_5s.next() => {
                    println!("timer");
                },
                Ok(evt) = eb_broadcast_rx.recv() => {
                    match evt {
                        Event::WsClientConnected{id: _, name: _} => {},
                        Event::WsClientDisconnected{id: _, name: _} => {},
                    }
                },
            }
        }
    });
    
    loop {
        println!("REDIS connecting");
        let client = redis::Client::open("redis://127.0.0.1:6379/").unwrap();
        match client.get_async_connection().await {
            Ok(mut conn) => {
                println!("REDIS connected");
                let mut pubsub = match conn.subscribe("tokio").await {
                    Ok(p) => p,
                    Err(e) => {
                        eprintln!("订阅失败: {}", e);
                        sleep(Duration::from_millis(1000)).await;
                        continue;
                    }
                };
                
                // 异步处理订阅消息,非阻塞
                while let Some(msg_result) = pubsub.next_message().await {
                    match msg_result {
                        Ok(msg) => {
                            if let Ok(payload) = msg.get_payload::<String>() {
                                let channel = msg.get_channel_name().to_string();
                                println!(" - 1 REDIS subscription event: channel: {} payload: {}", channel, payload);
                                if let Err(e) = mpsc_tx.send(msg) {
                                    eprintln!("发送到mpsc失败: {}", e);
                                }
                            }
                        }
                        Err(e) => {
                            eprintln!("订阅连接断开: {}", e);
                            break;
                        }
                    }
                }
            }
            Err(e) => {
                println!("REDIS连接失败: {}", e);
            }
        }
        sleep(Duration::from_millis(1000)).await;
    }
}

关键改进点

  • 用异步get_async_connection和subscribe替代同步API,释放Worker线程
  • 异步订阅流的next_message()非阻塞,允许Tokio正常调度select循环任务
  • 连接断开直接返回错误,无需依赖主动发命令触发断连

内容的提问来源于stack exchange,提问作者Daniel F

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 01:00:59