ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

C++基于消息队列的多线程实现示例代码

C++基于消息队列的多线程实现示例代码 前言先把消息队列这个词的身份说清楚因为它在 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_lockstd::mutex lock); template class Predicate void wait(std::unique_lockstd::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_guardstd::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 // 需要 C17 #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_lockstd::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::optionalT pop() { std::unique_lockstd::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_guardstd::mutex lk(mtx_); closed_ true; } not_empty_.notify_all(); not_full_.notify_all(); } std::size_t size() const { std::lock_guardstd::mutex lk(mtx_); return q_.size(); } private: const std::size_t cap_; mutable std::mutex mtx_; // mutablesize() 是 const 成员 std::condition_variable not_empty_; // 非空谓词专用 std::condition_variable not_full_; // 非满谓词专用 std::queueT q_; bool closed_ false; // 受 mtx_ 保护绝不能裸读 };每一行都值得推敲这里挑三处说明。第一closed_是普通bool而不是std::atomicbool——因为它的所有读写都在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() { BoundedQueueint queue(8); // 有界容量 8 constexpr int kProducers 2; constexpr int kConsumers 3; constexpr int kPerProducer 100; std::vectorstd::thread producers; std::vectorstd::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::optionalint 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 -stdc17 -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::mutexstd::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::atomicbool配release/acquire2裸等待cv_.wait(lk);醒来直接读队列cv_.wait(lk, pred);带谓词内部会循环重判3不持锁就改状态先q_.push(v)再cv_.notify_one()全程没锁在锁内改状态解锁后再通知4一个条件变量配两个谓词非空和非满共用一个condition_variablenotify_one唤醒了等错条件的一方每个谓词配一个自己的condition_variable5先查再取if (!q.empty()) { auto v q.pop(); }两步之间被别人抢走让pop()自己阻塞并返回可判空的结果6忘了joinstd::thread对象析构时仍处于 joinable 状态所有路径都要join()或从 C20 起用std::jthread7用detach逃课线程detach后继续访问栈上的队列对象主线程已经退出作用域不要detach确需后台线程就用shared_ptr管理生命周期8出队后还用旧引用T r q.front(); q.pop();之后继续读rpop时把元素移动出来不保留对内部存储的引用第 8 条是容器通用问题但在队列上特别容易犯front()返回的是队列内部元素的引用pop()之后那个元素已经析构继续读它是悬垂引用属于 UB。上面的实现里T value std::move(q_.front()); q_.pop();就是这个原因——先把值搬出来再丢弃队列里的副本。总结主题结论同步原语进程内线程队列用std::mutexstd::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是进程间的开销远高于进程内队列一句话线程消息队列的难点从来不在数据结构而在什么时候改状态、什么时候通知、什么时候唤醒谁、什么时候收工这四件事上。把这四条按上面的表做对队列就只剩十几行样板代码了。
返回列表