阻塞队列模型
内容介绍
C++ 高性能并发里,阻塞队列常用来把“任务生产”和“任务消费”解耦。它的目标不是模仿某种语言的语法,而是降低共享状态复杂度,并为背压、批处理、线程池提供基础结构。
生产者只负责提交数据或任务,消费者阻塞等待并处理。这样可以避免多个线程到处直接修改同一个对象。
简化阻塞队列
#include <condition_variable>
#include <iostream>
#include <mutex>
#include <optional>
#include <queue>
#include <stdexcept>
#include <thread>
template <class T>
class Channel {
public:
void send(T value) {
{
std::lock_guard<std::mutex> lock(m_);
if (closed_) {
throw std::runtime_error("send on closed queue");
}
q_.push(std::move(value));
}
cv_.notify_one();
}
std::optional<T> recv() {
std::unique_lock<std::mutex> lock(m_);
cv_.wait(lock, [&] {
return closed_ || !q_.empty();
});
if (q_.empty()) {
return std::nullopt;
}
T value = std::move(q_.front());
q_.pop();
return value;
}
void close() {
{
std::lock_guard<std::mutex> lock(m_);
closed_ = true;
}
cv_.notify_all();
}
private:
std::mutex m_;
std::condition_variable cv_;
std::queue<T> q_;
bool closed_ = false;
};
int main() {
Channel<int> ch;
std::thread producer([&] {
for (int i = 0; i < 5; ++i) {
ch.send(i);
}
ch.close();
});
std::thread consumer([&] {
while (auto value = ch.recv()) {
std::cout << *value << '\n';
}
});
producer.join();
consumer.join();
}最佳代码实践
- 用阻塞队列传递任务或数据,减少共享可变状态。
- 队列要有明确关闭语义,否则消费者可能永远等待。
recv返回optional<T>可以表达“通道关闭且无数据”。- 多生产者多消费者场景要清楚谁负责
close。 - 高吞吐场景应考虑有界队列和批量
pop,避免无边界积压和频繁唤醒。
常见错误用法
while (true) {
auto value = ch.recv();
// 没有处理通道关闭
}问题:队列关闭和普通数据必须有区分。
注意事项
- 这个阻塞队列是教学版,没有容量限制、超时、公平性等高级能力。
- 有界队列可以在队列满时阻塞发送方,适合做背压。
close后是否允许继续读取剩余数据,要在接口语义里写清楚。
学习路径
- 上一节:cppfuture和promise
- 当前阶段:cpp多线程
- 下一节:cppWaitGroup模型