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

线程暂停执行同步操作的低开销实现方案问询

线程同步场景的实现方案探讨

场景描述

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();
    }
}

用户问题

  1. 基于barrier和atomic的方案是否能有效解决该问题?
  2. 是否存在更优的实现方案?

问题解答

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 19:44:51