条件变量让线程等待受保护的状态满足某个条件,避免不停地轮询。正确用法包括三个相互配合的部分:保存状态的变量、保护状态的互斥量,以及等待/通知的条件变量。通知本身不保存任务,也不把互斥量解锁。
本页采用C++17,先说明等待语义,再把原有“货物交接、指定线程、交替和超时”的例子整理为能够结束的协议。
1. wait实际做了什么
跳转到“1. wait实际做了什么”std::condition_variable使用已经持锁的 std::unique_lock<std::mutex>。一次等待会:
- 原子地释放互斥量并进入等待,允许修改状态的线程取得锁。
- 因通知或虚假唤醒而离开等待。
- 重新取得互斥量,然后检查谓词;条件不满足就继续等待。
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的线程仍可能无限等待。
如果自己在循环里处理虚假唤醒,应保持同一个绝对截止时间;每醒一次又重新等待完整相对时长,可能使总等待反复延长。实际返回还要重新取得互斥量,调度和争锁也会使观察到的时长超过设定值。
5. 检查顺序
跳转到“5. 检查顺序”先确认谓词依赖哪些变量、它们是否都用同一锁保护,再检查每个状态变化后该通知谁。最后检查关闭能否唤醒所有等待者,以及条件变量和锁是否一直存在到线程不再使用它们。锁外轮询、把通知当消息、通知错等待群体和缺少退出协议,是这类程序的四个常见问题。
示例验证交接总数、针对某个编号的多次事件、关闭排空和超时谓词;有限运行通过不能穷举所有调度。完整任务队列的所有权与异常传递见线程池。