You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何正确使用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排查与修复

问题分析

程序偶发死锁的核心原因有两个:

  1. 队列并发访问无保护:主线程向std::queue push数据时未加锁,而工作线程操作队列时加了锁。std::queue本身不是线程安全的,并发读写会导致内部结构损坏,出现数据丢失或队列状态异常,进而引发工作线程因取不到数据永久阻塞、主线程因等待barrier同步永久阻塞的死锁。
  2. 同步逻辑错误:工作线程每处理一个数就调用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 20:07:33