无锁环形队列设计与实现(SPSC 与 MPMC)
生产者—消费者是并发程序中最常见的协作模式之一。用互斥锁保护队列,在竞争激烈时会带来线程阻塞与上下文切换的开销;无锁环形队列改用原子操作和内存序保证线程安全,读端与写端互不阻塞。本文从最简单的单生产者单消费者(SPSC)环形队列讲起,说明空满区分与内存序的配对方式;再分析多生产者多消费者(MPMC)的实现难点,给出 Dmitry Vyukov 提出的按槽序号(sequence)方案的完整实现。适合需要在 C++ 中手写或读懂无锁队列的开发者阅读。
SPSC:数据结构
单生产者单消费者是最简单的并发模型:生产者是写索引的唯一修改者,消费者是读索引的唯一修改者。这一「索引各归一线程」的性质,是 SPSC 队列可以用极少的同步机制实现的前提。
- 固定容量数组(buffer):存储元素的连续内存空间,大小为 capacity + 1。额外的一个空槽用于区分「满」和「空」两种状态。
- 读索引(head)与写索引(tail):
head指向下一个要读取的位置,tail指向下一个要写入的位置,均为std::atomic<size_t>。 - 循环递增:索引按模 capacity + 1 递增,实现缓冲区的循环利用。
用「空一格」方案区分空满:队列满的条件是 (tail + 1) % (capacity + 1) == head,队列空的条件是 head == tail。如果不用这个额外空槽,head == tail 将同时表示空和满,需要额外的标志位或计数器才能区分。
SPSC:操作流程与内存序
写入(push):
- 读取当前
tail; - 计算
next_tail = (tail + 1) % (capacity + 1); - 若
next_tail == head,队列已满,返回失败; - 将数据写入
buffer[tail]; - 以
memory_order_release更新tail = next_tail。
读取(pop):
- 读取当前
head; - 若
head == tail,队列为空,返回失败; - 读取
buffer[head]; - 以
memory_order_release更新head = (head + 1) % (capacity + 1)。
两个方向的 release/acquire 配对,是正确性的核心:
- tail 方向:生产者写完数据后
tail.store(release);消费者tail.load(acquire)读到新 tail 时,保证能看到对应槽位上已完整写入的数据。 - head 方向:消费者读完数据后
head.store(release);生产者head.load(acquire)读到新 head 后才复用该槽位,保证覆盖前消费者已经读完旧数据。
反过来,读写自己的索引(生产者读 tail、消费者读 head)用 memory_order_relaxed 即可:该索引只有一个写者,不存在与其他线程同步的问题。如果索引完全不用原子类型,读写索引构成数据竞争,在 C++ 标准下是未定义行为,可能出现撕裂读或索引更新丢失。
SPSC:完整实现
#include <atomic>
#include <vector>
#include <optional>
#include <cassert>
template<typename T>
class LockFreeRingBuffer {
public:
explicit LockFreeRingBuffer(size_t capacity)
: capacity_(capacity + 1), // 一个额外空间用来区分满和空
buffer_(capacity_),
head_(0),
tail_(0) {
assert(capacity > 0);
}
// 禁止复制
LockFreeRingBuffer(const LockFreeRingBuffer&) = delete;
LockFreeRingBuffer& operator=(const LockFreeRingBuffer&) = delete;
// 生产者调用:插入元素,成功返回true,满返回false
bool push(const T& item) {
size_t tail = tail_.load(std::memory_order_relaxed);
size_t next_tail = increment(tail);
if (next_tail == head_.load(std::memory_order_acquire)) {
// 缓冲区满
return false;
}
buffer_[tail] = item;
tail_.store(next_tail, std::memory_order_release);
return true;
}
// 生产者移动语义版本
bool push(T&& item) {
size_t tail = tail_.load(std::memory_order_relaxed);
size_t next_tail = increment(tail);
if (next_tail == head_.load(std::memory_order_acquire)) {
return false;
}
buffer_[tail] = std::move(item);
tail_.store(next_tail, std::memory_order_release);
return true;
}
// 消费者调用:取出元素,存在时返回 optional<T>,无数据时返回 empty optional
std::optional<T> pop() {
size_t head = head_.load(std::memory_order_relaxed);
if (head == tail_.load(std::memory_order_acquire)) {
// 缓冲区空
return std::nullopt;
}
T item = std::move(buffer_[head]);
head_.store(increment(head), std::memory_order_release);
return item;
}
// 是否空
bool empty() const {
return head_.load(std::memory_order_acquire) == tail_.load(std::memory_order_acquire);
}
// 容量(实际可用容量为 capacity_ - 1)
size_t capacity() const {
return capacity_ - 1;
}
private:
size_t increment(size_t idx) const noexcept {
return (idx + 1) % capacity_;
}
private:
const size_t capacity_;
std::vector<T> buffer_;
std::atomic<size_t> head_;
std::atomic<size_t> tail_;
}; 为什么 MPMC 复杂得多
一旦存在多个生产者或多个消费者,「索引各归一线程」不再成立:
- 同一索引被并发修改:多个生产者同时推进
tail,必须用 CAS(Compare-And-Swap)循环保证推进操作的原子与串行化,失败方重试。 - ABA 问题:一个值从 A 改为 B 又改回 A,依赖「值未变化」判断的 CAS 会误判,导致同一槽位被重复占用或覆盖。
- 槽位状态管理:多个生产者必须保证每次写入独占一个槽位;多个消费者必须保证每个元素只被消费一次。空、满、写入中、可读等状态需要额外的标记。
- 虚假共享(false sharing):多个线程频繁读写落在同一 CPU 缓存行上的计数器,会导致缓存行反复失效,显著影响性能。
链表型的经典解法是 Michael-Scott 队列(1996),但它引入了已出队节点的内存回收问题,需要延迟回收机制配合。环形数组方案则把问题转化为「每个槽位的状态管理」,回避了动态内存回收。
MPMC:Vyukov 的按槽序号方案
Dmitry Vyukov 在 1024cores.net 发表的 bounded MPMC queue 给出了一种被广泛借鉴的设计,核心是给每个槽位一个独立序号:
- 初始化:第
i个槽的sequence为i;两个全局位置计数enqueue_pos_(生产者推进)与dequeue_pos_(消费者推进)从 0 开始。 - 生产者抢占:要写入位置
pos时,检查槽序号seq与pos的差:diff == 0:该槽在当前这一圈空闲,用 CAS 把enqueue_pos_从pos抢到pos + 1,成功者独占该槽;diff < 0:消费者还没释放该槽(序号落后于位置),队列已满;diff > 0:其他生产者已经抢先推进了位置,重新加载enqueue_pos_再试。
- 写入与发布分离:抢占成功后写入数据,再把该槽
sequence发布为pos + 1——消费者只有在看到seq == pos + 1时才认为数据有效。 - 消费者抢占:要读取位置
pos时,检查seq == pos + 1表示数据有效,同样用 CAS 抢dequeue_pos_;读完后把sequence置为pos + capacity,即该槽在下一圈的位置,对未来的生产者重新可写。
这个设计的两个关键点:
- CAS 只抢位置计数,不携带数据。CAS 失败的线程只需重载位置重试;数据写入在独占槽位之后单独进行,用
release发布序号。 - 序号天然规避 ABA:
sequence等于「位置 + 已完成圈数 × capacity」,随圈数单调递增,同一槽位每一圈的序号都不同,旧状态不会被误认为当前状态。
需要说明的是,该算法严格意义上并非 lock-free:如果一个线程在抢占槽位之后、发布数据之前被长期挂起,等待该位置数据的消费者会停住。它是不使用互斥锁的非阻塞实现,工程语境中通常仍归入无锁队列;Vyukov 本人在原文中也注明了这一点。
MPMC:完整实现
#include <atomic>
#include <vector>
#include <cassert>
#include <optional>
template <typename T>
class MPMCRingBuffer {
public:
explicit MPMCRingBuffer(size_t capacity)
: buffer_(capacity), capacity_(capacity), index_mask_(capacity - 1)
{
// capacity必须是2的幂,这样方便用掩码计算index
assert((capacity >= 2) && ((capacity & (capacity - 1)) == 0));
for (size_t i = 0; i < capacity; ++i) {
buffer_[i].sequence.store(i, std::memory_order_relaxed);
}
enqueue_pos_.store(0, std::memory_order_relaxed);
dequeue_pos_.store(0, std::memory_order_relaxed);
}
// 禁止复制
MPMCRingBuffer(const MPMCRingBuffer&) = delete;
MPMCRingBuffer& operator=(const MPMCRingBuffer&) = delete;
bool push(const T& data) {
Cell* cell;
size_t pos = enqueue_pos_.load(std::memory_order_relaxed);
for (;;) {
cell = &buffer_[pos & index_mask_];
size_t seq = cell->sequence.load(std::memory_order_acquire);
intptr_t diff = (intptr_t)seq - (intptr_t)pos;
if (diff == 0) {
if (enqueue_pos_.compare_exchange_weak(pos, pos + 1, std::memory_order_relaxed)) {
break;
}
} else if (diff < 0) {
// 队列满
return false;
} else {
// 重新加载生产者位置,准备重试
pos = enqueue_pos_.load(std::memory_order_relaxed);
}
}
cell->data = data;
cell->sequence.store(pos + 1, std::memory_order_release);
return true;
}
std::optional<T> pop() {
Cell* cell;
size_t pos = dequeue_pos_.load(std::memory_order_relaxed);
for (;;) {
cell = &buffer_[pos & index_mask_];
size_t seq = cell->sequence.load(std::memory_order_acquire);
intptr_t diff = (intptr_t)seq - (intptr_t)(pos + 1);
if (diff == 0) {
if (dequeue_pos_.compare_exchange_weak(pos, pos + 1, std::memory_order_relaxed)) {
break;
}
} else if (diff < 0) {
// 队列空
return std::nullopt;
} else {
pos = dequeue_pos_.load(std::memory_order_relaxed);
}
}
T result = cell->data;
cell->sequence.store(pos + capacity_, std::memory_order_release);
return result;
}
private:
struct Cell {
std::atomic<size_t> sequence; // 序号,用于标志和版本控制槽状态
T data; // 数据存储
};
std::vector<Cell> buffer_; // 环形数组数据槽
size_t const capacity_; // 容量(2的幂)
size_t const index_mask_; // 用于快速取模(capacity_ - 1)
std::atomic<size_t> enqueue_pos_; // 生产者索引
std::atomic<size_t> dequeue_pos_; // 消费者索引
}; 实现细节上:容量要求为 2 的幂,这样 pos & index_mask_ 等价于取模且更快;槽的 sequence 用 acquire 读、release 写,与数据写入/读取构成发布—获取配对;生产者读取自己的 enqueue_pos_ 用 relaxed,因为 CAS 本身会仲裁位置归属。
要点
- SPSC 的最小正确实现只需要两个原子索引、一个额外空槽和两组 release/acquire 配对;能用 SPSC 的场景不必引入更复杂的机制。
- MPMC 的关键在于按槽独立序号:CAS 只用于抢占位置计数,数据写入与序号发布分离,序号随圈数单调递增从而规避 ABA。
- 「无锁」要区分语境:Vyukov 方案严格意义上不是 lock-free,但互不阻塞、无互斥锁,实践中满足多数高性能场景。
- 优先成熟库:无锁结构调试困难、内存序容易出错。
boost::lockfree::queue与boost::lockfree::spsc_queue、folly::MPMCQueue、moodycamel::ConcurrentQueue等实现经过充分验证,生产环境通常应先于手写实现考虑。