【问题标题】:How to implement a re-usable thread barrier with std::atomic如何使用 std::atomic 实现可重用的线程屏障
【发布时间】:2014-08-04 00:11:00
【问题描述】:

我有 N 个线程执行各种任务,这些线程必须定期与线程屏障同步,如下图所示,有 3 个线程和 8 个任务。 ||表示时间障碍,所有线程必须等到完成8个任务才能重新开始。

Thread#1  |----task1--|---task6---|---wait-----||-taskB--|          ...
Thread#2  |--task2--|---task5--|-------taskE---||----taskA--|       ...
Thread#3  |-task3-|---task4--|-taskG--|--wait--||-taskC-|---taskD   ...

我找不到可行的解决方案,认为 Semaphores http://greenteapress.com/semaphores/index.html 的小书很有启发性。我想出了一个使用下面显示的 std::atomic 的解决方案,它“似乎”正在使用三个 std::atomic。 我担心我的代码在极端情况下崩溃,因此引用了动词。那么您能否分享有关验证此类代码的建议?你有更简单的防呆代码吗?

std::atomic<int> barrier1(0);
std::atomic<int> barrier2(0);
std::atomic<int> barrier3(0);

void my_thread()
{

  while(1) {
    // pop task from queue
    ...
    // and execute task 
    switch(task.id()) {
      case TaskID::Barrier:
        barrier2.store(0);
        barrier1++;
        while (barrier1.load() != NUM_THREAD) {
          std::this_thread::yield();
        }
        barrier3.store(0);
        barrier2++;
        while (barrier2.load() != NUM_THREAD) {
          std::this_thread::yield();
        }
        barrier1.store(0);
        barrier3++;
        while (barrier3.load() != NUM_THREAD) {
          std::this_thread::yield();
        }
       break;
     case TaskID::Task1:
       ...
     }
   }
}

【问题讨论】:

  • 既然知道信号量,那么使用信号量有什么问题呢?
  • C++11 标准库中没有原生 std::semaphore,所以你必须使用 std::mutex 和 std::condition_variable。
  • 这有帮助吗? stackoverflow.com/questions/8115267/…我可以通过一票关闭重复,但我确定这是否是同一个问题。
  • @R.MartinhoFernandes 有趣的问题是我们是否需要旋转解决方案。通常情况下,屏障不会以这种方式实现,因为预计某些线程将不得不等待其他线程完成,而我们不想在此期间因忙于等待而阻塞核心。基于condition_variable 将等待线程发送到睡眠的解决方案在这里似乎更合适。
  • 你是对的,这是关键问题。我的建议是一个旋转解决方案,所以它可能没有那么有效。我将切换到找到 here 的 [condition_variable]。

标签: multithreading c++11


【解决方案1】:

Boost 提供 barrier implementation 作为 C++11 标准线程库的扩展。如果使用 Boost 是一种选择,那么您应该再看看。

如果您必须依赖标准库设施,您可以基于std::mutex 和std::condition_variable 推出自己的实现,而不会有太多麻烦。

class Barrier {
    int wait_count;
    int const target_wait_count;
    std::mutex mtx;
    std::condition_variable cond_var;

    Barrier(int threads_to_wait_for)
     : wait_count(0), target_wait_count(threads_to_wait_for) {}

    void wait() {
        std::unique_lock<std::mutex> lk(mtx);
        ++wait_count;
        if(wait_count != target_wait_count) {
            // not all threads have arrived yet; go to sleep until they do
            cond_var.wait(lk, 
                [this]() { return wait_count == target_wait_count; });
        } else {
            // we are the last thread to arrive; wake the others and go on
            cond_var.notify_all();
        }
        // note that if you want to reuse the barrier, you will have to
        // reset wait_count to 0 now before calling wait again
        // if you do this, be aware that the reset must be synchronized with
        // threads that are still stuck in the wait
    }
};

与基于原子的解决方案相比,此实现的优势在于,在condition_variable::wait 中等待的线程应该由操作系统的调度程序发送到睡眠状态,因此您不会因为等待线程在屏障上旋转而阻塞 CPU 内核。

关于重置障碍的几句话:最简单的解决方案是使用单独的reset() 方法并让用户确保永远不会同时调用reset 和wait。但在许多用例中,这对用户来说并不容易实现。

对于自重置屏障,您必须考虑等待计数的争用:如果在从wait 返回的最后一个线程之前重置等待计数,则某些线程可能会卡在屏障中。这里一个聪明的解决方案是不让终止条件依赖于等待计数变量本身。相反,您引入了第二个计数器,该计数器仅由调用notify 的线程增加。然后其他线程观察该计数器的变化以确定是否退出等待:

void wait() {
    std::unique_lock<std::mutex> lk(mtx);
    unsigned int const current_wait_cycle = m_inter_wait_count;
    ++wait_count;
    if(wait_count != target_wait_count) {
        // wait condition must not depend on wait_count
        cond_var.wait(lk, 
            [this, current_wait_cycle]() { 
                return m_inter_wait_count != current_wait_cycle;
            });
    } else {
        // increasing the second counter allows waiting threads to exit
        ++m_inter_wait_count;
        cond_var.notify_all();
    }
}

在所有线程在inter_wait_count 溢出之前离开等待的(非常合理的)假设下,此解决方案是正确的。

【讨论】:

  • boost Barrier 不可重复使用,作为注释包含在代码中的注释标记了将 wait_count 重置为 0 的难度。我在编写代码时遇到了同样的问题,即只有 2原子变量,存在死锁问题。
  • @user3636086 提升屏障是可重复使用的。它会在所有线程到达后自动重置,您可以根据需要再次调用 wait。
  • @user3636086 我补充了几句关于如何实现重置的内容。
【解决方案2】:

对于原子变量,使用其中三个作为屏障简直是矫枉过正,只会使问题复杂化。您知道线程的数量,因此您可以简单地在每次线程进入屏障时自动递增单个计数器,然后旋转直到计数器变得大于或等于 N。像这样:

void barrier(int N) {
    static std::atomic<unsigned int> gCounter = 0;
    gCounter++;
    while((int)(gCounter - N) < 0) std::this_thread::yield();
}

如果您的线程数不超过 CPU 内核数且预期等待时间较短,您可能需要删除对 std::this_thread::yield() 的调用。这个调用可能真的很昂贵(超过一微秒,我敢打赌,但我没有测量它)。根据您的任务规模,这可能很重要。

如果您想重复设置障碍,只需增加N 即可:

unsigned int lastBarrier = 0;
while(1) {
    switch(task.id()) {
        case TaskID::Barrier:
            barrier(lastBarrier += processCount);
            break;
    }
}

【讨论】:

  • 简单高效。我已经尝试过我的代码,工作正常,但由于线程池的性质,我不得不将 while(gCounter
  • 我故意使用gCounter &lt; N:据我所知,您的线程不能保证在屏障调用之间同步。因此,一个线程可能会提前运行,而另一个线程仍然在屏障内被抢占,并在后期线程离开屏障之前调用下一个屏障。在这种情况下,计数器可能会跳过后期线程正在等待的N,导致它死锁。我想,你对gCounter &lt; N 的问题是你想让柜台环绕吗?我已经编辑了我的答案以在这些情况下工作:当您期望溢出时,您需要使用无符号整数类型。
  • 请注意,如果您的线程是 Linux 上的高优先级实时 FIFO 线程,则此代码可能 a) 占用 cpu,b) 导致优先级反转。不要使用此代码。只需使用std::mutex + std::condition_variable。
  • @MaximYegorushkin 我不知道你用 a) 和 b) 指的是什么,但我的印象是你没有完全阅读我的答案:我特别说 yield() 电话除非每个核心最多有一个线程,否则不应删除。只要没有到达屏障的线程继续前进,这段代码就不会锁定。
  • @cmaster 在任何情况下都不会死锁,但它可能会大大降低系统速度,即使线程数等于或略小于物理内核数。忙碌等待的问题在于,您基本上是在欺骗系统,让系统认为您有 很多 工作要做。因此,就调度程序分配的 CPU 时间和 CPU 消耗的实际电力而言,您将获得大量资源。如果您决定将这些资源用于循环非常快在单个变量上,那是您的决定。
【解决方案3】:

我想指出,在@ComicSansMS 给出的解决方案中, 在执行cond_var.notify_all();之前,wait_count应该被重置为0

这是因为当第二次调用屏障时,如果 wait_count 未重置为 0,则 if 条件将始终失败。

【讨论】:

    猜你喜欢
    • 2016-02-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-05-24
    • 1970-01-01
    • 2018-05-11
    • 1970-01-01
    • 2015-08-25
    相关资源
    最近更新 更多