无法通过mpsc通道向select!循环发消息及Redis PUBSUB断线无通知问题
Redis集成大型消息Broker的两个问题解决方案
问题1:Redis PUBSUB断开连接无通知,pubsub_stream.next()永久阻塞
原因
你当前使用的是redis-rs的同步订阅API,同步方法会阻塞线程,且无内置心跳/超时机制感知连接断开。异步场景下用同步API会占用Worker线程,无法主动检测连接状态,导致失效连接上永久阻塞。
解决方案
改用redis-rs的异步订阅API,结合Tokio超时机制检测连接有效性:
- 使用
AsyncClient建立异步连接,通过subscribe获取异步订阅流 - 用
tokio::time::timeout包装next_message()调用,超时则判定连接失效 - 连接出错或超时直接进入重连逻辑
示例代码片段:
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!循环接收
原因
- 同步API阻塞Tokio Worker线程:你在异步任务中调用了同步的
connection.subscribe,该方法会永久占用Worker线程,导致tokio::spawn的select循环任务无法被调度执行。 - 同步回调中跨线程发送消息时,若
Msg未正确实现Sendtrait会导致发送失败,但你未看到错误日志,核心原因是线程阻塞。
解决方案
- 完全改用异步Redis API:替换同步连接/订阅方法为异步版本,避免阻塞Worker线程
- 调整任务结构:确保所有操作非阻塞,让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
相关产品推荐
相关产品推荐

