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

基于hiredis的redis-plus-plus订阅异常:Topic D无法订阅的问题

Redis-plus-plus 订阅HTTP触发的主题D失败问题排查与解决

核心原因

redis-plus-plus的Subscriber对象不支持多线程并发操作。你在RestServer的HTTP处理线程中调用subscribe("D"),同时消费线程在执行subscriber->consume()、主线程也调用过subscribe("C"),多线程对同一个Subscriber实例的并发操作会导致内部网络缓冲区的竞态问题,最终使得订阅D的指令没有被正确提交到Redis服务器。

为什么A/B/C能正常工作

  • A、B是在消费线程启动前调用的,属于单线程操作,没有竞态问题;
  • C是主线程延迟调用的,刚好消费线程处于sleep(10ms)的非阻塞阶段,此时对缓冲区的写入没有被干扰,后续consume()执行时将C的订阅指令发送给了Redis;
  • 而HTTP请求的触发时机是随机的,大概率会和消费线程的consume()操作并发,直接破坏了缓冲区的订阅指令,导致D的订阅请求丢失。

解决方案

方案1:添加互斥锁保护所有Subscriber操作

给RedisSubScribeServer新增一个互斥锁,确保所有对Subscriber的操作(包括subscribe和consume)都互斥执行:

class RedisSubScribeServer {
public: 
    // ... 原有代码
    void subscribe(std::string topic){
        std::lock_guard<std::mutex> lock(sub_mtx);
        subscriber->subscribe(std::move(topic));
    }

    void startServer() {
        {
            std::lock_guard<std::mutex> lock(sub_mtx);
            subscriber->subscribe("A");
            subscriber->subscribe("B");

            subscriber->on_message([this](std::string channel, std::string msg) {
                handleMessage(channel, msg);
            });
        }

        if (!consume_thread.joinable()) {
            consume_thread = std::thread([this]() {
                while (!stop_flag) {
                    try {
                        std::lock_guard<std::mutex> lock(sub_mtx);
                        if (subscriber) {
                            subscriber->consume();
                        }
                        std::this_thread::sleep_for(std::chrono::milliseconds(10));                          
                    } catch (const std::exception& e) {
                        std::cerr << e.what() << std::endl;
                        std::this_thread::sleep_for(std::chrono::seconds(1));
                    }
                }
            });
        }
        // ... 原有代码
    }
private:
    // ... 原有成员
    std::mutex sub_mtx; // 保护Subscriber的互斥锁
};

方案2:用线程安全队列转发订阅请求(更推荐)

将所有订阅请求放入线程安全队列,由消费线程统一处理订阅操作,彻底避免多线程直接操作Subscriber:

class RedisSubScribeServer {
public: 
    // ... 原有代码
    void subscribe(std::string topic){
        std::lock_guard<std::mutex> lock(queue_mtx);
        subscribe_queue.push(std::move(topic));
        cv.notify_one(); // 唤醒消费线程,及时处理订阅请求
    }

    void startServer() {
        subscriber->subscribe("A");
        subscriber->subscribe("B");

        subscriber->on_message([this](std::string channel, std::string msg) {
            handleMessage(channel, msg);
        });

        if (!consume_thread.joinable()) {
            consume_thread = std::thread([this]() {
                while (!stop_flag) {
                    try {
                        // 先处理队列中的订阅请求
                        process_subscribe_queue();

                        if (subscriber) {
                            // 用带超时的consume,避免一直阻塞,及时响应新的订阅请求
                            subscriber->consume(std::chrono::milliseconds(100));
                        }
                    } catch (const std::exception& e) {
                        std::cerr << e.what() << std::endl;
                        std::this_thread::sleep_for(std::chrono::seconds(1));
                    }
                }
            });
        }
        // ... 原有代码
    }
private:
    void process_subscribe_queue() {
        std::lock_guard<std::mutex> lock(queue_mtx);
        while (!subscribe_queue.empty()) {
            auto topic = std::move(subscribe_queue.front());
            subscribe_queue.pop();
            subscriber->subscribe(topic);
        }
    }

    // ... 原有成员
    std::queue<std::string> subscribe_queue;
    std::mutex queue_mtx;
    std::condition_variable cv;
};

验证

修改后重新启动服务:

  1. 通过HTTP请求订阅主题D;
  2. 执行redis-cli PUBSUB CHANNELS,确认D出现在订阅列表中;
  3. 执行redis-cli PUBLISH D "test",检查是否能收到消息回调。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 11:53:22