什么是生产消费者模式?

生产消费者模式是一种多线程协作的设计模式,核心思想是通过共享缓冲区(如队列)解耦生产者与消费者。生产者负责生成数据并放入缓冲区,消费者从缓冲区取出数据并处理。两者通过同步机制(如互斥锁、条件变量)协调工作,避免资源竞争或数据不一致。

生产消费者模式的特点

  1. 解耦性:生产者与消费者不直接交互,通过缓冲区降低耦合度。
  2. 异步性:生产者无需等待消费者处理即可继续生产,提高吞吐量。
  3. 平衡性:通过缓冲区大小调节生产与消费的速度差异,避免系统过载。
  4. 线程安全:通过同步机制确保缓冲区的数据访问安全。

适用场景

  1. 任务调度:如线程池中的任务队列,生产者提交任务,消费者执行任务。
  2. 数据流水线:日志处理、消息中间件等场景,数据生产与消费速率不一致。
  3. 事件驱动系统:GUI应用中的事件队列,用户输入(生产)与事件处理(消费)。

简单的生产-消费者模式例程

生产-消费者模式的核心构成主要有:

1.用于放置待处理任务的缓冲区(一般用线程安全的queue实现)

2.生产者线程:用于提交任务到缓冲区,若缓冲区已满则线程阻塞等待

3.消费者线程:用于将任务拿出缓冲区并执行,若缓冲区非空则阻塞等待

4.互斥锁(mutex): 用于保护对缓冲区的访问,保证同一时刻只有一个线程对缓冲区进行读(写)操                              作

5.条件变量(std::condition_variable):用于生产者和消费者之间的线程同步,即生产者线程在缓冲区未满时生产数据,否则等待消费者线程的通知。消费者线程在缓冲区不为空时消费数据,否则等待生产者线程的通知。

#include <iostream>
#include <queue>
#include <thread>
#include <mutex>
#include <condition_variable>

std::queue<int> buffer;   //缓冲区
const int max_size = 10;  //缓冲区最大值
std::mutex mtx;           //保护缓冲区的互斥锁
std::condition_variable cv_producer, cv_consumer;  //用于同步生产者和消费者线程的条件变量


void producer(int id) {
    for (int i = 0; i < 20; ++i) {
        std::unique_lock<std::mutex> lock(mtx);
        //条件变量退出等待条件,返回false继续等待,返回true退出等待
        cv_producer.wait(lock, [] { return buffer.size() < max_size; });
        
        buffer.push(i);
        std::cout << "Producer " << id << " produced " << i << std::endl;
        
        lock.unlock();
        //放入一个数据就唤醒一个消费者线程
        cv_consumer.notify_one();
    }
}

void consumer(int id) {
    for (int i = 0; i < 20; ++i) {
        std::unique_lock<std::mutex> lock(mtx);
        cv_consumer.wait(lock, [] { return !buffer.empty(); });
        
        int item = buffer.front();
        buffer.pop();
        std::cout << "Consumer " << id << " consumed " << item << std::endl;
        
        lock.unlock();
        //同理,消费一个数据,就唤醒一个生产者线程
        cv_producer.notify_one();
    }
}

int main() {
    std::thread producers[2];
    std::thread consumers[2];

    for (int i = 0; i < 2; ++i) {
        producers[i] = std::thread(producer, i + 1);
        consumers[i] = std::thread(consumer, i + 1);
    }

    for (int i = 0; i < 2; ++i) {
        producers[i].join();
        consumers[i].join();
    }

    return 0;
}

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐