【问题标题】:Correct way to wait a condition variable that is notified by several threads等待由多个线程通知的条件变量的正确方法
【发布时间】:2016-07-20 18:24:48
【问题描述】:

我正在尝试使用 C++11 并发支持来做到这一点。

我有一种工作线程的线程池,它们都做同样的事情,其中​​一个主线程有一个条件变量数组(每个线程一个,它们需要“启动”同步,即不提前运行一个周期他们的循环)。

    for (auto &worker_cond : cond_arr) {
        worker_cond.notify_one();
    }

那么这个线程必须等待池中每个线程的通知才能再次重新启动它的循环。这样做的正确方法是什么?有一个条件变量并等待某个整数,每个不是主线程的线程都会增加?类似的东西(仍在主线程中)

    unique_lock<std::mutex> lock(workers_mtx);
    workers_finished.wait(lock, [&workers] { return workers = cond_arr.size(); });

【问题讨论】:

标签: c++ multithreading c++11 concurrency


【解决方案1】:

我在这里看到两个选项:

选项 1:join()

基本上不是使用条件变量在线程中开始计算,而是为每次迭代生成一个新线程并使用join() 等待它完成。然后为下一次迭代生成新线程,依此类推。

选项 2:锁

只要其中一个线程仍在工作,您就不希望主线程发出通知。所以每个线程都有自己的锁,它在计算之前锁定它,然后解锁。您的主线程在调用 notify() 之前锁定所有这些,然后再解锁它们。

【讨论】:

  • 对于第一个,不会更容易使用。第二个不起作用,因为当主线程正在等待某个锁时,另一个具有不同锁的线程可能会循环多次,这是我不想要的!
  • 如果在每个周期后工作线程释放锁,则不会等待条件变量,并且只有在通知重新锁定并恰好执行一个周期时。
【解决方案2】:

我认为您的解决方案没有根本性的问题。

用workers_mtx 保护workers 并完成。

我们可以用计数信号量来抽象它。

struct counting_semaphore {
  std::unique_ptr<std::mutex> m=std::make_unique<std::mutex>();
  std::ptrdiff_t count = 0;
  std::unique_ptr<std::condition_variable> cv=std::make_unique<std::condition_variable>();

  counting_semaphore( std::ptrdiff_t c=0 ):count(c) {}
  counting_semaphore(counting_semaphore&&)=default;

  void take(std::size_t n = 1) {
    std::unique_lock<std::mutex> lock(*m);
    cv->wait(lock, [&]{ if (count-std::ptrdiff_t(n) < 0) return false; count-=n; return true; } );
  }
  void give(std::size_t n = 1) {
    {
      std::unique_lock<std::mutex> lock(*m);
      count += n;
      if (count <= 0) return;
    }
    cv->notify_all();
  }
};

take 带走count,如果不够就阻塞。

give 与count 相加,并通知是否有正数。

现在工作线程在两个信号量之间传送令牌。

std::vector< counting_semaphore > m_worker_start{count};
counting_semaphore m_worker_done{0}; // not count, zero
std::atomic<bool> m_shutdown = false;

// master controller:
for (each step) {
  for (auto&& starts:m_worker_start)
    starts.give();
  m_worker_done.take(count);
}

// master shutdown:
m_shutdown = true;
// wake up forever:
for (auto&& starts:m_worker_start)
  starts.give(std::size_t(-1)/2);

// worker thread:
while (true) {
  master->m_worker_start[my_id].take();
  if (master->m_shutdown) return;
  // do work
  master->m_worker_done.give();
}

或类似的。

live example.

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-08-04
    • 2012-04-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-02-22
    相关资源
    最近更新 更多