C++
手写有界阻塞队列
使用 mutex、condition_variable_any 和 stop_token 实现有容量、可关闭的阻塞队列。
发布于 2026年7月23日
手写有界阻塞队列
使用 mutex、condition_variable_any 和 stop_token 实现有容量、可关闭的阻塞队列。
本系列代码使用 C++20 和
oc::handmade命名空间,目标是解释实现机制、复杂度和工程边界,不是替代标准库。普通容器与缓存核心不内置互斥锁;这不代表 lock-free。
一、学习目标
- 实现背压和关闭状态机
- 正确处理虚假唤醒
- 让取消、关闭和数据竞争拥有确定语义
二、前置条件
熟悉 mutex、条件变量、谓词等待和 C++20 stop_token。
Linux/macOS:
cmake -S . -B build -DCMAKE_BUILD_TYPE=Debug
cmake --build build -j
ctest --test-dir build --output-on-failure
Windows PowerShell:
cmake -S . -B build -DCMAKE_BUILD_TYPE=Debug
cmake --build build --config Debug
ctest --test-dir build -C Debug --output-on-failure
三、问题与设计选择
一个互斥锁保护队列与 closed 状态;not_empty/not_full 分别唤醒消费者和生产者。close 后拒绝 push,但允许 drain 既有元素。
这里刻意保留一条边界:教学实现覆盖构造、复制移动、核心修改、查找和迭代契约,但不复刻标准库全部重载、ABI、constexpr、异构查找或节点句柄。
四、内存布局与核心不变量
0 <= size <= capacity;状态只从 open 变 closed;closed 且 empty 后 pop 永久失败。
每个修改操作都按“准备资源 → 构造新状态 → 提交连接或指针 → 清理旧状态”的顺序设计。提交点之前发生异常,应保持原对象可继续使用;无法提供强保证时,会在接口说明中明确基本保证。
五、核心实现
std::optional<T> pop(std::stop_token stop) {
std::unique_lock lock(mutex_);
if (!not_empty_.wait(lock, stop,
[&] { return closed_ || !queue_.empty(); }))
return std::nullopt;
if (queue_.empty()) return std::nullopt;
T value = std::move(queue_.front());
queue_.pop_front();
not_full_.notify_one();
return value;
}
上面先聚焦最容易写错的核心步骤;若本篇对应一个独立组件,下一节给出统一工程中的完整教学实现。代码没有放入 std 命名空间,避免未定义行为和名称冲突。
六、完整教学实现
下面是统一工程中经过 GCC、Clang、GoogleTest 和 Sanitizer 验证的完整组件。它依赖前序文章已经实现的公共类型以及头文件中的标准库 #include。
namespace oc::handmade {
template<class T>
class bounded_blocking_queue {
std::size_t capacity_;
std::deque<T> queue_;
bool closed_{};
mutable std::mutex mutex_;
std::condition_variable_any not_empty_;
std::condition_variable_any not_full_;
public:
explicit bounded_blocking_queue(std::size_t capacity) : capacity_(capacity) {
if (capacity == 0) throw std::invalid_argument("queue capacity must be positive");
}
bool push(T value) {
std::unique_lock lock(mutex_);
not_full_.wait(lock, [&] { return closed_ || queue_.size() < capacity_; });
if (closed_) return false;
queue_.push_back(std::move(value));
not_empty_.notify_one();
return true;
}
bool push(std::stop_token stop, T value) {
std::unique_lock lock(mutex_);
if (!not_full_.wait(lock, stop, [&] { return closed_ || queue_.size() < capacity_; }))
return false;
if (closed_) return false;
queue_.push_back(std::move(value));
not_empty_.notify_one();
return true;
}
std::optional<T> pop() {
std::unique_lock lock(mutex_);
not_empty_.wait(lock, [&] { return closed_ || !queue_.empty(); });
if (queue_.empty()) return std::nullopt;
T result = std::move(queue_.front());
queue_.pop_front();
not_full_.notify_one();
return result;
}
std::optional<T> pop(std::stop_token stop) {
std::unique_lock lock(mutex_);
if (!not_empty_.wait(lock, stop, [&] { return closed_ || !queue_.empty(); }))
return std::nullopt;
if (queue_.empty()) return std::nullopt;
T result = std::move(queue_.front());
queue_.pop_front();
not_full_.notify_one();
return result;
}
void close() noexcept {
{
std::scoped_lock lock(mutex_);
if (closed_) return;
closed_ = true;
}
not_empty_.notify_all();
not_full_.notify_all();
}
bool closed() const {
std::scoped_lock lock(mutex_);
return closed_;
}
std::size_t size() const {
std::scoped_lock lock(mutex_);
return queue_.size();
}
};
} // namespace oc::handmade
生产级标准库还要处理完整 allocator 传播、全部重载、ABI、调试迭代器和平台特化;这里保留的是能够独立推导核心数据结构的教学边界。
七、使用示例与输出
预期输出或状态:
容量 2 时第三个生产者等待;消费者取走一个元素后生产者继续;close 后剩余元素仍依次取出。
示例必须在文章对应的测试目标中实际编译。涉及顺序的输出只依赖接口明确承诺的顺序;无序容器不会把某次桶顺序写成稳定结果。
八、复杂度与失效规则
| 操作 | 复杂度 | 说明 |
|---|---|---|
| push/pop | 摊还 O(1) | 可能阻塞 |
| try_push/try_pop | O(1) | 立即返回 |
| close | O(等待者) | 唤醒所有线程 |
| size | O(1) | 需要加锁 |
复杂度中的 O(1) 若标记为“平均”或“摊还”,不能在面试中省略限定词。任何重新分配、节点删除、rehash 或缓存淘汰都必须单独说明迭代器、引用与指针是否失效。
九、异常安全与资源管理
- 获取资源后立即交给 RAII 对象或明确记录已构造数量。
- 用户类型构造、复制、移动、比较器和哈希器都可能抛异常。
- 只有在所有后续步骤不会失败时才修改不可回滚的链接。
- 析构、释放和关闭路径不得抛异常。
- 并发包装通过回调在锁内访问,避免返回保护对象的裸引用。
十、常见错误
1. 不用谓词等待条件变量
不用谓词等待条件变量会破坏本篇建立的契约。调试时先检查核心不变量,再缩小到触发该状态的最短操作序列。
2. close 只通知消费者不通知生产者
close 只通知消费者不通知生产者会破坏本篇建立的契约。调试时先检查核心不变量,再缩小到触发该状态的最短操作序列。
3. 持锁执行用户任务
持锁执行用户任务会破坏本篇建立的契约。调试时先检查核心不变量,再缩小到触发该状态的最短操作序列。
十一、面试追问
- 为什么需要两个条件变量?
- 关闭与取消有何区别?
- 如何证明没有丢失唤醒?
回答时先说数据结构不变量,再给复杂度,最后说明异常、迭代器或并发边界,通常比背诵结论更有说服力。
十二、练习与自测
- 实现超时 push/pop
- 写多生产者多消费者测试
- 画出 open/closed/drained 状态机
自测标准:能够不看代码画出内存或节点关系,解释一次成功操作和一次失败回滚,并写出至少一个会击穿错误实现的测试。