从一把大锁到无锁队列: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_acquirememory_order_release。我们从一个实际场景开始。

假设你正在写一个机器人感知系统,有两个线程:

线程职责特点
雷达驱动线程从网口/SDK 接收点云帧数据来得快,不能被卡住
感知算法线程对点云做滤波、聚类、检测计算重,耗时有波动

雷达每 100ms 产生一帧点云,也就是 10Hz

算法正常处理一帧要 80ms

这时系统看起来很稳:

雷达:每 100ms 产生一帧
算法:每 80ms 处理一帧

消费速度 > 生产速度

但真实系统不会永远这么规整。某一帧点云特别复杂,算法突然跑了 150ms

雷达:每 100ms 产生一帧
算法:某一帧处理 150ms

短时间内生产速度 > 消费速度

这时候有两个选择:

  1. 让雷达驱动线程等算法线程处理完;
  2. 让雷达驱动线程继续接收,把数据先放进一个中间缓冲区。

第一种看着简单,实际很危险。

驱动线程一旦被算法拖住,可能发生:

  • socket 缓冲区没人及时读,开始丢 UDP 包;
  • 驱动 SDK 内部缓冲区堆满,开始丢帧;
  • 时间戳越来越旧,后续融合模块拿到过期数据;
  • 控制链路出现不可预测延迟抖动。

在高频数据系统里,接收线程的第一职责不是处理数据,而是:

尽快把数据接住,然后立刻回去接下一份数据。

所以我们需要一个中间通道:

雷达驱动线程  --->  中间队列  --->  感知算法线程

驱动线程只负责把点云塞进队列。

算法线程自己从队列里取数据处理。

这个中间队列,就是本文的主角:SPSC 队列


1. SPSC 是什么:单生产者单消费者模型

SPSC 是一个缩写:

SPSC = Single Producer Single Consumer

翻译过来就是:

单生产者,单消费者

它描述的是一种非常具体的并发队列使用方式:

角色数量允许做什么
Producer1 个调用 push 写入数据
Consumer1 个调用 pop 读取数据

注意,这不是建议,而是铁律。

SPSC 队列只能有一个线程调用 push,只能有一个线程调用 pop

正确模型:

Producer Thread
      |
      | push
      v
+-------------+
| SPSC Queue  |
+-------------+
      |
      | pop
      v
Consumer Thread

典型场景:

场景生产者消费者
雷达点云雷达驱动线程点云算法线程
摄像头图像图像采集线程编码/识别线程
CAN 报文CAN 接收线程控制逻辑线程
音频采集音频采集线程音频处理线程
视频流视频采集线程视频编码线程

SPSC 解决的核心问题是:

一个固定的数据源,把高频数据交给另一个固定处理线程。

这里有两个关键限制:

  1. 一个固定的数据源;
  2. 一个固定的处理线程。

只要这两个条件成立,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 的职责很单纯:

  1. 看队列有没有满;
  2. 没满就把数据写入 buffer_
  3. 更新 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 的职责也很单纯:

  1. 看队列有没有空;
  2. 没空就从 buffer_ 读数据;
  3. 更新 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,通常能写出正确代码。但它有两个问题:

  1. 可能带来额外性能开销;
  2. 掩盖你真正需要的同步关系。

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)观察伪共享影响

注意几点:

  1. 编译要开优化,例如 -O2
  2. 数据量要足够大;
  3. 不要在核心循环里大量 std::cout,IO 会淹没队列开销;
  4. 测多次,看趋势,不要迷信单次结果;
  5. 不同 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 和平台测试。

但这些都是下一层优化。

先把本文这版吃透,才是真正进入无锁编程的起点。

更多推荐