跳转到内容
新建笔记

C++ 条件变量:谓词、交接、定向通知与关闭

条件变量让线程等待受保护的状态满足某个条件,避免不停地轮询。正确用法包括三个相互配合的部分:保存状态的变量、保护状态的互斥量,以及等待/通知的条件变量。通知本身不保存任务,也不把互斥量解锁。

本页采用C++17,先说明等待语义,再把原有“货物交接、指定线程、交替和超时”的例子整理为能够结束的协议。

std::condition_variable使用已经持锁的 std::unique_lock<std::mutex>。一次等待会:

  1. 原子地释放互斥量并进入等待,允许修改状态的线程取得锁。
  2. 因通知或虚假唤醒而离开等待。
  3. 重新取得互斥量,然后检查谓词;条件不满足就继续等待。

cv.wait(lock, predicate)等价于在持锁条件下反复执行 while (!predicate()) cv.wait(lock)。返回后仍持有互斥量,不是“通知后自动解锁所有线程”。同一个条件变量上的并发等待者必须使用同一个互斥量。C++工作草案:condition_variable

// 使用片段:mutex、cv和ready必须存在,并遵守同一同步协议。
std::unique_lock<std::mutex> lock(mutex); // 立即加锁。
cv.wait(lock, [&] { return ready; });
// 此处仍持锁,ready已经在锁内检查为true。

condition_variable_any可以配合满足要求的其他锁类型;不代表任意对象都能当锁使用。

接口意义
notify_one()使一个已阻塞等待者有机会醒来;不能选择其线程ID。
notify_all()使所有已阻塞等待者有机会醒来,各自仍需争锁并检查谓词。
wait无指定超时,直到谓词满足才正常继续。
wait_for相对时长等待。
wait_until按指定时钟等待至目标时刻。
cv_status枚举类型,含 timeout、no_timeout;不是业务状态对象。

2. 先修改状态,再通知

跳转到“2. 先修改状态,再通知”

生产者应在锁内修改谓词依赖的状态,再通知等待者。可以在仍持锁时通知,也可以解锁后通知;后者常减少被唤醒者马上再次争锁的机会,但不是唯一合法顺序。无论哪种写法,互斥量、状态和条件变量都必须在相关操作期间有效。

通知早于等待并不必然丢工作:谓词应观察持久状态。如果生产者已经设置 ready=true,稍后到来的消费者取得锁便能直接继续。反过来,只发通知、不记录状态,就无法区分“工作已经到达”和“还没有发生”。

一个货物槽:双向交接和关闭

跳转到“一个货物槽:双向交接和关闭”

以下完整程序的容量为1。生产者必须等待槽空,消费者必须等待有货;所有读取都持同一把锁,不能在锁外 while (cargo != 0)轮询普通变量。用 optional表达槽是否占用,因而数值0也可作为合法货物。

#include <cassert>
#include <condition_variable>
#include <mutex>
#include <optional>
#include <thread>
class Mailbox {
public:
bool send(int value)
{
std::unique_lock<std::mutex> lock(mutex_);
empty_.wait(lock, [&] { return closed_ || !slot_; });
if (closed_)
return false;
slot_ = value;
lock.unlock();
full_.notify_one();
return true;
}
std::optional<int> receive()
{
std::unique_lock<std::mutex> lock(mutex_);
full_.wait(lock, [&] { return closed_ || slot_.has_value(); });
if (!slot_)
return std::nullopt;
const int value = *slot_;
slot_.reset();
lock.unlock();
empty_.notify_one();
return value;
}
void close()
{
{
std::lock_guard<std::mutex> lock(mutex_);
closed_ = true;
}
empty_.notify_all();
full_.notify_all();
}
private:
std::mutex mutex_;
std::condition_variable empty_, full_;
std::optional<int> slot_;
bool closed_ = false;
};
int main()
{
Mailbox mailbox;
int sum = 0;
std::thread consumer([&] {
while (const auto value = mailbox.receive())
sum += *value;
});
try {
for (int i = 0; i != 10; ++i) {
const bool sent = mailbox.send(i);
assert(sent);
(void)sent;
}
} catch (...) {
mailbox.close();
consumer.join();
throw;
}
mailbox.close();
consumer.join();
assert(sum == 45);
assert(!mailbox.send(99));
assert(!mailbox.receive().has_value());
}

关闭会拒绝新货物,但消费者仍可取走已存入的最后一项;槽空且关闭才结束。两个条件变量分别对应“槽空”和“有货”,避免多个生产者/消费者使用同一个 notify_one时唤醒错误类别。析构前必须先关闭并等待所有使用者结束。

3. 指定某个工作线程:状态按对象分开

跳转到“3. 指定某个工作线程:状态按对象分开”

普通 notify_one不选择特定线程。如果每个工作线程都有自己的通知计数和条件变量,就可以针对一个逻辑编号发送工作。计数器比单个布尔值更适合保留多次通知;下例每个编号只配一个消费者。

#include <array>
#include <cassert>
#include <condition_variable>
#include <cstddef>
#include <limits>
#include <mutex>
#include <stdexcept>
#include <thread>
#include <vector>
class DirectedEvents {
public:
static constexpr std::size_t size = 5;
bool post(std::size_t id)
{
auto& event = events_.at(id); // 越界抛out_of_range。
{
std::lock_guard<std::mutex> lock(mutex_);
if (closed_)
return false;
if (pending_.at(id) == std::numeric_limits<std::size_t>::max())
throw std::overflow_error("too many pending events");
++pending_.at(id);
}
event.notify_one();
return true;
}
bool wait(std::size_t id)
{
auto& event = events_.at(id);
std::unique_lock<std::mutex> lock(mutex_);
event.wait(lock, [&] { return closed_ || pending_.at(id) != 0; });
if (pending_.at(id) == 0)
return false;
--pending_.at(id);
return true;
}
void close()
{
{
std::lock_guard<std::mutex> lock(mutex_);
closed_ = true;
}
for (auto& event : events_)
event.notify_all();
}
private:
std::mutex mutex_;
std::array<std::condition_variable, size> events_;
std::array<std::size_t, size> pending_{};
bool closed_ = false;
};
int main()
{
DirectedEvents events;
std::array<int, DirectedEvents::size> handled{};
std::vector<std::thread> workers;
workers.reserve(DirectedEvents::size);
try {
for (std::size_t id = 0; id != DirectedEvents::size; ++id)
workers.emplace_back([&, id] {
while (events.wait(id))
++handled[id];
});
const bool first = events.post(2);
const bool second = events.post(2);
assert(first && second);
(void)first;
(void)second;
} catch (...) {
events.close();
for (auto& worker : workers)
worker.join();
throw;
}
events.close();
for (auto& worker : workers)
worker.join();
assert((handled == std::array<int, 5>{0, 0, 2, 0, 0}));
}

这里的2是应用自定义索引,不是系统线程ID。只通知2号但对全部线程直接 join会让其他线程永远等待;所以关闭协议必须唤醒全部等待者。先收到关闭通知的2号仍会处理完剩余计数,然后退出。

如果目的是两个线程交替执行,核心状态可以改成 turn与 closed:线程0等待 closed || turn==0,处理后将 turn=1;线程1反向处理。达到约定轮数时设置关闭并通知全部线程。单独循环调用 notify_one或插入睡眠不保证交替,且没有终止条件的 while(true)不是完整示例。上面的单槽交接已经提供一种自然交替的生产/消费协议。

4. 超时返回不是任务取消

跳转到“4. 超时返回不是任务取消”

不带谓词的定时等待返回 cv_status,但无超时返回也可能只是虚假唤醒。带谓词重载返回最终谓词是否为true,不能把它简单当作“是否发生过一次notify”。

#include <cassert>
#include <chrono>
#include <condition_variable>
#include <mutex>
int main()
{
std::mutex mutex;
std::condition_variable changed;
bool ready = false;
std::unique_lock<std::mutex> lock(mutex);
const auto deadline = std::chrono::steady_clock::now()
+ std::chrono::milliseconds(2);
const bool available = changed.wait_until(lock, deadline, [&] {
return ready;
});
assert(!available && lock.owns_lock());
ready = true;
assert(changed.wait_for(lock, std::chrono::milliseconds(0), [&] {
return ready;
}));
}

本例没有生产者,因而第一段最终为false;已经为true的谓词不会因为0时长就失败。等待时间只约束这次等待,不会停止另一个线程的阻塞输入或I/O;超时后直接 join一个仍卡在 std::cin的线程仍可能无限等待。

如果自己在循环里处理虚假唤醒,应保持同一个绝对截止时间;每醒一次又重新等待完整相对时长,可能使总等待反复延长。实际返回还要重新取得互斥量,调度和争锁也会使观察到的时长超过设定值。

先确认谓词依赖哪些变量、它们是否都用同一锁保护,再检查每个状态变化后该通知谁。最后检查关闭能否唤醒所有等待者,以及条件变量和锁是否一直存在到线程不再使用它们。锁外轮询、把通知当消息、通知错等待群体和缺少退出协议,是这类程序的四个常见问题。

示例验证交接总数、针对某个编号的多次事件、关闭排空和超时谓词;有限运行通过不能穷举所有调度。完整任务队列的所有权与异常传递见线程池。