前言
先把"消息队列"这个词的身份说清楚,因为它在 C++ 语境下至少指两种完全不同的东西。一种是进程间通信(IPC)的消息队列——System V 的msgget/msgsnd/msgrcv,或者 POSIX 的mq_open/mq_send/mq_receive,它们由内核维护,能在不同进程之间传递消息。另一种是进程内、线程之间的消息队列,也就是一个"线程安全的阻塞队列",通常用std::mutex加std::condition_variable实现。本文的标题是"多线程实现",所以主角是第二种;第四部分会把第一种作为对照简单带过,避免概念混淆。
第二个误区更危险:很多人以为"线程间传标志"可以用volatile搞定。这是错的。volatile只告诉编译器"这个对象的读写不要被优化掉",它不提供原子性,也不建立 happens-before 关系,更不会阻止 CPU 乱序执行。用volatile做同步,程序里就存在数据竞争,属于未定义行为(UB),标准不保证任何结果。
本文给出一个可直接编译的有界阻塞队列(bounded blocking queue),讲清条件变量的正确用法、关闭语义与线程收尾顺序,并对比无锁队列和 IPC 消息队列的适用场景。
一、条件变量:为什么必须配谓词循环
没有条件变量时,消费者只能忙等:
// ❌ 忙等:把一整个核心烧在空转上,还挡住了别的线程 while (queue.empty()) { /* spin */ }std::condition_variable::wait解决的正是这件事:它原子地"解锁 + 挂起",被唤醒时再重新加锁,等待期间不占 CPU。它的两个重载是:
void wait(std::unique_lock<std::mutex>& lock); template <class Predicate> void wait(std::unique_lock<std::mutex>& lock, Predicate pred);第一个是裸等待,第二个带谓词。永远用第二个,原因有两层:
- 虚假唤醒(spurious wakeup):标准明确允许
wait在没有notify的情况下返回。裸wait醒来后如果不重新检查条件,就会读到空队列。 - 唤醒不等于条件成立:即使是被
notify_one唤醒,在你重新拿到锁之前,可能有另一个消费者抢先取走了数据。醒来必须重新判断谓词。
带谓词的重载等价于while (!pred()) wait(lock);,把这两件事一起解决了。与它配对的另一半是通知必须发生在状态变更之后、且变更要在锁内完成:
// ❌ 状态变更没持锁,且先改状态再通知的顺序毫无保障 void push_bad(int v) { q_.push(v); // 数据竞争:别的线程可能同时读 q_ cv_.notify_one(); // 可能在对端进入 wait 之前就发出,通知直接丢掉 } // ✅ 持锁改状态,出了锁再通知 void push_good(int v) { { std::lock_guard<std::mutex> lk(mtx_); q_.push(v); } // 先解锁 cv_.notify_one(); // 再通知,此时等待者一定能看到新状态 }notify_one还是notify_all,判断标准很简单:这次状态变化能满足几个等待者的谓词。生产一条消息只可能让一个消费者满足"非空",所以用notify_one;而"关闭队列"这件事让所有等待者都要醒来退出,必须用notify_all。
二、有界阻塞队列的完整实现
有了上面的规则,队列的骨架就很直白了。要点是:有界(避免生产过快打爆内存)、双条件变量("非空"和"非满"各自一个,避免notify_one叫错人)、关闭语义(让消费者能优雅退出)。
// bounded_queue.h #pragma once #include <condition_variable> #include <cstddef> #include <mutex> #include <optional> // 需要 C++17 #include <queue> #include <utility> // 有界阻塞队列:多生产者多消费者共用;关闭后 pop 会排空并返回空 template <typename T> class BoundedQueue { public: explicit BoundedQueue(std::size_t capacity) : cap_(capacity) {} // 队列持有锁与条件变量,禁止拷贝 BoundedQueue(const BoundedQueue&) = delete; BoundedQueue& operator=(const BoundedQueue&) = delete; // 阻塞入队。队列已关闭时返回 false,调用方应停止生产。 bool push(T value) { std::unique_lock<std::mutex> lk(mtx_); not_full_.wait(lk, [this] { return closed_ || q_.size() < cap_; }); if (closed_) return false; q_.push(std::move(value)); lk.unlock(); // 先解锁再通知,减少被唤醒者的无效等待 not_empty_.notify_one(); return true; } // 阻塞出队。返回空(std::nullopt)表示队列已关闭且已排空。 std::optional<T> pop() { std::unique_lock<std::mutex> lk(mtx_); not_empty_.wait(lk, [this] { return closed_ || !q_.empty(); }); if (q_.empty()) return std::nullopt; // 只可能是"已关闭且排空" T value = std::move(q_.front()); q_.pop(); lk.unlock(); not_full_.notify_one(); return value; } // 关闭:把标志置位后唤醒所有等待者,让它们各自退出 void close() { { std::lock_guard<std::mutex> lk(mtx_); closed_ = true; } not_empty_.notify_all(); not_full_.notify_all(); } std::size_t size() const { std::lock_guard<std::mutex> lk(mtx_); return q_.size(); } private: const std::size_t cap_; mutable std::mutex mtx_; // mutable:size() 是 const 成员 std::condition_variable not_empty_; // "非空"谓词专用 std::condition_variable not_full_; // "非满"谓词专用 std::queue<T> q_; bool closed_ = false; // 受 mtx_ 保护,绝不能裸读 };每一行都值得推敲,这里挑三处说明。第一,closed_是普通bool而不是std::atomic<bool>——因为它的所有读写都在mtx_保护之下,加上原子性纯属多余,反而容易让人误以为"某些地方可以不持锁访问"。第二,mutable std::mutex mtx_中的mutable是必要的:size()是 const 成员函数,而加锁会调用mtx_的非 const 成员lock()。第三,两个条件变量用的是两个不同的等待谓词,这一点在下一节会展开。
用std::mutex的时候不需要任何std::atomic,也不需要手写内存序:锁的获取操作天然是 acquire 语义、释放是 release 语义,临界区内的读写与外部的读写之间建立了完整的 happens-before 关系。只有当你真的去掉锁去做无锁结构时,才必须自己用release/acquire建立这层关系。
三、生产者/消费者示例与收尾顺序
// main.cpp #include "bounded_queue.h" #include <iostream> #include <optional> #include <thread> #include <vector> int main() { BoundedQueue<int> queue(8); // 有界:容量 8 constexpr int kProducers = 2; constexpr int kConsumers = 3; constexpr int kPerProducer = 100; std::vector<std::thread> producers; std::vector<std::thread> consumers; for (int p = 0; p < kProducers; ++p) { // p 按值捕获,queue 按引用捕获:引用的是主线程栈上的对象 producers.emplace_back([&queue, p] { for (int i = 0; i < kPerProducer; ++i) { const int msg = p * kPerProducer + i; if (!queue.push(msg)) break; // 队列被关闭,停止生产 } }); } for (int c = 0; c < kConsumers; ++c) { consumers.emplace_back([&queue] { // optional 的 explicit operator bool 在条件位置可以隐式转换 while (const std::optional<int> msg = queue.pop()) { if (*msg % 100 == 0) { std::cout << "got " << *msg << '\n'; } } }); } // 收尾顺序:先等生产者结束,再关闭队列,最后等消费者排空退出 for (std::thread& t : producers) t.join(); queue.close(); for (std::thread& t : consumers) t.join(); std::cout << "left " << queue.size() << '\n'; return 0; }编译:g++ -std=c++17 -O2 -Wall -Wextra -pthread main.cpp -o mq_demo。注意-pthread不能省,它不只是链接 libpthread,在 Linux 上还决定了一些与线程相关的宏定义。
这段程序里最关键的不是队列本身,而是收尾顺序。close()必须在所有生产者join之后调用:如果提前关闭,生产者会发现push返回false而丢数据;如果忘了调用close(),消费者会永远卡在wait里,join也就永远不返回。而queue这个对象是主线程栈上的,所有线程都按引用使用它,所以必须在全部join完成之后它才能离开作用域。
另外提一句std::cout:标准保证多个线程并发调用同一个流对象的格式化输出函数不会产生数据竞争,但输出内容可能交错("got 100"和"got 200"的字符混在一起)。要保证整行原子,得自己加一把输出锁。
四、另外两种"消息队列"
无锁队列。用std::atomic加 CAS 实现的无锁队列避免了锁竞争和线程挂起,代价是复杂度陡增。最容易踩的坑是ABA 问题:线程 A 读到栈顶指针 P 的下一步操作是 CAS,期间线程 B 把 P 弹出、又有一个新节点恰好分配到同一地址 P 并重新入栈;A 的 CAS 比较指针相等就成功了,但它基于的"P 的 next 没变"这个前提已经不成立,链表被破坏。常见对策是给指针打包一个版本号(tagged pointer)、使用 Hazard Pointer,或者干脆用有界数组加原子下标实现环形缓冲。若要自己写无锁结构,release/acquire的配对是硬要求:
| memory_order | 语义 | 典型用途 |
|---|---|---|
relaxed | 只保证该操作本身原子,不建立任何跨线程顺序 | 统计计数 |
acquire | 用于读。它之后的读写不会被重排到它之前 | 读到"就绪标志"后再读数据 |
release | 用于写。它之前的读写不会被重排到它之后 | 写完数据后再置"就绪标志" |
acq_rel | 用于读改写,同时具备两者的语义 | fetch_add做引用计数 |
seq_cst | 默认值,在以上基础上再加一个全局单一总顺序 | 需要最直观推理时 |
再说一次:这些只在无锁代码里才需要你操心。用std::mutex+std::condition_variable的队列,锁已经把一切顺序问题包办了。而volatile在这张表里根本没有位置——它既不原子也不排序,不能用于线程同步。
进程间消息队列。跨进程传递消息要用内核提供的设施。System V 那一套是msgget建队列、msgsnd发送、msgrcv接收、msgctl控制,头文件<sys/msg.h>,消息缓冲区必须以long mtype开头(接收时可以按类型筛选)。POSIX 那一套是mq_open、mq_send、mq_receive、mq_close、mq_unlink,头文件<mqueue.h>,句柄类型是mqd_t,还支持优先级(mq_send的最后一个参数),老版本的 glibc 需要额外链接-lrt。它们的开销远大于进程内队列——每次收发都是一次系统调用加一次数据拷贝,所以即使在同一进程内,也不要用 IPC 队列代替线程队列。最后要说清:C++ 标准库至今没有任何消息队列设施,需要跨进程队列的话,Boost.Interprocess 提供了boost::interprocess::message_queue。
常见坑点
| # | 场景 | ❌ 错误做法 | ✅ 正确写法 |
|---|---|---|---|
| 1 | 用volatile做线程同步 | volatile bool ready;然后一个线程写、一个线程读 | 用std::mutex保护,或用std::atomic<bool>配release/acquire |
| 2 | 裸等待 | cv_.wait(lk);醒来直接读队列 | cv_.wait(lk, pred);带谓词,内部会循环重判 |
| 3 | 不持锁就改状态 | 先q_.push(v)再cv_.notify_one(),全程没锁 | 在锁内改状态,解锁后再通知 |
| 4 | 一个条件变量配两个谓词 | "非空"和"非满"共用一个condition_variable,notify_one唤醒了等错条件的一方 | 每个谓词配一个自己的condition_variable |
| 5 | 先查再取 | if (!q.empty()) { auto v = q.pop(); },两步之间被别人抢走 | 让pop()自己阻塞并返回可判空的结果 |
| 6 | 忘了join | std::thread对象析构时仍处于 joinable 状态 | 所有路径都要join();或从 C++20 起用std::jthread |
| 7 | 用detach逃课 | 线程detach后继续访问栈上的队列对象,主线程已经退出作用域 | 不要detach;确需后台线程就用shared_ptr管理生命周期 |
| 8 | 出队后还用旧引用 | T& r = q.front(); q.pop();之后继续读r | pop时把元素移动出来,不保留对内部存储的引用 |
第 8 条是容器通用问题,但在队列上特别容易犯:front()返回的是队列内部元素的引用,pop()之后那个元素已经析构,继续读它是悬垂引用,属于 UB。上面的实现里T value = std::move(q_.front()); q_.pop();就是这个原因——先把值搬出来,再丢弃队列里的副本。
总结
| 主题 | 结论 |
|---|---|
| 同步原语 | 进程内线程队列用std::mutex+std::condition_variable,不需要std::atomic |
| 等待方式 | 一律用带谓词的wait,醒来必然重判条件 |
| 通知时机 | 在锁内改状态,解锁后再notify;一对一用notify_one,全体退出用notify_all |
| 条件变量数量 | 一个谓词一个变量,别让两种等待共用同一个 |
volatile | 不能用于线程同步:不原子、不排序、不建立 happens-before |
| 收尾顺序 | 生产者join→close()→ 消费者join,全部结束后队列对象才可析构 |
| IPC 队列 | SysV 的msgsnd/msgrcv、POSIX 的mq_send/mq_receive是进程间的,开销远高于进程内队列 |
一句话:线程消息队列的难点从来不在数据结构,而在"什么时候改状态、什么时候通知、什么时候唤醒谁、什么时候收工"这四件事上。把这四条按上面的表做对,队列就只剩十几行样板代码了。