从一把大锁到无锁队列:C++ SPSC Ring Buffer 保姆级实战
从一把大锁到无锁队列:C++ SPSC Ring Buffer 保姆级实战
用雷达点云场景吃透
std::atomic、内存序、环形缓冲区、伪共享,以及单生产者单消费者模型的真实边界。
很多人学 C++ 并发,一开始都是这么写的:
std::mutex mtx;
std::queue<Data> q;
void push(Data data) {
std::lock_guard<std::mutex> lock(mtx);
q.push(data);
}
bool pop(Data& data) {
std::lock_guard<std::mutex> lock(mtx);
if (q.empty()) {
return false;
}
data = q.front();
q.pop();
return true;
}
这段代码没有错。
它甚至是很多业务系统里最应该优先选择的写法:简单、清楚、不容易出事故。
但如果你进入的是机器人、自动驾驶、音视频、工业控制、实时传感器这类系统,一把大锁就会开始露出问题。
雷达点云一帧一帧从网口进来,摄像头图像一张一张从 SDK 回调里吐出来,CAN 报文以毫秒级频率持续上报。接收线程最怕的不是“处理慢一点”,而是被别的线程拖住,错过了下一批数据。
这篇文章只讲一件事:
如何从业务痛点出发,完整理解并手写一个 C++ 无锁 SPSC Ring Buffer 队列。
这不是为了炫技。SPSC 是无锁编程最适合入门的一类结构:它足够简单,能把 std::atomic、内存序、缓存行这些硬核概念串起来;它也足够实用,在高频数据通道里经常出现。
0. 前言:为什么我开始研究 SPSC 无锁队列?
先别急着看 memory_order_acquire 和 memory_order_release。我们从一个实际场景开始。
假设你正在写一个机器人感知系统,有两个线程:
| 线程 | 职责 | 特点 |
|---|---|---|
| 雷达驱动线程 | 从网口/SDK 接收点云帧 | 数据来得快,不能被卡住 |
| 感知算法线程 | 对点云做滤波、聚类、检测 | 计算重,耗时有波动 |
雷达每 100ms 产生一帧点云,也就是 10Hz。
算法正常处理一帧要 80ms。
这时系统看起来很稳:
雷达:每 100ms 产生一帧
算法:每 80ms 处理一帧
消费速度 > 生产速度
但真实系统不会永远这么规整。某一帧点云特别复杂,算法突然跑了 150ms:
雷达:每 100ms 产生一帧
算法:某一帧处理 150ms
短时间内生产速度 > 消费速度
这时候有两个选择:
- 让雷达驱动线程等算法线程处理完;
- 让雷达驱动线程继续接收,把数据先放进一个中间缓冲区。
第一种看着简单,实际很危险。
驱动线程一旦被算法拖住,可能发生:
- socket 缓冲区没人及时读,开始丢 UDP 包;
- 驱动 SDK 内部缓冲区堆满,开始丢帧;
- 时间戳越来越旧,后续融合模块拿到过期数据;
- 控制链路出现不可预测延迟抖动。
在高频数据系统里,接收线程的第一职责不是处理数据,而是:
尽快把数据接住,然后立刻回去接下一份数据。
所以我们需要一个中间通道:
雷达驱动线程 ---> 中间队列 ---> 感知算法线程
驱动线程只负责把点云塞进队列。
算法线程自己从队列里取数据处理。
这个中间队列,就是本文的主角:SPSC 队列。
1. SPSC 是什么:单生产者单消费者模型
SPSC 是一个缩写:
SPSC = Single Producer Single Consumer
翻译过来就是:
单生产者,单消费者
它描述的是一种非常具体的并发队列使用方式:
| 角色 | 数量 | 允许做什么 |
|---|---|---|
| Producer | 1 个 | 调用 push 写入数据 |
| Consumer | 1 个 | 调用 pop 读取数据 |
注意,这不是建议,而是铁律。
SPSC 队列只能有一个线程调用
push,只能有一个线程调用pop。
正确模型:
Producer Thread
|
| push
v
+-------------+
| SPSC Queue |
+-------------+
|
| pop
v
Consumer Thread
典型场景:
| 场景 | 生产者 | 消费者 |
|---|---|---|
| 雷达点云 | 雷达驱动线程 | 点云算法线程 |
| 摄像头图像 | 图像采集线程 | 编码/识别线程 |
| CAN 报文 | CAN 接收线程 | 控制逻辑线程 |
| 音频采集 | 音频采集线程 | 音频处理线程 |
| 视频流 | 视频采集线程 | 视频编码线程 |
SPSC 解决的核心问题是:
一个固定的数据源,把高频数据交给另一个固定处理线程。
这里有两个关键限制:
- 一个固定的数据源;
- 一个固定的处理线程。
只要这两个条件成立,SPSC 就非常舒服。
如果多个线程都要写入同一个队列,那不是 SPSC。
如果多个线程都要从同一个队列取任务,那也不是 SPSC。
SPSC 快,是因为它把问题限制得足够小。
2. 为什么不用 mutex:一把锁背后的真实代价
初学者看到无锁队列时,第一反应通常是:
为什么不用
std::lock_guard?C++ 都给了自动锁,安全又省事。
这个问题问得对。
如果你写的是普通业务系统,std::mutex + std::queue 经常就是正确答案。比如:
- 后台任务队列;
- 低频日志缓存;
- 用户请求排队;
- 配置更新;
- 初始化阶段的数据交换。
这些场景里,队列本身通常不是性能瓶颈。用锁换简单性,完全合理。
但高频数据通道不一样。
先看一个有锁队列:
#include <mutex>
#include <optional>
#include <queue>
template <typename T>
class MutexQueue {
public:
void push(const T& item) {
// lock_guard 是 RAII 锁:
// 构造时加锁,离开作用域时自动解锁。
std::lock_guard<std::mutex> lock(mutex_);
queue_.push(item);
}
std::optional<T> pop() {
std::lock_guard<std::mutex> lock(mutex_);
if (queue_.empty()) {
return std::nullopt;
}
T item = queue_.front();
queue_.pop();
return item;
}
private:
std::mutex mutex_;
std::queue<T> queue_;
};
这段代码的优缺点很清楚:
| 维度 | 表现 |
|---|---|
| 简单性 | 很好 |
| 通用性 | 多生产者、多消费者都能用 |
| 正确性 | 锁保护整个临界区,不容易写错 |
| 高频低延迟 | 有竞争时可能抖动 |
互斥锁真正麻烦的地方,不是“加锁”这个动作本身,而是竞争时会发生什么。
当一个线程抢不到锁时,操作系统可能会把它挂起:
线程抢锁失败
-> 进入等待状态
-> 操作系统调度其他线程
-> 锁释放后再唤醒
-> 发生上下文切换
-> 线程继续执行
上下文切换不是免费的。它涉及寄存器保存、线程状态切换、调度器介入、缓存热度下降。
在普通服务端业务里,几微秒到十几微秒可能没什么感觉。
但在实时传感器接收链路里,这种延迟不是固定的,最麻烦的是不可预测。
对比一下:
| 方案 | 同步方式 | 竞争时行为 | 延迟特征 | 适合场景 |
|---|---|---|---|---|
mutex + queue | 互斥锁保护临界区 | 可能阻塞、挂起、唤醒 | 延迟不稳定 | 通用业务队列 |
| SPSC 无锁队列 | 原子指针 + 内存序 | 不阻塞线程 | 延迟更可控 | 高频一写一读通道 |
这里还有一个必须纠正的误区:
std::atomic不是“自动上锁”。
std::atomic 不是用来保护一大段复杂逻辑的。它保护的是某个变量本身的原子读写,并且配合内存序建立线程间可见性。
在 SPSC 里,我们不会用锁保护整个队列,而是把队列拆成几个非常简单的状态:
- 写指针
head_; - 读指针
tail_; - 固定数组
buffer_。
然后利用“一写一读”的角色限制,把复杂互斥问题降级成“一个线程写某个指针,另一个线程读这个指针”。
这才是 SPSC 能无锁的根本原因。
3. Ring Buffer:SPSC 队列的骨架
无锁不是第一步。
第一步是选对存储结构。
SPSC 队列常用的底层结构是 Ring Buffer,也叫环形缓冲区。
你可以把它理解成一个首尾相接的数组:
下标: 0 1 2 3 4 5 6 7
+---+---+---+---+---+---+---+---+
| | | | | | | | |
+---+---+---+---+---+---+---+---+
^ ^
tail head
它有两个核心指针:
| 指针 | 含义 | 谁推进 |
|---|---|---|
head_ | 下一个写入位置 | 生产者 |
tail_ | 下一个读取位置 | 消费者 |
写入一个元素后,head_ 往前走。
读取一个元素后,tail_ 往前走。
走到数组末尾后,通过取模回到开头:
next = (current + 1) % CAPACITY;
这就是“环形”的来源。
3.1 为什么不用无限增长的 std::queue?
实时系统不喜欢无限增长。
无限增长意味着:
- 消费者慢时内存持续上涨;
- 延迟越来越大;
- 系统没有明确背压点;
- 最后可能不是丢一帧,而是整个进程被拖死。
固定容量 Ring Buffer 的态度更硬:
我最多就存这么多。
满了就告诉你失败。
你必须设计丢帧、报警、降级策略。
这正是实时系统需要的。
队列满不是异常,它是信号:
消费者已经跟不上生产者了。
3.2 如何判断队列空?
当 head_ == tail_ 时,说明没有可读数据:
head == tail -> empty
图示:
下标: 0 1 2 3
+---+---+---+---+
| | | | |
+---+---+---+---+
^
|
head/tail
3.3 如何判断队列满?
常见做法是牺牲一个格子。
当 head_ 的下一个位置等于 tail_,说明队列满:
next_head = (head + 1) % CAPACITY;
if (next_head == tail) {
// full
}
也就是说:
next_head == tail -> full
为什么要牺牲一个格子?
因为如果不牺牲,当数组完全写满时,head_ 也可能等于 tail_。这会和空队列状态冲突。
| 状态 | 条件 |
|---|---|
| 空 | head == tail |
| 满 | (head + 1) % CAPACITY == tail |
代价是:容量为 CAPACITY 的数组,最多只能存 CAPACITY - 1 个元素。
这个代价通常值得,因为逻辑简单、不容易出错。
4. 从零定义 LockFreeSPSCQueue
现在开始写真正的队列。
先给出类骨架:
#include <atomic>
#include <cstddef>
#include <vector>
template <typename T, std::size_t CAPACITY>
class LockFreeSPSCQueue {
public:
LockFreeSPSCQueue()
: buffer_(CAPACITY), head_(0), tail_(0) {
static_assert(CAPACITY >= 2, "CAPACITY must be at least 2");
}
private:
std::vector<T> buffer_;
// head_ 是写指针:只有生产者线程会修改它,消费者线程只读取它。
alignas(64) std::atomic<std::size_t> head_;
// tail_ 是读指针:只有消费者线程会修改它,生产者线程只读取它。
alignas(64) std::atomic<std::size_t> tail_;
};
先看成员:
| 成员 | 作用 |
|---|---|
buffer_ | 固定容量的底层存储 |
head_ | 写指针,表示下一个写入位置 |
tail_ | 读指针,表示下一个读取位置 |
alignas(64) | 尽量让两个原子变量分开到不同缓存行 |
为什么 head_ 和 tail_ 是 std::atomic<std::size_t>?
因为它们会被两个线程同时访问:
| 变量 | 生产者 | 消费者 |
|---|---|---|
head_ | 写 | 读 |
tail_ | 读 | 写 |
虽然不会有两个线程同时写同一个指针,但会有一个线程写、另一个线程读。
普通变量在这种情况下会形成数据竞争,C++ 标准层面就是未定义行为。
所以必须用 std::atomic。
5. push:生产者如何写入数据
push 的职责很单纯:
- 看队列有没有满;
- 没满就把数据写入
buffer_; - 更新
head_,告诉消费者“新数据来了”。
完整代码:
bool push(const T& item) {
// 当前写位置。
// 只有生产者线程会修改 head_,所以生产者自己读取 head_ 时,
// 不需要和其他线程建立同步关系,只需要原子性即可。
std::size_t current_head = head_.load(std::memory_order_relaxed);
// 计算下一个写位置,走到末尾后回绕到 0。
std::size_t next_head = (current_head + 1) % CAPACITY;
// 生产者需要读取 tail_ 来判断队列是否满。
// tail_ 是消费者发布的读位置,这里使用 acquire,
// 确保能看到消费者对 tail_ 的释放更新。
std::size_t current_tail = tail_.load(std::memory_order_acquire);
// 如果下一个写位置追上了读位置,说明队列满。
if (next_head == current_tail) {
return false;
}
// 真正写入数据。
// 这一步必须发生在 head_.store(release) 之前。
buffer_[current_head] = item;
// 发布新的写位置。
// release 保证:上面对 buffer_ 的写入,不会被重排到这句之后。
// 消费者一旦通过 acquire 看到新的 head_,就必须能看到对应的数据。
head_.store(next_head, std::memory_order_release);
return true;
}
5.1 执行流程拆解
假设容量为 8,当前状态:
tail = 2
head = 5
说明 [2, 5) 之间有数据,head = 5 是下一个写位置。
生产者执行:
current_head = 5
next_head = 6
current_tail = 2
因为 next_head != current_tail,队列没满。
于是写入:
buffer[5] = item
head = 6
消费者看到 head 从 5 变成 6,就知道新数据已经可读。
5.2 为什么先写 buffer,再更新 head?
这一点是 SPSC 的命门。
正确顺序:
写 buffer
发布 head
如果顺序反了:
发布 head
写 buffer
消费者可能看到 head_ 已经前进,于是立刻去读 buffer_,但此时数据还没写进去。
那消费者读到的可能是旧数据,甚至是未初始化数据。
head_.store(..., std::memory_order_release) 就是告诉编译器和 CPU:
前面的写入必须先完成,才能把新的 head 发布出去。
6. pop:消费者如何读取数据
pop 的职责也很单纯:
- 看队列有没有空;
- 没空就从
buffer_读数据; - 更新
tail_,告诉生产者“这个格子我读完了”。
完整代码:
bool pop(T& item) {
// 当前读位置。
// 只有消费者线程会修改 tail_,所以消费者自己读取 tail_ 时,
// 不需要和其他线程建立同步关系,只需要原子性即可。
std::size_t current_tail = tail_.load(std::memory_order_relaxed);
// 消费者需要读取 head_ 来判断队列是否为空。
// head_ 是生产者发布的写位置,这里使用 acquire,
// 确保如果看到了新的 head_,也能看到生产者在 release 前写入的 buffer_ 数据。
std::size_t current_head = head_.load(std::memory_order_acquire);
// 读写位置相等,说明没有可读数据。
if (current_tail == current_head) {
return false;
}
// 读取数据。
// 这一步必须发生在 head_.load(acquire) 之后。
item = buffer_[current_tail];
// 计算下一个读位置。
std::size_t next_tail = (current_tail + 1) % CAPACITY;
// 发布新的读位置。
// release 保证:读取 buffer_ 的动作不会被重排到更新 tail_ 之后。
// 生产者看到新的 tail_ 后,就可以认为这个槽位已经被消费完。
tail_.store(next_tail, std::memory_order_release);
return true;
}
6.1 执行流程拆解
假设当前状态:
tail = 2
head = 6
说明 [2, 6) 之间有数据。
消费者执行:
current_tail = 2
current_head = 6
因为 current_tail != current_head,队列不空。
于是读取:
item = buffer[2]
tail = 3
生产者看到 tail 从 2 变成 3,就知道 buffer[2] 这个位置可以重新写了。
6.2 为什么读完后再更新 tail?
正确顺序:
读 buffer
发布 tail
如果顺序反了:
发布 tail
读 buffer
生产者可能看到 tail_ 已经前进,于是认为旧槽位可以复用,直接写入新数据。
消费者这时再去读,就可能读到被覆盖的数据。
所以 tail_.store(..., std::memory_order_release) 的意义是:
我真的读完了,才告诉生产者这个格子可以复用。
7. 完整版 LockFreeSPSCQueue
把前面的代码合起来:
#include <atomic>
#include <cstddef>
#include <vector>
template <typename T, std::size_t CAPACITY>
class LockFreeSPSCQueue {
public:
LockFreeSPSCQueue()
: buffer_(CAPACITY), head_(0), tail_(0) {
static_assert(CAPACITY >= 2, "CAPACITY must be at least 2");
}
bool push(const T& item) {
std::size_t current_head = head_.load(std::memory_order_relaxed);
std::size_t next_head = (current_head + 1) % CAPACITY;
std::size_t current_tail = tail_.load(std::memory_order_acquire);
if (next_head == current_tail) {
return false;
}
buffer_[current_head] = item;
head_.store(next_head, std::memory_order_release);
return true;
}
bool pop(T& item) {
std::size_t current_tail = tail_.load(std::memory_order_relaxed);
std::size_t current_head = head_.load(std::memory_order_acquire);
if (current_tail == current_head) {
return false;
}
item = buffer_[current_tail];
std::size_t next_tail = (current_tail + 1) % CAPACITY;
tail_.store(next_tail, std::memory_order_release);
return true;
}
bool empty() const {
std::size_t current_tail = tail_.load(std::memory_order_acquire);
std::size_t current_head = head_.load(std::memory_order_acquire);
return current_tail == current_head;
}
private:
std::vector<T> buffer_;
alignas(64) std::atomic<std::size_t> head_;
alignas(64) std::atomic<std::size_t> tail_;
};
避坑提示
这个版本要求
T可以默认构造,因为std::vector<T> buffer_(CAPACITY)会构造出CAPACITY个元素。工业级实现会进一步使用未初始化存储、placement new、析构管理来支持更复杂的类型。本文先聚焦 SPSC 的并发机制,不把内存管理复杂度混进来。
8. 内存序专题:为什么 acquire/release 是 SPSC 的灵魂
如果你只看代码,很容易觉得内存序是“玄学参数”:
std::memory_order_relaxed
std::memory_order_acquire
std::memory_order_release
其实这三个参数解决的是三个不同层面的问题。
| 内存序 | 直觉解释 | 在本文队列中的用途 |
|---|---|---|
relaxed | 只保证这个变量的原子读写,不建立同步关系 | 自己读自己负责修改的指针 |
release | 发布前面的写入/读取,不允许关键操作被重排到发布之后 | 更新 head_ 或 tail_ |
acquire | 获取对方发布的状态,看到状态后也能看到状态之前的数据 | 读取对方的 head_ 或 tail_ |
8.1 SPSC 里有两组同步关系
第一组:生产者发布数据给消费者。
// Producer
buffer_[current_head] = item;
head_.store(next_head, std::memory_order_release);
// Consumer
std::size_t current_head = head_.load(std::memory_order_acquire);
item = buffer_[current_tail];
含义:
生产者先写 buffer,再 release 发布 head。
消费者 acquire 读取 head 后,才能安全读取 buffer。
第二组:消费者释放槽位给生产者。
// Consumer
item = buffer_[current_tail];
tail_.store(next_tail, std::memory_order_release);
// Producer
std::size_t current_tail = tail_.load(std::memory_order_acquire);
buffer_[current_head] = item;
含义:
消费者先读完 buffer,再 release 发布 tail。
生产者 acquire 读取 tail 后,才能安全复用槽位。
8.2 为什么不能全部 relaxed?
relaxed 只保证原子变量本身读写不撕裂、不形成数据竞争。
但它不保证:
buffer_写入一定发生在head_更新之前;- 消费者看到
head_更新后一定看到对应buffer_; - 消费者读取
buffer_一定发生在tail_更新之前。
如果全部用 relaxed,你可能在 x86 上短时间看不出问题,因为 x86 的内存模型相对强。
但 C++ 代码不是只写给 x86 的。编译器优化、ARM 平台、不同优化级别,都可能把这个问题暴露出来。
底层原理补充
多线程正确性不能靠“我本机跑了没事”证明。错误内存序最麻烦的地方是低概率、难复现、依赖平台。越是高性能代码,越不能用运气换正确性。
8.3 为什么不全部 seq_cst?
默认的原子操作通常是 memory_order_seq_cst,也就是顺序一致性。
它最容易理解:所有线程看到一个全局一致的原子操作顺序。
那为什么不用它?
可以用。
初学阶段全部用 seq_cst,通常能写出正确代码。但它有两个问题:
- 可能带来额外性能开销;
- 掩盖你真正需要的同步关系。
SPSC 的同步关系很明确:
- 生产者发布
head_,消费者获取head_; - 消费者发布
tail_,生产者获取tail_。
用 acquire/release 可以精确表达这件事。
对比:
| 写法 | 正确性 | 性能 | 表达意图 |
|---|---|---|---|
全部 seq_cst | 通常正确 | 可能更重 | 不够精确 |
| 精确 acquire/release | 正确 | 更轻 | 清楚表达发布/获取关系 |
全部 relaxed | 不可靠 | 看似最快 | 错误风险高 |
9. alignas(64) 与伪共享:性能为什么会被缓存行拖垮
现在解释前面埋下的 alignas(64)。
alignas(64) std::atomic<std::size_t> head_;
alignas(64) std::atomic<std::size_t> tail_;
CPU 访问内存时,不是每次只搬一个变量。
它通常按 Cache Line 搬运数据。常见缓存行大小是 64 字节。
也就是说,哪怕你只读一个 std::size_t,CPU 也可能把它所在的整条缓存行搬到核心的缓存里。
如果 head_ 和 tail_ 紧挨着,可能落在同一条缓存行:
同一条 Cache Line:
+-------------------------------+
| head_ | tail_ | 其他数据 ... |
+-------------------------------+
生产者频繁写 head_。
消费者频繁写 tail_。
虽然它们写的是不同变量,但如果这两个变量在同一条缓存行,两个 CPU 核心就会为了这条缓存行反复做一致性同步。
这叫 False Sharing,伪共享。
它“伪”在:
- 程序逻辑上没有共享同一个变量;
- CPU 缓存层面却共享了同一条缓存行。
结果就是性能抖动。
使用 alignas(64) 的目标是让两个高频写入的原子变量尽量分开:
Cache Line A:
+-------------------------------+
| head_ |
+-------------------------------+
Cache Line B:
+-------------------------------+
| tail_ |
+-------------------------------+
这样生产者写 head_ 时,不会频繁影响消费者写 tail_ 所在缓存行。
避坑提示
alignas(64)是性能优化,不是正确性条件。去掉它,队列语义上仍然应该正确,但高频场景下性能可能明显变差。
10. 实战场景:雷达点云如何使用 SPSC
现在把队列放回真实业务。
我们写一个可编译运行的示例:
- 生产者线程模拟雷达,每
100ms产生一帧; - 消费者线程模拟算法,每帧处理
80ms; - 队列容量设置为 8;
- 如果队列满,生产者丢弃当前帧;
- 如果队列空,消费者短暂休眠,避免 CPU 空转。
完整代码:
#include <atomic>
#include <chrono>
#include <cstddef>
#include <cstdint>
#include <iostream>
#include <thread>
#include <vector>
template <typename T, std::size_t CAPACITY>
class LockFreeSPSCQueue {
public:
LockFreeSPSCQueue()
: buffer_(CAPACITY), head_(0), tail_(0) {
static_assert(CAPACITY >= 2, "CAPACITY must be at least 2");
}
bool push(const T& item) {
std::size_t current_head = head_.load(std::memory_order_relaxed);
std::size_t next_head = (current_head + 1) % CAPACITY;
std::size_t current_tail = tail_.load(std::memory_order_acquire);
if (next_head == current_tail) {
return false;
}
buffer_[current_head] = item;
head_.store(next_head, std::memory_order_release);
return true;
}
bool pop(T& item) {
std::size_t current_tail = tail_.load(std::memory_order_relaxed);
std::size_t current_head = head_.load(std::memory_order_acquire);
if (current_tail == current_head) {
return false;
}
item = buffer_[current_tail];
std::size_t next_tail = (current_tail + 1) % CAPACITY;
tail_.store(next_tail, std::memory_order_release);
return true;
}
private:
std::vector<T> buffer_;
alignas(64) std::atomic<std::size_t> head_;
alignas(64) std::atomic<std::size_t> tail_;
};
struct PointCloudFrame {
std::uint64_t timestamp_ms{};
int frame_id{};
// 为了让示例更接近“较大数据帧”,这里放一点模拟数据。
// 真实项目里通常是点云 vector、图像 buffer 或 SDK 数据结构。
float points[1024][3]{};
};
std::uint64_t now_ms() {
return static_cast<std::uint64_t>(
std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::steady_clock::now().time_since_epoch())
.count());
}
int main() {
LockFreeSPSCQueue<PointCloudFrame, 8> queue;
std::atomic<bool> running{true};
std::atomic<int> dropped_frames{0};
std::thread producer([&]() {
int frame_id = 0;
while (running.load(std::memory_order_relaxed)) {
PointCloudFrame frame;
frame.timestamp_ms = now_ms();
frame.frame_id = frame_id++;
// 模拟填充点云数据。
frame.points[0][0] = static_cast<float>(frame.frame_id);
if (queue.push(frame)) {
std::cout << "[Driver] push frame_id = "
<< frame.frame_id << '\n';
} else {
// 队列满,说明算法线程消费不过来。
// 实时系统里常见策略是丢弃当前帧,避免阻塞驱动线程。
dropped_frames.fetch_add(1, std::memory_order_relaxed);
std::cout << "[Driver] queue full, drop frame_id = "
<< frame.frame_id << '\n';
}
// 模拟雷达 10Hz。
std::this_thread::sleep_for(std::chrono::milliseconds(100));
}
});
std::thread consumer([&]() {
PointCloudFrame frame;
while (running.load(std::memory_order_relaxed)) {
if (!queue.pop(frame)) {
// 队列空时不要疯狂 while 空转。
// 这里用 1ms sleep 做一个简单退避。
std::this_thread::sleep_for(std::chrono::milliseconds(1));
continue;
}
std::cout << "[Algorithm] process frame_id = "
<< frame.frame_id << '\n';
// 模拟算法耗时。
std::this_thread::sleep_for(std::chrono::milliseconds(80));
}
});
std::this_thread::sleep_for(std::chrono::seconds(5));
running.store(false, std::memory_order_relaxed);
producer.join();
consumer.join();
std::cout << "Dropped frames = "
<< dropped_frames.load(std::memory_order_relaxed)
<< '\n';
return 0;
}
Linux 编译:
g++ -std=c++17 -O2 -pthread spsc_demo.cpp -o spsc_demo
./spsc_demo
Windows MinGW/MSYS2 编译:
g++ -std=c++17 -O2 spsc_demo.cpp -o spsc_demo.exe
spsc_demo.exe
10.1 怎么观察队列满?
把消费者算法耗时从 80ms 改成 150ms:
std::this_thread::sleep_for(std::chrono::milliseconds(150));
这时生产者每 100ms 产生一帧,消费者每 150ms 才处理一帧。
一段时间后队列会满,生产者会打印:
[Driver] queue full, drop frame_id = ...
这不是队列坏了。
这是队列在告诉你:
算法线程跟不上输入数据了。
这时候工程上要选策略:
| 策略 | 适合场景 |
|---|---|
| 丢弃当前帧 | 实时感知,更关注最新数据 |
| 覆盖最旧帧 | 只关心最近状态 |
| 报警降级 | 安全关键链路 |
| 扩大缓冲 | 短时间突刺可接受 |
| 降低生产频率 | 可控制数据源频率 |
不要幻想“队列无限大就解决了”。无限大只是把丢帧变成了延迟爆炸。
11. SPSC 的边界:为什么只能一个生产者、一个消费者
这是整篇文章最容易被误用的地方。
SPSC 不能多个生产者。
也不能多个消费者。
11.1 多个生产者会发生什么?
看 push 的核心逻辑:
std::size_t current_head = head_.load(std::memory_order_relaxed);
std::size_t next_head = (current_head + 1) % CAPACITY;
buffer_[current_head] = item;
head_.store(next_head, std::memory_order_release);
假设两个生产者线程 A 和 B 同时调用 push。
当前:
head = 5
可能发生:
线程 A 读取 current_head = 5
线程 B 读取 current_head = 5
线程 A 写 buffer[5] = A_data
线程 B 写 buffer[5] = B_data
线程 A store head = 6
线程 B store head = 6
结果:
- A 的数据被 B 覆盖;
- 两个生产者都以为自己写成功;
head_只前进了一格;- 队列状态已经错了。
std::atomic 没有自动解决这个问题。
因为这里缺的不是“原子读写”,而是多个生产者抢占写槽位的协议。
那需要 CAS,也就是 compare_exchange,这已经不是 SPSC 的设计了。
11.2 多个消费者会发生什么?
看 pop 的核心逻辑:
std::size_t current_tail = tail_.load(std::memory_order_relaxed);
item = buffer_[current_tail];
tail_.store(next_tail, std::memory_order_release);
假设两个消费者线程 C 和 D 同时调用 pop。
当前:
tail = 3
可能发生:
线程 C 读取 current_tail = 3
线程 D 读取 current_tail = 3
线程 C 读取 buffer[3]
线程 D 读取 buffer[3]
线程 C store tail = 4
线程 D store tail = 4
结果:
- 同一份数据被消费两次;
- 两个消费者都以为自己拿到了任务;
- 队列语义已经错了。
如果消费者竞争更复杂,还可能出现越空读取、顺序错乱等问题。
所以 SPSC 的铁律必须写在脑子里:
只能一个线程 push
只能一个线程 pop
如果你需要别的模型:
| 模型 | 含义 |
|---|---|
| MPSC | 多生产者,单消费者 |
| SPMC | 单生产者,多消费者 |
| MPMC | 多生产者,多消费者 |
这些不是 SPSC 稍微改改参数就能得到的。它们需要不同的并发协议。
12. Benchmark 设计:如何证明它真的更快
性能测试不能只写一句“无锁更快”。
你应该至少比较:
- 总耗时;
- 吞吐量;
- 丢帧数;
- 消费延迟;
- CPU 占用;
- 队列满/空次数;
- 去掉
alignas(64)后的变化。
一个基础 Benchmark 可以这样设计:
| 实验 | 目的 |
|---|---|
| SPSC 传递 100 万个整数 | 测无锁队列吞吐 |
mutex + queue 传递 100 万个整数 | 做基线对比 |
| 消费者故意变慢 | 观察队列满和丢帧 |
去掉 alignas(64) | 观察伪共享影响 |
注意几点:
- 编译要开优化,例如
-O2; - 数据量要足够大;
- 不要在核心循环里大量
std::cout,IO 会淹没队列开销; - 测多次,看趋势,不要迷信单次结果;
- 不同 CPU、系统负载、编译器会影响结果。
你可以先做一个整数吞吐测试,再换成点云结构体测试。
整数测试看队列本身开销。
大结构体测试更接近真实业务,但会混入拷贝成本。
如果真实项目里帧很大,更常见的做法是传递指针、智能指针、buffer 句柄,而不是每次拷贝整帧。
13. 最后的检查清单
写完一个 SPSC 队列后,别急着说自己懂无锁了。先对照下面这张表自查。
| 检查项 | 是否必须满足 |
|---|---|
是否只有一个生产者线程调用 push | 必须 |
是否只有一个消费者线程调用 pop | 必须 |
head_ / tail_ 是否是 std::atomic | 必须 |
写入 buffer_ 是否发生在发布 head_ 之前 | 必须 |
读取 buffer_ 是否发生在发布 tail_ 之前 | 必须 |
| 读取对方指针是否使用 acquire | 建议/必要同步 |
| 发布自己指针是否使用 release | 建议/必要同步 |
| 队列满是否有明确策略 | 必须 |
| 队列空是否避免无意义 CPU 空转 | 工程上必须考虑 |
| 是否误用于 MPSC/SPMC/MPMC | 绝对不能 |
14. 总结:这篇文章真正要掌握什么
SPSC 无锁队列不是魔法。
它的核心逻辑可以压缩成几句话:
Ring Buffer 负责存数据。
head_ 表示下一个写位置。
tail_ 表示下一个读位置。
生产者只写 head_。
消费者只写 tail_。
atomic 保证跨线程读写指针安全。
release/acquire 保证数据和指针发布顺序正确。
alignas(64) 减少伪共享带来的性能抖动。
再换成工程语言:
SPSC 的性能来自边界。它快,是因为它只处理一写一读。
不要把它当通用队列。
不要多个线程一起 push。
不要多个线程一起 pop。
不要用“我本机跑了没出错”来证明无锁代码正确。
真正可靠的并发代码,必须能说清楚:
- 谁写?
- 谁读?
- 哪个状态被发布?
- 哪个线程获取这个发布?
- 队列满了怎么办?
- 队列空了怎么办?
- 这个结构的适用边界在哪里?
如果这些问题都能答上来,SPSC 这块地基就算打稳了。
15. 思维导图复盘
SPSC 无锁队列
├── 背景痛点
│ ├── 高频数据接收
│ ├── mutex 阻塞
│ ├── 上下文切换
│ └── 驱动线程不能被算法拖死
├── 核心模型
│ ├── Single Producer
│ ├── Single Consumer
│ ├── push 只由生产者调用
│ └── pop 只由消费者调用
├── 数据结构
│ ├── Ring Buffer
│ ├── buffer_
│ ├── head_
│ ├── tail_
│ ├── 空判断 head == tail
│ └── 满判断 next_head == tail
├── 并发机制
│ ├── std::atomic
│ ├── relaxed
│ ├── acquire
│ ├── release
│ └── happens-before
├── 性能优化
│ ├── Cache Line
│ ├── False Sharing
│ └── alignas(64)
├── 实战使用
│ ├── LiDAR 点云
│ ├── 摄像头帧
│ ├── CAN 数据
│ └── 音视频流
└── 使用边界
├── 不能多个生产者
├── 不能多个消费者
├── 队列满要设计策略
└── 不适合复杂通用任务队列
附录:一份更适合保存的完整头文件
如果你想把队列单独保存成一个头文件,可以用下面这个版本。
文件名建议:
lock_free_spsc_queue.hpp
代码:
#pragma once
#include <atomic>
#include <cstddef>
#include <vector>
template <typename T, std::size_t CAPACITY>
class LockFreeSPSCQueue {
public:
LockFreeSPSCQueue()
: buffer_(CAPACITY), head_(0), tail_(0) {
static_assert(CAPACITY >= 2, "CAPACITY must be at least 2");
}
LockFreeSPSCQueue(const LockFreeSPSCQueue&) = delete;
LockFreeSPSCQueue& operator=(const LockFreeSPSCQueue&) = delete;
bool push(const T& item) {
const std::size_t current_head =
head_.load(std::memory_order_relaxed);
const std::size_t next_head =
increment(current_head);
const std::size_t current_tail =
tail_.load(std::memory_order_acquire);
if (next_head == current_tail) {
return false;
}
buffer_[current_head] = item;
head_.store(next_head, std::memory_order_release);
return true;
}
bool pop(T& item) {
const std::size_t current_tail =
tail_.load(std::memory_order_relaxed);
const std::size_t current_head =
head_.load(std::memory_order_acquire);
if (current_tail == current_head) {
return false;
}
item = buffer_[current_tail];
tail_.store(increment(current_tail), std::memory_order_release);
return true;
}
bool empty() const {
const std::size_t current_tail =
tail_.load(std::memory_order_acquire);
const std::size_t current_head =
head_.load(std::memory_order_acquire);
return current_tail == current_head;
}
bool full() const {
const std::size_t current_head =
head_.load(std::memory_order_acquire);
const std::size_t next_head =
increment(current_head);
const std::size_t current_tail =
tail_.load(std::memory_order_acquire);
return next_head == current_tail;
}
private:
static constexpr std::size_t increment(std::size_t value) {
return (value + 1) % CAPACITY;
}
private:
std::vector<T> buffer_;
alignas(64) std::atomic<std::size_t> head_;
alignas(64) std::atomic<std::size_t> tail_;
};
这份代码的定位是教学版和轻量实践版。
如果要走工业级,还可以继续增强:
- 支持 move-only 类型;
- 使用未初始化存储避免默认构造成本;
- 对容量为 2 的幂时用位运算替代取模;
- 添加批量 push/pop;
- 添加统计指标;
- 添加可选等待策略;
- 加入更严格的 benchmark 和平台测试。
但这些都是下一层优化。
先把本文这版吃透,才是真正进入无锁编程的起点。
更多推荐
所有评论(0)