线程暂停执行同步操作的低开销实现方案问询
线程同步场景的实现方案探讨
场景描述
N个线程同时对一个数据结构执行小增量操作,但偶尔需要执行同步动作——所有线程需暂停,等待某一线程完成该同步操作后再继续。要求无需同步时对线程影响尽可能小。
尝试的实现方案
方案1:基于barrier和atomic的实现
#include <vector> #include <thread> #include <atomic> #include <barrier> int main() { const size_t nr_threads = 10; std::vector<std::thread> threads; std::barrier barrier { nr_threads }; std::atomic_bool sync_required { false }; auto rare_condition = []() { return std::rand() == 42; }; for (int i = 0; i < nr_threads; ++i) { threads.emplace_back([&, i]() { while (true) { if (sync_required) { if (i == 0) { barrier.arrive_and_wait(); sync_required = false; // solo synchronized work barrier.arrive_and_wait(); } else { barrier.arrive_and_wait(); barrier.arrive_and_wait(); } } // standard loop body ... // sometimes a global synchronized action is required if (rare_condition()) sync_required = true; } }); } // eventually ... treads quit for (auto& thread : threads) { thread.join(); } }
方案2:基于shared_mutex的尝试方案
#include <array> #include <atomic> #include <thread> #include <vector> #include <barrier> #include <shared_mutex> int main() { const size_t nr_threads = 10; std::vector<std::thread> threads; std::shared_mutex sync_mtx; std::atomic_bool sync_required { false }; auto rare_condition = []() { return std::rand() == 42; }; for (int i = 0; i < nr_threads; ++i) { threads.emplace_back([&, i]() { std::shared_lock shared_lock { sync_mtx }; while (true) { // very rarely another thread requires all the others to stop for a bit if (sync_required) { if (i == 0) { // unlock shared, but lock unique, seems a little odd but neccesary shared_lock.unlock(); { std::unique_lock unique_lock{ shared_lock }; sync_required = false; // solo sync work } shared_lock.lock(); } else { shared_lock.unlock(); // todo: need condition variable, which adds more complexity to this solution? shared_lock.lock(); } } // sometimes a global syncronized action is required sync_required = sync_required || rare_condition(); } }); } // eventually ... treads quit for (auto& thread : threads) { thread.join(); } }
用户问题
- 基于barrier和atomic的方案是否能有效解决该问题?
- 是否存在更优的实现方案?
问题解答
1. 基于barrier和atomic的方案有效性分析
该方案存在明显正确性缺陷,无法可靠满足需求:
- 竞态条件:当多个线程同时检测到
sync_required为true时,线程0可能还在执行同步工作,其他线程就已经完成两次屏障等待并回到主循环继续修改数据,破坏了"所有线程暂停等待同步完成"的核心要求。 - 重复触发逻辑混乱:如果同步过程中又有线程设置
sync_required为true,会导致屏障等待逻辑冲突,无法正确处理连续同步请求。 - 单点依赖风险:若线程0在同步过程中异常退出,其他线程会永久阻塞在屏障等待上,缺乏容错性。
2. 更优实现方案推荐
针对"无同步时开销极小,偶尔全局同步"的需求,推荐以下几种经过验证的方案:
方案一:atomic_flag + 条件变量(低常态开销)
常态下线程仅需一次轻量原子检查,触发同步时通过条件变量避免忙等:
#include <vector> #include <thread> #include <atomic> #include <mutex> #include <condition_variable> std::atomic_flag sync_pending = ATOMIC_FLAG_INIT; std::mutex sync_mtx; std::condition_variable sync_cv; bool sync_completed = false; void thread_func(int thread_id) { while (true) { // 常态下仅原子检查,开销极低 if (sync_pending.test(std::memory_order_acquire)) { std::unique_lock<std::mutex> lock(sync_mtx); // 等待同步触发信号 sync_cv.wait(lock, []{ return sync_pending.test(std::memory_order_relaxed); }); if (thread_id == 0) { // 执行同步工作 // ... sync_completed = true; sync_cv.notify_all(); // 重置状态 sync_pending.clear(std::memory_order_release); sync_completed = false; } else { // 等待线程0完成同步 sync_cv.wait(lock, []{ return sync_completed; }); } } // 执行增量操作 // ... // 触发同步的条件 if (std::rand() == 42) { sync_pending.test_and_set(std::memory_order_release); sync_cv.notify_all(); } } } int main() { const size_t nr_threads = 10; std::vector<std::thread> threads; for (int i = 0; i < nr_threads; ++i) { threads.emplace_back(thread_func, i); } for (auto& t : threads) t.join(); return 0; }
优势:常态下开销接近无锁;同步时通过条件变量避免CPU空转,资源占用低。
方案二:修正后的barrier方案
针对原barrier方案的缺陷,用原子操作确保同步状态只被处理一次:
#include <vector> #include <thread> #include <atomic> #include <barrier> std::atomic_bool sync_required = false; // 初始化可重用屏障,线程数为10 std::barrier sync_barrier{10}; void thread_func(int thread_id) { while (true) { // 原子交换确保仅一次处理同步请求 if (sync_required.exchange(false, std::memory_order_acquire)) { // 等待所有线程到达屏障,完成暂停 sync_barrier.arrive_and_wait(); if (thread_id == 0) { // 执行同步工作 // ... } // 等待同步工作完成,所有线程继续执行 sync_barrier.arrive_and_wait(); } // 增量操作 // ... if (std::rand() == 42) { sync_required.store(true, std::memory_order_release); } } }
优势:利用barrier的原生同步能力,逻辑清晰;修正后避免了竞态条件,可处理重复同步请求。
方案三:优化后的shared_mutex方案
针对原shared_mutex方案的逻辑错误,补充条件变量实现正确同步:
#include <vector> #include <thread> #include <atomic> #include <shared_mutex> #include <condition_variable> std::shared_mutex rw_mtx; std::atomic_bool sync_needed = false; std::mutex cv_mtx; std::condition_variable sync_cv; void thread_func(int thread_id) { while (true) { std::shared_lock<std::shared_mutex> r_lock(rw_mtx); if (sync_needed) { r_lock.unlock(); std::unique_lock<std::shared_mutex> w_lock(rw_mtx); if (thread_id == 0) { // 执行同步工作 // ... sync_needed = false; sync_cv.notify_all(); } else { // 等待同步完成 std::unique_lock<std::mutex> cv_lock(cv_mtx); sync_cv.wait(cv_lock, []{ return !sync_needed; }); } r_lock.lock(); } // 增量操作(共享锁保证线程安全) // ... if (std::rand() == 42) { sync_needed = true; } } }
优势:共享锁天然保护增量操作的线程安全,适合本身需要线程安全访问数据结构的场景。
内容的提问来源于stack exchange,提问作者Aedoro
相关产品推荐
相关产品推荐

