阻塞队列模型

内容介绍

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 后是否允许继续读取剩余数据,要在接口语义里写清楚。

学习路径