尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

手写线程安全消息队列:条件变量与生产者消费者模型实战

发布时间:2026/9/9 15:12:30

资讯中心
01
ARTICLE

手写线程安全消息队列:条件变量与生产者消费者模型实战

手写线程安全消息队列:条件变量与生产者消费者模型实战
1. 从裸奔的共享变量到消息队列先看一个真实案例我接手过一个采集程序最初版本用的是最直白的做法一个生产者线程不断往全局链表里塞数据消费者线程定时醒来遍历链表为了防冲突给链表挂了一把大锁。前两个月一切正常后来数据量上来问题陆续暴露消费者为了不漏数据只能高频轮询CPU空转严重生产者偶尔拿到锁却发现链表被消费者清空了数据没丢但时延忽高忽低最头疼的是业务逻辑后来要求同一批数据必须一起处理链表锁完全没法表达这种批量语义。后面被迫重构把共享变量锁这种裸奔方式换成了消息队列。这里说的消息队列不是中间件那种跨进程的MQ而是线程内部、进程地址空间里的消息传递机制。改完以后最直观的变化生产者只需要入队消费者只需要出队两边完全不感知对方的存在同步、解耦、削峰一次解决。Linux下实现线程间消息队列的常见途径就两条一条是内核提供的POSIX消息队列mq_open/mq_send/mq_receive另一条是自己用互斥锁条件变量在用户态实现。本篇标题带着(1)我计划从实用角度出发先用mutex condition_variable手写一个够用、可扩展的线程消息队列把核心机制吃透后面再展开批量消息、优先级、超时、无锁等进阶话题。这个系列定位是实用功能代码集所以不会去贴几千行的大框架而是给可以直接搬进项目的核心代码和相关注意事项。适合谁看如果你正在写多线程采集、日志异步落盘、任务分发这类程序或者面试前想踏踏实实搞懂条件变量和生产者消费者模型这篇值得花二十分钟读完代码可以直接抄进你的工程改造。2. 为什么不自造轮子也要先搞懂这套组合拳2.1 说在前面为什么不直接用std::queue配锁一段烂大街的代码如下面的写法不少第一次写多线程的人第一反应都是这样std::queueint q; std::mutex mtx; void producer() { std::lock_guardstd::mutex lock(mtx); q.push(1); } void consumer() { std::lock_guardstd::mutex lock(mtx); if (!q.empty()) { int v q.front(); q.pop(); // 处理v } }问题马上来了。消费者这端当队列为空时它什么都做不了只能反复lock、检查、unlock这就是空转。更隐蔽的问题是代码根本没处理生产者唤醒消费者这件事。如果消费者在q.empty()的瞬间被切走生产者又push了新数据消费者醒来后可能一直等不到下一次调度——虽然现实中因为锁的存在这种情况的概率不算高但谁都不想靠运气写并发代码。条件变量condition_variable就是用来解决等待通知这两个核心问题的。生产者入队后notify消费者在条件不满足时把自己挂起等通知来了再醒来既避免了空转又避免了唤醒丢失——当然前提是使用姿势要对。2.2 一个容易栽的细节unique_lock和lock_guard的取舍条件变量的等待操作必须搭配unique_lock这是C标准写死的。原因很朴素wait()内部需要原子的完成把线程阻塞和释放互斥锁两个动作等线程被唤醒后再自动重新持有锁。lock_guard的锁粒度是构造析构固定的没法在中间释放再获取所以只能由unique_lock来承担。这个点属于解释过一万遍但还是有人踩的基础知识如果你之前只写过lock_guard需要先转换一下思维。3. 手写一个线程安全消息队列核心实现与逐步拆解3.1 接口设计的几个关键取舍我实现的这个队列基础版本只提供入队、出队、停止三个语义没做批量接口刻意保持小。下面的结构体是整个队列的核心template typename T class ThreadSafeQueue { public: explicit ThreadSafeQueue(size_t capacity) : capacity_(capacity), stopped_(false) {} bool push(T value); bool pop(T value); void stop(); size_t size() const; private: mutable std::mutex mtx_; std::condition_variable not_empty_cv_; std::condition_variable not_full_cv_; std::queueT queue_; size_t capacity_; bool stopped_; };几个接口为什么要这样设计入队用右值引用T是为了减少复制。如果队列里装的是大对象比如带缓冲区的数据包每次push都深拷贝一次开销非常明显后面在优化部分我会专门测试带move和不带move的差别。pop的语义是阻塞直到拿到一个元素或者队列被停止。返回值用boolfalse代表队列已经停止且没有残留数据调用方可以借此退出循环。停止机制是这个队列和教科书版本最大的差异点很多线上事故都源于消费者线程无法优雅退出我单独用一节讲。3.2 入队实现notify放在锁内还是锁外bool push(T value) { std::unique_lockstd::mutex lock(mtx_); not_full_cv_.wait(lock, [this]() { return stopped_ || queue_.size() capacity_; }); if (stopped_) { return false; } queue_.emplace(std::forwardT(value)); not_empty_cv_.notify_one(); return true; }wait(lock, predicate)是条件变量的标准写法它等价于while (!predicate()) { cv.wait(lock); }重点是自我唤醒后要重新检查条件也就是所谓的伪唤醒spurious wakeup防护。用if判断条件的写法在极端情况下会出问题虽然概率小但多线程程序出事就是事故所以一律用带谓词的wait重载。关于notify_one是放在锁内还是解锁后再调用C标准没有强制要求但实测在锁内notify有个微小的问题被唤醒的线程要等当前线程释放锁才能继续多了一层无谓的调度延迟。如果是高吞吐场景建议这样写bool push(T value) { { std::unique_lockstd::mutex lock(mtx_); not_full_cv_.wait(lock, [this]() { return stopped_ || queue_.size() capacity_; }); if (stopped_) { return false; } queue_.emplace(std::forwardT(value)); } not_empty_cv_.notify_one(); return true; }这段是经典优化后的版本。先加锁、等条件、入队然后立刻解锁解锁后再通知消费者。这样消费者被唤醒时锁已经释放了可以直接进入临界区延迟更低。需要注意的坑是必须保证notify时队列状态确实变化了不能想着我统一在函数末尾notify一次如果wait条件不满足早退后续就会漏通知造成生产者/消费者互相傻等。3.3 弹出实现pop的阻塞与超时语义bool pop(T value) { std::unique_lockstd::mutex lock(mtx_); not_empty_cv_.wait(lock, [this]() { return stopped_ || !queue_.empty(); }); if (stopped_ queue_.empty()) { return false; } value std::move(queue_.front()); queue_.pop(); not_full_cv_.notify_one(); return true; }这里有两个细节很多人会忽视。一个是pop返回值的判定必须同时判断stopped_和queue_.empty()。因为停止信号发出时队列里可能还有残留数据正确的语义是先把存量数据消费完再退出所以stop()之后生产者虽然被拦住了但消费者的循环仍然能取完最后的元素取完下次进入pop时queue_空了才返回false。另一个是pop里notify的时机。为什么出队后要notify生产者因为队列容量是有限的消费者拿走一个元素生产者可能正卡在not_full等待上必须通知它有空位了否则队列满了以后生产者和消费者都会卡死。这个点反映的是生产者消费者模型里两类条件变量各自独立、各管各的。如果想支持非阻塞式弹出可以加一个tryPopbool tryPop(T value) { std::lock_guardstd::mutex lock(mtx_); if (queue_.empty()) { return false; } value std::move(queue_.front()); queue_.pop(); not_full_cv_.notify_one(); return true; }使用场景是消费者线程除了消费消息还需要定期做其他事情比如心跳不能无限期阻塞在pop上那就用tryPop轮询配合sleep或者后面我们再扩展带超时的pop版本。3.4 stop和析构最容易被忽略的停止机制void stop() { { std::lock_guardstd::mutex lock(mtx_); stopped_ true; } not_empty_cv_.notify_all(); not_full_cv_.notify_all(); }朴素版本常常漏了stop方面的工作。没有stop线程退出就只能靠发一个特殊消息或者干脆detach让线程自生自灭前者污染业务逻辑后者在程序退出时容易崩溃。stop加了两道保险把stopped_置为true然后同时唤醒两个条件的等待者。注意这里是notify_all而不是notify_one因为可能同时有多个生产者和多个消费者在等待只唤醒一个会导致其他人无法退出。析构函数里也需要调用stop()防止派生类先析构了条件变量消费者还挂在上面~ThreadSafeQueue() { stop(); }这个类析构时thread可能还在pop里wait如果不通知程序会直接卡在析构处——实际调试时的表现就是主线程退出时hang住用gdb看到一堆线程在__pthread_cond_wait里。4. 生产消费者场景的完整构建与实测验证4.1 一个能跑的完整示例代码只看不用是学不会的下面给出一个可以直接编译运行的示例包含两个生产者线程、两个消费者线程以及主线程在3秒后触发停止#include condition_variable #include cstdio #include mutex #include queue #include thread #include chrono #include vector template typename T class ThreadSafeQueue { public: explicit ThreadSafeQueue(size_t capacity) : capacity_(capacity), stopped_(false) {} bool push(T value) { { std::unique_lockstd::mutex lock(mtx_); not_full_cv_.wait(lock, [this]() { return stopped_ || queue_.size() capacity_; }); if (stopped_) { return false; } queue_.emplace(std::forwardT(value)); } not_empty_cv_.notify_one(); return true; } bool pop(T value) { std::unique_lockstd::mutex lock(mtx_); not_empty_cv_.wait(lock, [this]() { return stopped_ || !queue_.empty(); }); if (stopped_ queue_.empty()) { return false; } value std::move(queue_.front()); queue_.pop(); not_full_cv_.notify_one(); return true; } void stop() { { std::lock_guardstd::mutex lock(mtx_); stopped_ true; } not_empty_cv_.notify_all(); not_full_cv_.notify_all(); } size_t size() const { std::lock_guardstd::mutex lock(mtx_); return queue_.size(); } private: mutable std::mutex mtx_; std::condition_variable not_empty_cv_; std::condition_variable not_full_cv_; std::queueT queue_; size_t capacity_; bool stopped_; }; void producer(ThreadSafeQueueint q, int id) { for (int i 0; i 10; i) { bool ok q.push(i id * 100); if (ok) { printf(producer[%d] push %d, queue size%zu\n, id, i id * 100, q.size()); } else { printf(producer[%d] stopped, give up\n, id); break; } std::this_thread::sleep_for(std::chrono::milliseconds(30)); } } void consumer(ThreadSafeQueueint q, int id) { while (true) { int value 0; bool ok q.pop(value); if (!ok) { printf(consumer[%d] queue stopped, exit\n, id); break; } printf(consumer[%d] pop %d, queue size%zu\n, id, value, q.size()); std::this_thread::sleep_for(std::chrono::milliseconds(50)); } } int main() { ThreadSafeQueueint q(5); std::vectorstd::thread producers; std::vectorstd::thread consumers; for (int i 0; i 2; i) { producers.emplace_back(producer, std::ref(q), i 1); consumers.emplace_back(consumer, std::ref(q), i 1); } std::this_thread::sleep_for(std::chrono::seconds(3)); q.stop(); for (auto t : producers) t.join(); for (auto t : consumers) t.join(); printf(main exit, final queue size%zu\n, q.size()); return 0; }这个示例的关键点在于验证两个行为第一消费者不会空转测试时可以在top里看到消费线程处于睡眠状态CPU占用率几乎为零第二主线程调用stop后所有生产者消费者线程都能在短时间内退出程序能顺利跑完不会卡死。提示编译时务必加 -pthreadg 11 以上直接g -stdc11 -O2 -pthread test.cpp -o test即可。4.2 实测中需要盯的三个观察点跑这个程序我建议不要把目光只放在能跑上要盯三个细节第一个是队列长度变化。把printf打好能看到q.size()在0到5之间波动一旦到达5生产者就会阻塞在not_full上消费者的速度决定了能否把队列填满这直观体现了流量控制的效果。如果去掉capacity限制生产速度远大于消费时内存会无上限增长这是分布式采集程序最常见的隐患之一。第二个是消费顺序。生产者的i取值是0-9消费者的打印顺序并不严格按照生产者入队顺序排列因为两个生产者是并发的谁先抢到锁不一定。如果你的业务要求严格的全局顺序这个队列就不适用了需要用带序号的消息或者单一生产者模型这是消息队列选型时最容易忽略的点。第三个是stop后队列里的残留数据。我在测试时故意让生产者比消费者快得多程序退出前队列里还剩几条数据。观察到的行为是stop后生产者们立即退出但消费者们仍会继续pop直到队列被清空后才返回false退出。这正是我们期待的语义——先消费完存量再干干净净地终止。5. 容量、唤醒与线程安全的三个细节补课5.1 为什么需要两个条件变量而不是一个很多初版实现只用一个条件变量然后无论入队还是出队都notify这一个能工作但效率很差。考虑一种情况队列满了所有生产者都在not_full上等待此时一个消费者拿走一个元素它呼叫notify本意是通知生产者有空位了但用单一条件变量时消费者自己也在这个变量上等待notify完全可能唤醒另一个消费者而消费者发现队列是空的或者数据被其他消费者抢走了只能再次睡过去。这说明单一条件变量会造成错误的唤醒和惊群问题。用两个独立条件变量生产者等not_full消费者等not_empty各等各的唤醒才能精确送达。这一设计是教科书级别的在实际多线程项目里尤其重要因为它直接决定了高并发下锁竞争的激烈程度。5.2 队列容量和真满/假满的判断capacity语义上代表着内存占用的上限。在push的wait里条件是queue_.size() capacity_注意这里用的是小于也就是说容量为5时队列最多能放5个元素。想要允许队列装到6个判断就得改成这种边界条件一旦写错配合生产者消费者速度差异往往在压力测试时才暴露。另外queue_.size()返回的是size_t类型是无符号数和capacity_比较时类型不一致但两者都是非负的不存在负数比较的坑。不过如果后续扩展成支持无限容量这里条件会变成永真式就失去控流的意义了。生产环境里我一般建议初始容量保守一点比如处理网络包的线程队列设成几千个处理磁盘IO的设成几百个再根据实际q.size()的长期均值调整。5.3 多消费者场景下的惊群与负载均衡上面代码用的 multiple consumers 模式下每次pop只会唤醒一个消费者天然避免了惊群。但要注意notify_one选择唤醒哪个线程由操作系统调度决定不能保证唤醒的一定是等待最久的线程。如果希望多个消费者更公平地分摊任务需要用更复杂的调度策略比如每个消费者维护一个专属子队列生产者按某种策略轮询或负载最小的方式选择队列投递。这种设计在单队列模式下是做不到的。另外pop操作本身的移动赋值value std::move(queue_.front())意味着队列里存储的对象必须是可移动构造/可移动赋值的。如果你的T类型是const成员或者禁用了move这里直接编译不过。C11之后大多数标准库类型string、vector、shared_ptr都是可移动的所以实际影响不大但如果存放自定义对象需要确认一下。6. 从实测到工程锁粒度、move语义与更优的消息传递方案6.1 push一个对象要复制多少次聊聊move语义很多初学者写push时会把参数写成const T然后调用q.push(item)看起来没毛病但里外里会多出好几次拷贝。以std::string为例std::string msg hello; q.push(msg); // msg被拷贝进临时对象再移动进队列 q.push(std::move(msg)); // 直接把msg的资源转移进队列msg变成空串 q.push(hello std::to_string(id)); // 临时对象直接被移动零拷贝我的push签名用了T配合std::forward保证了对象以移动的方式进入容器。实测对比一个包含10KB缓冲区的自定义结构普通拷贝版每次push要拷贝10KB内存改用move后只需要拷贝几十字节的指针和长度性能差距在几十倍以上。这对网络数据包、日志行这类大对象尤其重要。有个使用上的细节调用q.push(std::move(msg))之后msg就处于有效但未指定的状态不能再依赖它的内容。如果后续还需要使用msg建议先复制一份再move或者干脆传const引用交给队列内部拷贝。这是move语义的天生代价谈不上坑但新手容易踩。6.2 批量接口与锁粒度优化一个预告性思路单条push/pop在高吞吐下最明显的瓶颈是锁竞争。每一笔消息都要经历加锁-等待条件-入队/出队-解锁-通知的完整流程当生产者和消费者数量都很多时锁的竞争会非常激烈。一种优化思路是把锁粒度放大用批量来摊薄锁开销template typename InputIt bool pushBulk(InputIt first, InputIt last);接口实现时一次性拿锁把一批元素全部入队然后只做一次notify。消费者端的对应优化是批量弹出也就是一次取走队列中所有可用的元素处理完这一批后再回到wait状态。这里有个收益点批量操作让同一批消息的处理共享了一次锁的开销同时队列的吞吐量会因为锁竞争减少而显著提升这个方向我计划在(2)里专门展开。6.3 自旋锁和条件变量怎么选不是所有场景都适合用条件变量。如果队列里消息到达的频率极高每个元素之间的间隔只有几十纳秒条件变量上下文切换的开销反而会成为瓶颈因为在互斥量和条件变量上进行线程阻塞唤醒的代价大约在几微秒级别。这种情况下可以把阻塞等待换成自旋等待——忙等一个原子变量。Linux下常用的自旋锁方案是用std::atomic_flag或者干脆用while (flag.test_and_set(std::memory_order_acquire));。自旋锁适合临界区极短且锁持有时间远小于线程切换时间的场景但如果等待时间长了CPU空转损耗同样可怕。实际工程选择时我个人的经验法则是消息到达平均间隔少于10微秒的值得考虑自旋锁更多的情况用互斥锁条件变量锁定胜局因为它在等待时不消耗CPU。7. 排查卡死问题一个完整的实战排查链路7.1 现象程序卡住不动了怎么定位某天你的程序突然不干活了top里进程还在但业务停止推进日志也不打印。结合消息队列的使用场景我倾向于按下面的顺序排查先看CPU。如果进程占用率在多个线程间跳来跳去说明线程在空转很可能进入了忙等循环检查有没有把阻塞写法误写成while(queue.empty());这种自旋。如果CPU很低线程基本都睡着了那就是真的阻塞了大概率所有线程都挂在了某个wait上。再用gdb挂上去gdb -p pid thread apply all bt如果多个线程的调用栈同时卡在__pthread_cond_wait或者std::condition_variable::wait说明大家都在等通知那问题就变成谁该通知但没通知。我见过三种常见原因生产者线程崩了或者提前退出了没人再入队消费者等空了。某个生产者线程长期持有mutex不释放比如在持锁期间调用了阻塞IO其他线程全在mutex上排队。条件变量的notify和wait之间出现了竞争用法不对例如wait前先检查了一次条件但没在锁内检查。其中第三种最隐蔽很多人会在wait外面先写if (queue_.empty())再进入wait这在多线程下是经典的竞态条件。正确做法就是上面代码那样把条件判断完全交给wait的谓词处理。7.2 队列泄漏consumer少了怎么办检查程序里new出来的消息是否只入队不出队malloc的缓冲是否只分配不释放这种队列泄漏和内存泄漏很像只是泄漏的是逻辑消息而不是内存。最直观的验证方法程序跑一段时间后看q.size()是否持续增长不回落。如果一直增长要么生产者速度远超消费者容量不够引起的内存增长要么消费者逻辑里遇到某些消息不处理就丢弃了导致队列里的对象永远得不到释放。对于前者可以把队列容量调小让生产者在队列满时阻塞更久观察消费速度能不能跟上。对于后者需要给消息加一个累计计数器消费端每处理一条就自增和生产端入队数对比差多少一目了然。7.3 一个亲测有效的验证手法打点记录唤醒次数排查完之后怎么确信问题真的修好了除了跑通业务我建议给队列加一对调试计数器一个统计pop成功次数一个统计push成功次数并配合时间戳打点打印。写一个小脚本每100毫秒cat一次/proc/ /status里的voluntary_ctxt_switches值对比正常和异常场景下的线程切换频率。手动模拟停住场景逐步回放排查链路比盲目改代码高效得多。8. 再说几个绕不开的使用注意点用完这个队列有一些边角问题我实际踩过顺手汇总在这里能帮你少走弯路。第一pop里边执行value std::move(queue_.front())对T类型是有要求的。如果T是std::unique_ptr这种只移动类型它的默认拷贝构造是被删除的直接写在代码里能编过但万一某个分支触发了拷贝路径比如value是const的编译直接报错。我给生产环境写的模板参数一律要求T可移动而且显式用static_assert(std::is_move_constructibleT::value)卡住把错误尽早在编译期暴露出来。第二用std::ref(q)把队列传给线程时生命周期管理一定要小心。让线程持有队列的引用或裸指针调用方必须保证队列存活时间超过所有线程。我的习惯是把队列和线程一起打包进一个管理类管理类析构时先stop再join顺序不能反一旦先join消费者线程再stop消费者会永远挂在pop上先stop再join则丝滑退出。第三停止后不应该允许继续push。代码里stopped_一旦为truepush即使等到了not_full也会返回false这是有意的设计。但调用方要记得响应这个false不然可能产生想退出退不掉的循环重试。我通常在生产者的循环里加一个if (!q.push(...)) break;双方配合才能干净结束。第四条件变量的谓词尽量不要做耗时操作。我看到有人的谓词里直接调了queue_.size() 1000000这种比较本身没问题如果有人把size()实现成O(n)的遍历那每次wait都会反复计算性能会比较难看。C标准里queue的size()不同实现有差异不过常见的std::queue基于dequesize()是O(1)。即便如此也不建议在谓词里做业务级计算比如统计所有消息的大小总和这个动作放在消费端做更合理。9. 进步路线这个队列后面还能怎么扩展本系列标题是(1)说明后面还有计划。我梳理一下这个线程消息队列可扩展的几个方向有些是我自己项目里打磨过的有些是看到的业界方案供你参考一是批量弹出接口。一次取走队列中全部或最多N个元素把锁竞争的开销摊薄到每个元素上这在高吞吐采集场景收益很大。实现上注意返回一个std::vector 移动整个vector出去比逐个move再push高效得多但要小心弹出后通知生产者的次数建议只notify一次。二是超时接口。pop支持wait_for或wait_until这样消费者可以最多等100毫秒没数据就去干点别的或等到某个绝对时间点再回来。实现逻辑是把wait(lock, pred)换成wait_for(lock, timeout, pred)注意条件变量和超时版本之间是谓词满足或超时任一即可返回的关系别把两个逻辑搞混。三是多队列路由。消息队列从一个队列一对消费者变成多个队列按业务类型分发。这背后就是Linux网络协议栈里per-CPU队列思想的线程版只是把per-CPU换成了per-worker或per-priority。核心收益是减少了锁竞争代价是分发逻辑复杂了。对于日志系统按级别分开处理、对于网络包按连接分流的场景都很值得做。四是无锁队列。boost::lockfree::queue和moodycamel::ConcurrentQueue是Linux下常见的两个无锁实现。无锁队列不是万能药它牺牲了有界等待的确定性换来极低的锁开销。我实测过moodycamel的ConcurrentQueue在单生产者单消费者场景下吞吐比mutex条件变量高出一个数量级但它的内存模型和ABA问题的处理更复杂编码稍有不慎就把自己坑了。作为系列第一篇主要把mutex条件变量这套基础打牢。后续我会按计划更新批量接口、性能测试对比、以及更高阶的无锁方案。如果你对哪个方向更感兴趣可以按照自己的需要先自行扩展期待在评论区看到你的实现思路和踩坑记录。最后分享一个我自己贴了三年多的调试习惯给每个线程起好名字用pthread_setname_np或者在类里维护一个name_成员打印日志时带上线程名和队列残量。线上查问题时一眼扫日志就能看出是哪个线程卡住、队列是在堆积还是被清空比重新拉gdb要快得多。习惯虽小排查效率至少提升一个档次。
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。