使用redisrs异步订阅Redis时重复接收两条消息问题排查
Redis订阅重复推送消息问题排查与解决方案
问题现象
基于Rust实现的Redis订阅任务中,当Redis频道有消息时,会向Socket.IO客户端重复推送两条消息,redis_subscription_handler中衍生的异步任务循环会执行两次。
完整实现代码
use chrono::Local; use futures_util::StreamExt as _; use serde_json::{json, Value}; use socketioxide::extract::{Data, SocketRef}; use std::sync::Arc; use tokio::sync::Mutex; use tokio::time::Duration; use tracing::{error, info}; use crate::redis::Redis; #[derive(serde::Deserialize, Debug)] pub struct Message { #[serde(rename = "roomId")] room_id: String, message: String, } pub struct SocketHandlers; impl SocketHandlers { pub async fn handle_sockets(socket: SocketRef, redis: Arc<Mutex<Redis>>) { let redis_clone = Arc::clone(&redis); socket.on( "join-room", move |socket: SocketRef, Data::<String>(room)| { let redis_clone = Arc::clone(&redis_clone); async move { handle_join_room(socket, room, redis_clone).await; } }, ); let redis_clone = Arc::clone(&redis); socket.on( "send-message", move |socket: SocketRef, Data::<Value>(message)| { let redis_clone = Arc::clone(&redis_clone); async move { handle_send_message(socket, message, redis_clone).await; } }, ); let redis_clone = Arc::clone(&redis); if let Err(err) = redis_subscription_handler(redis_clone, socket, Duration::from_secs(1)).await { error!("Error in subscription handler: {}", err); } } } async fn handle_join_room(socket: SocketRef, room: String, redis: Arc<Mutex<Redis>>) { info!("Joining room: {:?}", room); let _ = socket.leave_all(); if let Err(e) = socket.join(room.clone()) { error!("Error joining room: {:?}", e); } let message = json!({ "room_id": room, "message": format!("{} joined the room", socket.id), "sender": "Server", "time": Local::now().to_string() }); let str_msg = match serde_json::to_string(&message) { Ok(v) => v, Err(e) => { error!("Error parsing message: {:?}", e); return; } }; let redis_guard = redis.lock().await; if let Err(e) = redis_guard.publish_message("Messages", &str_msg).await { error!("Error publishing message to Redis: {:?}", e); } } async fn handle_send_message(socket: SocketRef, message: Value, redis: Arc<Mutex<Redis>>) { info!("Message: {:?}", message); //parse message with serde_json let parsed_message: Message = match serde_json::from_value::<Message>(message.clone()) { Ok(v) => v, Err(e) => { error!("Error parsing message: {:?}", e); return; } }; let message = json!({ "room_id": parsed_message.room_id, "message": parsed_message.message, "sender": socket.id, "time": Local::now().to_string() }); let str_msg = match serde_json::to_string(&message) { Ok(v) => v, Err(e) => { error!("Error parsing message: {:?}", e); return; } }; let redis_guard = redis.lock().await; if let Err(e) = redis_guard.publish_message("Messages", &str_msg).await { error!("Error publishing message to Redis: {:?}", e); } } pub async fn redis_subscription_handler( redis: Arc<Mutex<Redis>>, socket: SocketRef, delay: Duration, ) -> Result<(), Box<dyn std::error::Error>> { let redis_guard = redis.lock().await; let mut pubsub = redis_guard .client .get_tokio_connection() .await .unwrap() .into_pubsub(); let _ = pubsub.subscribe("Messages").await; // Spawn an asynchronous task tokio::spawn(async move { // Inside the spawned task, use a loop to continuously process messages while let Some(msg) = pubsub.on_message().next().await { // Parse the payload let payload: Value = match serde_json::from_str(&msg.get_payload::<String>().unwrap()) { Ok(payload) => payload, Err(err) => { error!("Error parsing JSON payload: {}", err); continue; // Skip to the next iteration if parsing fails } }; info!("Received message: {:?}", payload); // Extract room_id from payload let room_id = match payload["room_id"].as_str() { Some(room_id) => room_id.to_string(), None => { error!("Error: 'room_id' field missing in payload"); continue; // Skip to the next iteration if 'room_id' is missing } }; // Emit the message to the socket socket.within(room_id).emit("message", payload).unwrap(); // Introduce a delay between processing each message tokio::time::sleep(delay).await; } }); Ok(()) }
问题原因分析
- 重复订阅Redis频道:每次客户端连接时,
handle_sockets方法都会调用redis_subscription_handler,该函数会创建新的PubSub连接并订阅"Messages"频道。若客户端触发多次连接(比如Socket.IO重连机制),会导致多个订阅任务同时运行,每条消息被多次接收推送。 - Socket连接生命周期管理缺失:未对Socket连接的订阅状态做标记,同一客户端连接可能被重复绑定订阅任务。
- Redis连接未复用:每次订阅都创建新的PubSub连接,增加了重复订阅的概率。
解决方案
方案1:全局单一PubSub连接分发消息
将Redis订阅逻辑从客户端连接处理中剥离,用全局PubSub连接统一接收消息,再通过Socket.IO房间机制分发,避免每个客户端创建独立订阅:
// 应用启动时初始化全局订阅 pub async fn init_global_redis_subscription(redis: Arc<Mutex<Redis>>) { let redis_guard = redis.lock().await; let mut pubsub = redis_guard .client .get_tokio_connection() .await .unwrap() .into_pubsub(); let _ = pubsub.subscribe("Messages").await; tokio::spawn(async move { while let Some(msg) = pubsub.on_message().next().await { let payload: Value = match serde_json::from_str(&msg.get_payload::<String>().unwrap()) { Ok(payload) => payload, Err(err) => { error!("Error parsing JSON payload: {}", err); continue; } }; info!("Received message: {:?}", payload); let room_id = match payload["room_id"].as_str() { Some(room_id) => room_id.to_string(), None => { error!("Error: 'room_id' field missing in payload"); continue; } }; // 全局广播到指定房间 socketioxide::Socket::broadcast().within(room_id).emit("message", payload).unwrap(); } }); }
修改handle_sockets,移除其中的redis_subscription_handler调用,仅在应用启动时初始化一次全局订阅。
方案2:标记Socket订阅状态,避免重复创建任务
在Socket连接生命周期内,用扩展数据标记是否已初始化订阅,确保每个Socket只绑定一次订阅任务:
// 修改handle_sockets方法 pub async fn handle_sockets(socket: SocketRef, redis: Arc<Mutex<Redis>>) { // 检查是否已创建订阅,未创建则初始化 if socket.extensions().get::<bool>().is_none() { socket.extensions().insert(true); let redis_clone = Arc::clone(&redis); if let Err(err) = redis_subscription_handler(redis_clone, socket.clone(), Duration::from_secs(1)).await { error!("Error in subscription handler: {}", err); } } // 原有join-room、send-message逻辑保持不变 let redis_clone = Arc::clone(&redis); socket.on( "join-room", move |socket: SocketRef, Data::<String>(room)| { let redis_clone = Arc::clone(&redis_clone); async move { handle_join_room(socket, room, redis_clone).await; } }, ); let redis_clone = Arc::clone(&redis); socket.on( "send-message", move |socket: SocketRef, Data::<Value>(message)| { let redis_clone = Arc::clone(&redis_clone); async move { handle_send_message(socket, message, redis_clone).await; } }, ); }
同时需确保订阅任务在Socket断开时自动取消,避免资源泄漏。
方案3:过滤发布者自身的消息
Redis PubSub默认会让发布者收到自己发布的消息,若业务无需此逻辑,可在推送时过滤:
// 在订阅任务的推送步骤添加过滤 let sender_id = payload["sender"].as_str().unwrap_or(""); if socket.id() != sender_id { socket.within(room_id).emit("message", payload).unwrap(); }
内容的提问来源于stack exchange,提问作者Gaurav D. Lohar
相关产品推荐
相关产品推荐

