欢迎来到并发编程的“诸神之战”——MPMC(Multi-Producer Multi-Consumer)。
这是并发领域最难啃的骨头,也是区分“普通程序员”和“系统架构师”的分水岭。你之前的 SPSC 是“单行道”,MPSC 是“多入口单出口”,而 MPMC 是**“多入口多出口的疯狂立交桥”**。
要把 MPMC 写好,靠 CAS 硬刚是不行的,那会让 CPU 缓存行“炸裂”。我们需要引入工业级杀器:LMAX Disruptor(无锁环形缓冲区)。
1. MPMC 的“两头肉搏”到底有多痛?
在 MPMC 中,所有生产者和消费者都盯着同一个头指针(Head)和同一个尾指针(Tail)。
- 生产者肉搏:N 个生产者疯狂 CAS 抢
Head(分配槽位)。每次 CAS 都是 RFO(Read For Ownership),导致Head所在的缓存行在 N 个核心间“乒乓”。 - 消费者肉搏:M 个消费者疯狂 CAS 抢
Tail(消费进度)。同样导致缓存行在 M 个核心间“乒乓”。 - 头尾互搏:更恐怖的是,
Head和Tail如果在同一个缓存行(False Sharing),生产者抢 Head 还会顺便把 Tail 的缓存行弄脏,直接拖垮消费者。
结果:吞吐量随着核心数增加不升反降,延迟飙升到微秒级。
2. 顶级框架 LMAX Disruptor 的“核心理念”
LMAX 是金融交易系统(每秒处理 600 万笔订单),它解决 MPMC 的思路极其高明,核心就两招:
- 干掉竞争(CAS):用顺序分配 + 栅栏(Barrier)代替互斥。生产者不抢同一个指针,而是各自申请一段连续的空间。
- 干掉伪共享:缓存行填充(Cache Line Padding),让 Head、Tail、Cursor 老死不相往来。
3. 核心机制:申请(Claim)与 发布(Publish)
Disruptor 不直接操作指针,而是通过序号(Sequence):
- Ring Buffer:一个巨大的数组(容量 2 的 N 次方,用位运算取模)。
- 生产者:
- 申请(Claim):调用
next(n),申请连续 n 个槽位。它只需要原子 CAS 修改ProducerCursor(类似于fetch_add),拿到一段独占的序号范围[lo, hi]。 - 填充数据:在内存中写入数据(此时消费者看不到,因为还没发布)。
- 发布(Publish):更新
PublishedCursor(带内存屏障 Release),告诉消费者:“这段序号的数据好了”。
- 申请(Claim):调用
- 消费者:
- 等待:通过**读屏障(Acquire)**读取
PublishedCursor。 - 批量拉取:如果
PublishedCursor > ConsumerCursor,拉取所有新数据。 - 处理:处理数据。
- 更新进度:更新
ConsumerCursor(标记已经消费到这里)。
- 等待:通过**读屏障(Acquire)**读取
4. 实战 C++ 极简核心逻辑
template<typename T, size_t SIZE = 1024> // SIZE 必须是 2 的幂
class MPMCQueue {
static_assert((SIZE & (SIZE - 1)) == 0, "SIZE must be power of 2");
static constexpr size_t MASK = SIZE - 1;
// 关键:强制 64 字节对齐,避免伪共享
alignas(64) std::atomic<size_t> m_producer_cursor {0}; // 发布序号
alignas(64) std::atomic<size_t> m_consumer_cursor {0}; // 消费序号
alignas(64) T m_buffer[SIZE]; // 存储真实数据
public:
bool enqueue(const T& item) {
// 1. 申请序号(唯一一次原子操作)
size_t pos = m_producer_cursor.fetch_add(1, std::memory_order_acquire);
size_t index = pos & MASK;
// 2. 检查是否覆盖了消费者(环形缓冲区判满)
// 注意:这里必须使用 acquire 读消费者,防止读到过期数据
if (m_consumer_cursor.load(std::memory_order_acquire) + SIZE <= pos) {
// 队列满了 -> 回退(注意:实际工业实现需要有回滚机制,此处简化为返回 false)
// 备注:Disruptor 使用多生产者时,这里通过等待策略解决
return false;
}
// 3. 写入数据(无锁,直接内存拷贝)
m_buffer[index] = item;
// 4. 发布:虽然 enqueue 只有一个屏障,但为了保证数据可见性,需要 Release 语义
// 实际为:m_producer_cursor 的 fetch_add 已经带了 Release,但为了极致读性能,需要让消费者能看到数据。
// 此处为了逻辑完整,我们用原子线程栅栏保证写入完成
std::atomic_thread_fence(std::memory_order_release);
return true;
}
bool dequeue(T& item) {
size_t pos = m_consumer_cursor.load(std::memory_order_relaxed);
// 快速判空(无竞争)
if (pos >= m_producer_cursor.load(std::memory_order_acquire)) {
return false;
}
size_t index = pos & MASK;
item = m_buffer[index]; // 读取数据
// 更新消费者游标(必须使用 Release 让生产者知道空间释放了)
m_consumer_cursor.store(pos + 1, std::memory_order_release);
return true;
}
};
代码分析:上面的代码为了展示原理省去了多生产者回滚逻辑,但它揭示了 Disruptor 最核心的**“以空间换时间”**思维。
5. 灵魂法宝:等待策略(Wait Strategy)
在 Disruptor 中,如果生产者太快,消费者跟不上,不能盲目自旋(那样会把 CPU 烧干)。LMAX 根据不同场景提供了 等待策略:
- BusySpinWaitStrategy(极致低延迟):疯狂循环,不退让。代价:烧 CPU,只适合核心独占的场景。
- YieldingWaitStrategy(推荐):自旋 100 次后,调用
std::this_thread::yield()(让出时间片)。 - SleepingWaitStrategy(省电):自旋 +
yield+ 最终sleep。适合日志、非实时系统。
6. 硬件视角:为什么 CAS 在 MPMC 中不是瓶颈?
很多人以为 Disruptor 快是因为“无锁”,其实是**“无竞争”**。
- 在 Disruptor 中,只有一个
m_producer_cursor被所有生产者 CAS 修改。 - 关键优化:
fetch_add是在同一个缓存行上累加,相比compare_exchange_weak(包含失败重试),fetch_add更高效,因为它不包含复杂的“比较-交换-失败重试”循环。
总结脑图(供笔记)
- SPSC:零 CAS,仅屏障,极限吞吐。
- MPSC:单 CAS + 批量提交,生产端收束。
- MPMC:环形缓冲 + 序号分配(Disruptor)。
- 规避竞争:生产者和消费者只修改各自的 Cursor,减少“乒乓”。
- 数据隔离:缓存行填充,彻底消除伪共享。
- 结论:你无法消灭 CAS,但可以将 CAS 频率降到 O(1),并将 CAS 的缓存行与数据彻底分离。