基于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; };
验证
修改后重新启动服务:
- 通过HTTP请求订阅主题D;
- 执行
redis-cli PUBSUB CHANNELS,确认D出现在订阅列表中; - 执行
redis-cli PUBLISH D "test",检查是否能收到消息回调。
内容的提问来源于stack exchange,提问作者user29615459
相关产品推荐
相关产品推荐

