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

如何实现多线程安全消息队列的多路阻塞等待?

实现多线程安全消息队列的监听函数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, &notified]() {
        std::lock_guard lock(mtx);
        notified = true;
        cv.notify_one();
    };

    (queues.register_on_push_callback(callback), ...);

    // 阻塞等待直到任意队列触发回调
    std::unique_lock lock(mtx);
    cv.wait(lock, [&notified](){ 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 13:24:50