如何正确使用C++ std::barrier?及相关程序偶发死锁bug排查
一、如何正确使用C++的std::barrier?
std::barrier是C++20引入的同步原语,用于协调一组线程在某个同步点汇合,所有线程到达后再继续执行。正确使用方式如下:
- 初始化:创建时指定参与同步的线程总数(阈值),例如:
// 11个线程(10个工作线程+1个主线程)参与同步 std::barrier sync_barrier(11); - 同步等待:线程调用
arrive_and_wait()后会阻塞,直到所有指定数量的线程都调用该方法,随后所有线程同时解除阻塞,barrier自动重置,可重复使用。 - 动态调整线程数:若某个线程不再参与后续同步,调用
arrive_and_drop(),这会将阈值永久减1,后续同步只需剩余线程到达即可。 - 注意事项:
- 确保所有参与同步的线程都能正常到达同步点,避免因线程异常退出导致barrier永远无法满足阈值条件,引发永久阻塞。
- 不要用同一个barrier同步不同规模的线程组,除非通过
arrive_and_drop()提前调整阈值。 - 同步逻辑要匹配业务需求,避免错误的同步时机导致线程执行顺序混乱。
二、偶发死锁bug排查与修复
问题分析
程序偶发死锁的核心原因有两个:
- 队列并发访问无保护:主线程向
std::queuepush数据时未加锁,而工作线程操作队列时加了锁。std::queue本身不是线程安全的,并发读写会导致内部结构损坏,出现数据丢失或队列状态异常,进而引发工作线程因取不到数据永久阻塞、主线程因等待barrier同步永久阻塞的死锁。 - 同步逻辑错误:工作线程每处理一个数就调用
arrive_and_wait(),而主线程仅每10个数调用一次。这种设计导致barrier同步时机完全混乱,大量工作线程会因主线程未同步而提前阻塞,偶发出现同步计数不匹配,触发死锁。
修复方案
1. 保护队列的并发访问
主线程push数据时必须加锁,确保队列操作的线程安全:
// main函数中的push逻辑修改 std::lock_guard<std::mutex> lock(m); procs.push(p); cnd.notify_one();
2. 调整同步逻辑,匹配业务需求
业务要求主线程每生成10个数后,等待所有工作线程处理完这批数再继续。因此需要调整barrier的使用时机:
- 工作线程处理完当前批次的所有可处理数据后,再调用
arrive_and_wait()。 - 主线程生成完10个数后,调用
arrive_and_wait()等待所有工作线程处理完毕。
修复后的完整代码
#include <iostream> #include <thread> #include <queue> #include <mutex> #include <condition_variable> #include <barrier> #include <vector> #include <atomic> thread_local int k=0; void test(std::stop_token t, std::barrier<> &b, std::queue<int> &p, std::mutex& m, std::condition_variable& cnd, std::atomic<int>& processed) { int i=-9999; std::unique_lock<std::mutex> uk(m, std::defer_lock); while(true) { uk.lock(); cnd.wait(uk, [&]{return !p.empty() || t.stop_requested();}); if(t.stop_requested()) { uk.unlock(); break; } // 处理当前批次内的可用数据 int cnt = 0; while(!p.empty()) { i = p.front(); p.pop(); cnt++; std::cout << "thread: " << std::this_thread::get_id() << " k: " << ++k << " i: "<< i << std::endl; processed.fetch_add(1); } uk.unlock(); // 等待批次处理完成的同步 b.arrive_and_wait(); } std::cout << "finished" << std::endl; } int main(){ std::barrier workDone(11); // 10个工作线程+主线程 std::vector<std::jthread> threads; std::queue<int> procs; std::mutex m ; std::stop_source r ; std::stop_token t{r.get_token()} ; std::condition_variable cnd ; std::atomic<int> processed(0); for(int i=0;i<10;i++) { threads.emplace_back(test, t, std::ref(workDone), std::ref(procs), std::ref(m), std::ref(cnd), std::ref(processed)); } for(int batch=0; batch<3; batch++) { // 生成当前批次的10个数 std::lock_guard<std::mutex> lock(m); for(int num=1; num<=10; num++) { int p = batch*10 + num; procs.push(p); } cnd.notify_all(); // 唤醒所有工作线程处理 // 等待所有工作线程完成当前批次处理 workDone.arrive_and_wait(); // 重置计数器,准备下一批次 processed.store(0); } r.request_stop(); cnd.notify_all(); std::cout << "end" << std::endl ; }
修复说明
- 队列操作全程加锁,彻底解决并发访问导致的数据竞争问题。
- 每批次生成10个数后,主线程唤醒所有工作线程处理,然后通过barrier等待所有工作线程完成当前批次处理,确保业务逻辑的正确性。
- 使用原子变量
processed可辅助验证批次处理的完成情况,进一步增强逻辑可靠性。
内容的提问来源于stack exchange,提问作者depeng ji
相关产品推荐
相关产品推荐

