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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 09:45:05