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

如何实现无锁入队、可阻塞出队的线程间数据交换队列?

解决方案:无锁发送+阻塞接收的动态扩容队列

针对你的需求(渲染线程无锁发送、接收线程休眠等待、动态扩容、支持boost::intrusive_ptr),以下是基于STL和Boost的两种可行实现方案:

方案一:Boost.Lockfree容器 + 条件变量(适配动态扩容)

Boost.Lockfree的栈/队列是固定大小的,我们可以通过无锁栈存储动态块的方式实现动态扩容,同时保证发送端的无锁特性:

  • 发送端:维护一个无锁栈存储数据块,绝大多数情况下直接无锁写入当前块;仅当块满时才分配新块(罕见操作,可能触发内存管理器锁)并无锁推入栈,写入后调用非阻塞的条件变量通知。
  • 接收端:循环尝试从栈中取块处理,无数据时通过条件变量休眠,收到通知后重试。

代码示例

#include <boost/lockfree/stack.hpp>
#include <condition_variable>
#include <mutex>
#include <vector>
#include <boost/intrusive_ptr.hpp>

// T为你的数据类型,需已实现intrusive_ptr的引用计数操作
template<typename T>
class LockfreeBlockingQueue {
public:
    using DataBlock = std::vector<boost::intrusive_ptr<T>>;
    static constexpr size_t BLOCK_CAPACITY = 256; // 单块初始容量,可按需调整

    void push(boost::intrusive_ptr<T> item) {
        std::lock_guard<std::mutex> block_lock(block_mtx_);
        // 当前块满时,分配新块并推入无锁栈
        if (current_block_->size() >= BLOCK_CAPACITY) {
            current_block_.reset(new DataBlock);
            stack_.push(current_block_.get());
        }
        current_block_->push_back(std::move(item));
        cv_.notify_one(); // 非阻塞通知,不影响渲染线程性能
    }

    boost::intrusive_ptr<T> pop() {
        while (true) {
            DataBlock* block = nullptr;
            // 先尝试从无锁栈中取块
            if (stack_.pop(block)) {
                std::unique_ptr<DataBlock> block_guard(block);
                while (!block->empty()) {
                    auto item = std::move(block->back());
                    block->pop_back();
                    return item;
                }
            } else {
                // 栈为空时检查当前块是否有剩余数据
                std::lock_guard<std::mutex> block_lock(block_mtx_);
                if (!current_block_->empty()) {
                    auto item = std::move(current_block_->back());
                    current_block_->pop_back();
                    return item;
                }
                // 无数据时进入休眠
                std::unique_lock<std::mutex> wait_lock(wait_mtx_);
                cv_.wait(wait_lock, [this]() {
                    // 双重检查:栈或当前块是否有数据
                    if (!stack_.empty()) return true;
                    std::lock_guard<std::mutex> block_lock(block_mtx_);
                    return !current_block_->empty();
                });
            }
        }
    }

private:
    boost::lockfree::stack<DataBlock*> stack_{boost::lockfree::capacity(16)}; // 存储块的无锁栈,初始容量足够即可
    std::unique_ptr<DataBlock> current_block_{new DataBlock};
    std::mutex block_mtx_; // 仅在切换/检查当前块时使用,竞争极少
    std::mutex wait_mtx_;
    std::condition_variable cv_;
};

核心特点

  • 发送端仅在块满时触发一次锁操作,其余时间完全无锁,符合渲染线程的低延迟要求。
  • 接收端无数据时休眠,彻底避免忙循环浪费CPU。
  • 动态扩容通过分配新块实现,属于罕见操作,内存管理器的锁可接受。
  • 天然支持boost::intrusive_ptr,向量存储该类型时会正确执行析构逻辑。

方案二:原子链表 + 条件变量

如果不想依赖Boost.Lockfree,可以用std::atomic实现无锁链表,结合条件变量实现阻塞等待:

  • 发送端:用CAS操作无锁插入节点到链表头部,插入后调用条件变量通知。
  • 接收端:无锁弹出链表头部节点,无节点时休眠等待。

代码示例

#include <atomic>
#include <condition_variable>
#include <mutex>
#include <boost/intrusive_ptr.hpp>

template<typename T>
struct QueueNode {
    boost::intrusive_ptr<T> data;
    std::atomic<QueueNode*> next{nullptr};
};

template<typename T>
class AtomicBlockingQueue {
public:
    void push(boost::intrusive_ptr<T> item) {
        QueueNode<T>* new_node = new QueueNode<T>{std::move(item)};
        // CAS无锁插入链表头部
        QueueNode<T>* old_head = head_.load(std::memory_order_relaxed);
        do {
            new_node->next.store(old_head, std::memory_order_relaxed);
        } while (!head_.compare_exchange_weak(old_head, new_node,
                    std::memory_order_release, std::memory_order_relaxed));
        cv_.notify_one(); // 非阻塞通知
    }

    boost::intrusive_ptr<T> pop() {
        while (true) {
            QueueNode<T>* old_head = head_.load(std::memory_order_relaxed);
            if (!old_head) {
                // 无数据时进入休眠
                std::unique_lock<std::mutex> lock(wait_mtx_);
                cv_.wait(lock, [this]() { return head_.load(std::memory_order_relaxed) != nullptr; });
                continue;
            }
            // CAS无锁弹出头部节点
            if (head_.compare_exchange_weak(old_head, old_head->next.load(std::memory_order_relaxed),
                        std::memory_order_acquire, std::memory_order_relaxed)) {
                auto data = std::move(old_head->data);
                delete old_head;
                return data;
            }
        }
    }

    ~AtomicBlockingQueue() {
        // 清理剩余节点
        QueueNode<T>* node = head_.load(std::memory_order_relaxed);
        while (node) {
            QueueNode<T>* next = node->next.load(std::memory_order_relaxed);
            delete node;
            node = next;
        }
    }

private:
    std::atomic<QueueNode<T>*> head_{nullptr};
    std::mutex wait_mtx_;
    std::condition_variable cv_;
};

核心特点

  • 发送端完全无锁(除了内存分配的罕见锁),CAS操作性能极高,适合渲染线程。
  • 接收端无数据时休眠,避免无效CPU占用。
  • 天然支持动态扩容,每个节点单独分配,无需预设队列大小。
  • 兼容boost::intrusive_ptr,节点销毁时会正确处理引用计数。

关键注意事项

  • 虚假唤醒处理:接收端的条件变量等待必须配合循环检查,避免系统虚假唤醒导致错误返回空数据。
  • 内存序正确性:原子操作的内存序(如release/acquire)需保证线程间数据可见性,避免逻辑错误。
  • 内存泄漏防护:两种方案都需确保数据块/节点被正确释放,避免内存泄漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 04:55:22