如何实现多线程安全消息队列的多路阻塞等待?
实现多线程安全消息队列的监听函数waitMsg
现有一个线程安全的消息队列实现(如下代码),需要实现一个waitMsg模板函数,使调用线程阻塞直到传入的任意一个队列不为空。
原线程安全消息队列代码
#pragma once #include <queue> #include <mutex> #include <condition_variable> template<class ElementType> class MsgQueue { public: /** Construction **/ MsgQueue() = default; // Forbid copying and moving MsgQueue(const MsgQueue &) = delete; MsgQueue(MsgQueue&&) = delete; MsgQueue &operator=(const MsgQueue &) = delete; MsgQueue &operator=(MsgQueue&&) = delete; public: /** Methods **/ void push(const ElementType &element) { { std::lock_guard lock(m_lock); m_messages.push(element); } m_notifier.notify_one(); } void push(ElementType &&element) { { std::lock_guard lock(m_lock); m_messages.push(std::move(element)); } m_notifier.notify_one(); } template<class... Args> void emplace(Args&&... args) { { std::lock_guard lock(m_lock); m_messages.emplace(std::forward<Args>(args)...); } m_notifier.notify_one(); } [[nodiscard]] ElementType &front() { std::unique_lock lock(m_lock); if(m_messages.empty()) { m_notifier.wait(lock, [&]{ return !m_messages.empty(); }); } return m_messages.front(); } void pop() { if(!m_messages.empty()) { std::lock_guard lock(m_lock); m_messages.pop(); } } [[nodiscard]] bool empty() const { return m_messages.empty(); } [[nodiscard]] size_t size() const { return m_messages.size(); } private: /** Members **/ std::queue<ElementType> m_messages; std::condition_variable m_notifier; std::mutex m_lock; };
单个队列的阻塞等待示例
MsgQueue<std::variant<int, float, std::string>> msgQueue; const auto &msg = msgQueue.front();
注意:原
front()函数存在悬空引用风险:函数返回时std::unique_lock销毁释放锁,其他线程可能立即调用pop()移除队列首元素,导致返回的引用指向已销毁的内存。同时原empty()和size()方法未加锁,存在线程安全问题,后续方案中已同步修复。
方案一:基于回调机制的实现(无额外线程开销)
这种方案通过给MsgQueue添加回调注册功能,当队列新增消息时触发回调,直接唤醒waitMsg线程。
步骤1:修改MsgQueue类
添加回调相关成员与方法,同时修复线程安全问题:
#pragma once #include <queue> #include <mutex> #include <condition_variable> #include <functional> #include <vector> template<class ElementType> class MsgQueue { public: /** Construction **/ MsgQueue() = default; // Forbid copying and moving MsgQueue(const MsgQueue &) = delete; MsgQueue(MsgQueue&&) = delete; MsgQueue &operator=(const MsgQueue &) = delete; MsgQueue &operator=(MsgQueue&&) = delete; public: /** Methods **/ // 注册回调:队列新增元素时触发 void register_on_push_callback(std::function<void()> callback) { std::lock_guard lock(m_callback_lock); m_callbacks.push_back(std::move(callback)); } // 清理所有注册的回调 void clear_callbacks() { std::lock_guard lock(m_callback_lock); m_callbacks.clear(); } void push(const ElementType &element) { { std::lock_guard lock(m_lock); m_messages.push(element); } notify_all(); } void push(ElementType &&element) { { std::lock_guard lock(m_lock); m_messages.push(std::move(element)); } notify_all(); } template<class... Args> void emplace(Args&&... args) { { std::lock_guard lock(m_lock); m_messages.emplace(std::forward<Args>(args)...); } notify_all(); } // 修复后的front():返回元素副本避免悬空引用 [[nodiscard]] ElementType front() { std::unique_lock lock(m_lock); m_notifier.wait(lock, [&]{ return !m_messages.empty(); }); auto elem = m_messages.front(); return elem; } void pop() { std::lock_guard lock(m_lock); if(!m_messages.empty()) { m_messages.pop(); } } [[nodiscard]] bool empty() const { std::lock_guard lock(m_lock); return m_messages.empty(); } [[nodiscard]] size_t size() const { std::lock_guard lock(m_lock); return m_messages.size(); } private: /** Members **/ void notify_all() { m_notifier.notify_one(); // 触发所有注册的回调 std::lock_guard lock(m_callback_lock); for(auto &cb : m_callbacks) { if(cb) cb(); } } std::queue<ElementType> m_messages; std::condition_variable m_notifier; std::mutex m_lock; // 回调相关成员 std::vector<std::function<void()>> m_callbacks; std::mutex m_callback_lock; };
步骤2:实现waitMsg函数
利用条件变量等待回调触发:
#include <condition_variable> #include <mutex> template<class... Queues> void waitMsg(Queues&... queues) { std::condition_variable cv; std::mutex mtx; bool notified = false; // 给每个队列注册回调:触发时标记状态并唤醒条件变量 auto callback = [&cv, &mtx, ¬ified]() { std::lock_guard lock(mtx); notified = true; cv.notify_one(); }; (queues.register_on_push_callback(callback), ...); // 阻塞等待直到任意队列触发回调 std::unique_lock lock(mtx); cv.wait(lock, [¬ified](){ return notified; }); // 清理回调,避免后续push触发无效通知 (queues.clear_callbacks(), ...); }
方案二:基于std::async的实现(最小化修改原队列)
如果不想大幅修改MsgQueue核心逻辑,可以用std::async为每个队列启动独立等待任务,等待第一个任务完成即返回。
步骤1:给MsgQueue添加等待方法
在原MsgQueue的public方法区添加:
void wait_until_non_empty() { std::unique_lock<std::mutex> lock(m_lock); m_notifier.wait(lock, [this](){ return !m_messages.empty(); }); }
步骤2:实现waitMsg函数
#include <future> #include <iterator> #include <chrono> template<class... Queues> void waitMsg(Queues&... queues) { // 为每个队列启动异步等待任务 auto futures = { std::async(std::launch::async, &Queues::wait_until_non_empty, &queues)... }; // 等待第一个任务完成(C++20可用std::chrono::forever,旧标准可替换为足够长的超时时间) std::wait_for(std::begin(futures), std::end(futures), std::chrono::forever); }
内容的提问来源于stack exchange,提问作者Caglayan Dokme
相关产品推荐
相关产品推荐

