CHAPTER 26 / Linux 与并发系统
条件变量、有界队列与关闭协议
队列为空时消费者应该等什么,服务关闭后又怎样保证所有等待者能退出?
这一章要弄清楚
- 用谓词等待而非依赖通知次数
- 为有界队列设计背压和 drain 关闭语义
- 覆盖空队列关闭与被阻塞生产者
先备知识:线程、互斥锁与死锁
C++20 / macOS 与 Linux;本章的硬件模型只推演逻辑,不代表设备性能。
通知不是一份可积攒的许可
设生产者提交工作,消费者处理工作。队列空时,让消费者不停检查会浪费 CPU;条件变量允许它暂停,等可能改变条件的事件。关键是等待“队列非空或已经关闭”这个状态谓词,而不是等待“有人发过通知”。通知本身不保存工作,也没有保证下一次 wait 能消费先前通知。队列和关闭标记才是持久状态,必须由同一 mutex 保护。
先区分mutex与这次持锁责任(补充,另估20分钟)
std::unique_lock<std::mutex> lock(m);先对仍存活的mutex m加锁;成功后lock拥有这次锁。原文件的std::unique_lock lock(m)是同一具体类型的简写:这里从构造实参m的mutex类型推导类模板实参,不是创建一种无类型的锁。unique_lock允许显式lock.unlock():释放m并记录自己不再持锁;析构只在仍持锁时解锁,所以不会因为手动unlock后离开作用域而再解一次。不能在不持锁时再次调用unlock。与第25章的lock_guard相比,这个可改变的持锁状态才适合交给wait暂时释放、再重新获得。unique_lock合同
沿22-a的容量3队列作给定状态推演:push(9)先取得锁,发现size=3而等待;等待期间释放m,消费者才能把size改成2。该push被唤醒并重新取得m后,若队列未关闭且仍有空位,就写入9,size回到3;显式unlock后不再持锁,函数退出也不重复解锁。若取得锁时size=0且未关闭,谓词已真,无需实际睡眠就能入队。若观察到closed而提前return,没有执行显式unlock,析构负责释放。下面再看wait如何反复检查这个谓词。
wait(lock, predicate) 的逻辑相当于循环检查谓词:为假时原子地释放锁并进入等待;被唤醒后重新获取锁,再检查。原子释放并等待避免检查为空与真正睡眠之间丢失关键唤醒。虚假唤醒是接口允许的;即使没有虚假唤醒,另一消费者也可能先取走元素,所以不能把 while 改成 if。通知只说明“值得再检查”,并不承诺条件现在仍然成立。
有界容量就是最小的背压
若生产者每秒产生一千个任务,而消费者只能处理五百个,无界队列会持续增长。有界队列设容量 C:满时生产者等待 closed || size < C,空时消费者等待 closed || !empty。这把处理速度反馈给上游,叫背压(backpressure)。容量零要么被明确设计成同步 rendezvous,要么拒绝;本章选择拒绝,避免永远等不到空位。
实现通常使用两个条件变量:not_empty 唤醒消费者,not_full 唤醒生产者。状态修改发生在锁内;随后可以解锁再通知,使被唤醒者更有机会直接拿锁。通知放在锁内也可以正确,选择应考虑对象生命周期及整体协议。队列被销毁前,必须保证没有线程还在访问它,不能以一次 notify_all 代替 join。
一项新工作与所有等待者的终止是不同通知(补充,另估15分钟)
notify_one()从当前阻塞在这个条件变量上的线程中唤醒一个(若有),不保证选到哪一个;notify_all()唤醒所有当前等待者。没有等待者时,两者都不会为未来存下一张通知票。醒来仍须重新取得锁、检查谓词,通知不替它取得锁或消费工作。条件变量通知合同
给定空队列和两个等待消费者,生产者入队7后调用available.notify_one;即使两个消费者因其他允许的唤醒原因都重新检查,这次只有一项7可被取走,另一个看到空队列仍应等待。pop移走一项后通知space,给满队列上等待的生产者重新检查空间。close则要让两类等待者都观察到终止状态,因此分别对available与space调用notify_all;只叫醒一个可能留下再也没有新状态变化可等的线程。通知数量不是入队数量,也不是完成数量。
先定义关闭的含义
本章采用 drain 语义:close 后拒绝新的 push;已经入队的元素继续被 pop;队列空且关闭时 pop 返回空结果表示终止。另一种 cancel 语义会丢弃未处理元素,需要另写合同。不要让 bool false 同时表示“暂时为空”“永久关闭”“任务值为假”,接口应让调用者能够判断下一步。
关闭标记在锁内从 false 变为 true,随后对两个条件变量都 notify_all。只通知消费者会让满队列上的生产者永远等待。重复 close 是安全的幂等操作。等待者醒来仍在锁内检查 closed,生产者不能在关闭后趁空位把新任务塞进去。消费者则先取已有任务,再在空状态返回终止。
两类测试解决两类误解
第一个例子用容量 3 的真实线程队列传递 1 到 100,消费者累加后得到 5050,再通过关闭退出。它检查数据路径和正常排空。第二个例子专门把生产者堵在容量为 1 的满队列上,再关闭队列,验证 push 返回失败且已存的 7 仍可取出。它通过受锁保护的 waiting 标记确定进入等待阶段,而不是用 sleep 猜调度。
超时测试可以发现挂起,但不能靠长时间等待证明协议正确。画出每个等待谓词,以及每次状态变化会唤醒哪类线程,往往更容易发现遗漏。实际队列还需要决定异常、取消、移动失败、多个生产者与消费者的策略;本章 int 队列刻意避开抛异常的元素移动,让同步协议先清楚可见。
异常退出先关闭,再等待
元素是 int 不代表入队永远不失败,deque 扩容仍可能抛异常。下载程序用 CloseAndJoin 守卫保证生产端失败时先 close,让空队列上等待的消费者醒来,再 join;如果只 join 而不关闭,消费者可能永远等新任务。附带回归通过同一把队列锁观察消费者的等待状态,再显式抛出教学异常,外层捕获后检查消费者已退出且累加仍为零。测试钩子只服务本例的单消费者验证,不是通用队列接口,也没有使用 sleep 或真的制造内存耗尽。
满队列到关闭
容量 1,元素 7 占用槽位,9 还没入队。
持锁设置 closed,然后通知所有等待者。
生产者拒绝 9;消费者仍能取 7,然后结束。
阅读完整推演文字
- 容量已满
队列:[7];生产者:等待空位;closed:false
容量 1,元素 7 占用槽位,9 还没入队。
- 发起关闭
队列:[7];通知:notify_all;closed:true
持锁设置 closed,然后通知所有等待者。
- 排空退出
push(9):false;pop:7;下次 pop:结束
生产者拒绝 9;消费者仍能取 7,然后结束。
跟着例子,走完一遍
容量三的生产者消费者队列
生产者按顺序放入 1..100;消费者累加,close 后排空退出。
- push 在满时等待,pop 在空时等待。
- 每次出队通知 not_full,每次入队通知 not_empty。
- close 唤醒两边,消费者排空后结束。
#include <condition_variable>
#include <deque>
#include <iostream>
#include <mutex>
#include <optional>
#include <stdexcept>
#include <thread>
class Queue {
std::mutex m; std::condition_variable available,space,consumer_state;
std::deque<int> data; bool closed=false,consumer_waiting=false;
public:
bool push(int x) {
std::unique_lock lock(m);
space.wait(lock,[&]{return closed || data.size()<3;});
if(closed) return false;
data.push_back(x); lock.unlock(); available.notify_one(); return true;
}
std::optional<int> pop() {
std::unique_lock lock(m);
consumer_waiting=!closed && data.empty();
if(consumer_waiting) consumer_state.notify_all();
available.wait(lock,[&]{return closed || !data.empty();});
consumer_waiting=false;
if(data.empty()) return std::nullopt;
int x=data.front(); data.pop_front(); lock.unlock(); space.notify_one(); return x;
}
void close() {
{std::lock_guard lock(m); closed=true;}
available.notify_all(); space.notify_all();
}
// Test hook for this one-consumer example: observe wait while holding the same mutex.
void wait_for_consumer_for_test() {
std::unique_lock lock(m);
consumer_state.wait(lock,[&]{return consumer_waiting;});
}
};
struct InjectedFailure {};
void consume(int& sum,bool& finished,bool inject_failure=false) {
Queue q; std::thread consumer;
struct CloseAndJoin {
Queue& queue; std::thread& thread;
~CloseAndJoin() {queue.close(); if(thread.joinable()) thread.join();}
} close_and_join{q,consumer};
consumer=std::thread([&]{while(auto x=q.pop()) sum+=*x; finished=true;});
if(inject_failure) {
q.wait_for_consumer_for_test();
// Explicit producer failure injection; no real allocation failure is claimed.
throw InjectedFailure{};
}
for(int i=1;i<=100;++i) if(!q.push(i)) throw std::runtime_error("unexpected close");
q.close(); q.close(); consumer.join();
if(sum!=5050 || q.push(101) || q.pop()) throw std::runtime_error("queue contract");
}
int main() {
int sum=0; bool finished=false; consume(sum,finished);
if(!finished) throw std::runtime_error("consumer completion");
int empty_sum=0; bool empty_finished=false,caught=false;
try {consume(empty_sum,empty_finished,true);} catch(const InjectedFailure&) {caught=true;}
if(!caught || !empty_finished || empty_sum!=0) throw std::runtime_error("close before join");
std::cout<<"sum="<<sum<<'\n';
}sum=5050
状态谓词保证虚假唤醒与先到通知都不影响正确性。
在本机运行这个例子
下载后,在文件所在目录执行。需要支持 C++20 的编译器;POSIX 示例还需要章节说明中的系统条件。
clang++ -std=c++20 -Wall -Wextra -Wpedantic -Werror -pthread 22-a.cpp -o example && ./example预期标准输出:
sum=5050
关闭唤醒满队列上的生产者
一个槽已有 7;生产者尝试再放 9 并等待。主线程确认它进入等待后关闭。
- waiting 标志只用于本例确定测试交错。
- 关闭后 notify_all,生产者醒来检查 closed。
- push 被拒绝,旧元素 7 保持不变。
#include <condition_variable>
#include <iostream>
#include <mutex>
#include <stdexcept>
#include <thread>
int main() {
std::mutex m; std::condition_variable space,entered;
bool closed=false,waiting=false,rejected=false; int slot=7;
std::thread producer([&] {
std::unique_lock lock(m); waiting=true; entered.notify_one();
space.wait(lock,[&]{return closed || slot==0;});
if(closed) rejected=true; else slot=9;
});
{
std::unique_lock lock(m); entered.wait(lock,[&]{return waiting;}); closed=true;
}
space.notify_all(); producer.join();
if(!rejected || slot!=7) throw std::runtime_error("shutdown");
std::cout<<"push_rejected="<<rejected<<" retained="<<slot<<'\n';
}push_rejected=1 retained=7
关闭是谓词的一部分,生产者不用等消费者腾空才有机会退出。
在本机运行这个例子
下载后,在文件所在目录执行。需要支持 C++20 的编译器;POSIX 示例还需要章节说明中的系统条件。
clang++ -std=c++20 -Wall -Wextra -Wpedantic -Werror -pthread 22-b.cpp -o example && ./example预期标准输出:
push_rejected=1 retained=7
关闭只通知 not_empty
消费者退出了,满队列上的生产者仍等着空位,join 永远不返回。
修正思路:push 和 pop 的等待谓词都包含 closed;close 更新状态后唤醒两组等待者,并明确排空语义。
轮到你动手
先写预测或代码,再按需打开提示。完整答案用于对照自己的推理。
练习 1
为队列添加 capacity 参数并拒绝 0。设计四个关闭边界案例。
给我一点提示
- 区分空队列关闭和满队列关闭。
- close 重复调用不应改变语义。
查看答案与推理
构造时 capacity==0 抛 invalid_argument。测试空队列 close 后 pop 结束;有元素 close 后排空;满队列阻塞 push 被 close 唤醒并失败;重复 close 无异常,关闭后所有 push 失败。
练习 2
为何不能在 wait 前用 if(data.empty()),醒来就直接取 front?
给我一点提示
- 假设两个消费者被同时唤醒。
- 考虑没有元素也可能醒来的接口合同。
查看答案与推理
某消费者先取完唯一元素,另一个拿锁后已没有元素;还存在虚假唤醒。用谓词 wait,只有 closed 或非空才继续,并在 closed 且空时返回终止。
把理解说出来
先用中文讲清因果,再用英文回答。问题依据技能主题编写,并非公司内部题库。
条件变量保存通知次数吗?
参考回答 / English answer
不保存一般通知次数。共享状态保存条件,通知只促使等待线程重新检查。
A condition variable does not store notification credits. The protected state records the condition, and notifications prompt a recheck.为什么 wait 必须配谓词或 while?
参考回答 / English answer
允许虚假唤醒,其他线程也可能先消费资源。醒来后必须持锁重新验证条件。
Wakeups do not guarantee that the condition still holds. I recheck under the mutex because of spurious wakeups and competing consumers.wait 如何避免检查后、睡前的丢唤醒?
参考回答 / English answer
wait 原子地释放传入 mutex 并进入等待;生产者持同一锁改变谓词,封闭这一竞争窗口。
Wait atomically releases the mutex and begins waiting. The producer changes the predicate under the same mutex.drain 与 cancel 关闭有什么差别?
参考回答 / English answer
drain 处理已有任务再退出;cancel 丢弃未开始任务。新提交处理、返回值和异常通知必须明确。
Drain preserves queued work before termination. Cancel discards pending work, so the API must specify how callers learn that outcome.析构函数只调用 close 就安全了吗?
参考回答 / English answer
不一定,访问队列的线程可能还在运行。对象销毁前需要外部所有者 join 或其他生命周期保证。
Close changes the protocol state but does not finish every thread. The owner must ensure all users have stopped before destruction.为什么要限制队列容量?
参考回答 / English answer
生产速率长期大于消费速率会持续占内存;容量限制让上游等待或拒绝,形成背压。
A bounded queue prevents unlimited backlog growth. Blocking or rejecting producers propagates backpressure to the source.继续查证
- OSTEP:Condition Variables ↗
30.1 定义、30.2 生产者消费者
- C++ draft:condition_variable ↗
wait、谓词等待与 notify_all
公开资料用于查证;本章图解和例题是独立教学内容。CPU 逻辑模型不能证明设备性能。