You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

生产者-消费者模型中条件判断与线程同步问题求助

生产者消费者问题优化解答

问题描述

我有两类生产者:一类从文本文件读取偶数数据,另一类通过rand()函数生成奇数数据;同时有两类消费者,分别消费偶数和奇数数据。已使用bufferEmpty、bufferFull两个信号量及mutexBuffer互斥锁实现同步,但仍有以下疑问:

  1. 应在何处应用条件判断?
  2. 如何暂停当前执行缓冲区操作的线程?
  3. 如何避免线程遗漏数据?

当前代码

#include <iostream>
#include <fstream>
#include <pthread.h>
#include <semaphore.h>
#include <cstdlib>
#include <unistd.h>

using namespace std;

#define BUFFER_SIZE 10
int buffer[BUFFER_SIZE];
int count = 0;
int readcount = 0;
sem_t bufferEmpty, bufferFull;
pthread_mutex_t mutexBuffer;

void* producer1(void* args){
    fstream obj;
    cout<<"I will read data from OS_A2_Q1.txt"<<endl;
    obj.open("OS_A2_Q1.txt",ios::in);
    if (!obj){
        cout<<"Error opening file"<<endl;
        exit(0);
    }
    else{
        int data;
        while (1) {
            obj >>data;
            //cout<<data<<endl;
            cout<<"I am producer 1"<<endl;
            sem_wait(&bufferEmpty);
            pthread_mutex_lock(&mutexBuffer);               
            buffer[count] = data;
            count++;
            pthread_mutex_unlock(&mutexBuffer);
            sem_post(&bufferFull);
            if (obj.eof())
                break;
            sleep(1);
        }
        cout<<"I have successfully completed reading data from file"<<endl;
    }
    obj.close();
}
void* producer2(void* args){
    cout<<"I will insert 100 random numbers into the buffer"<<endl;
    int x,i=0;
    
    while (i < 100){
        while (1){
            x = rand() % 100;
            if ((x % 2 != 0) && (x != 0))
                break;
        }
        i++;
        cout<<"inside prod 2"<<endl;
        sem_wait(&bufferEmpty);
        pthread_mutex_lock(&mutexBuffer);               
        buffer[count] = x;
        count++;
        pthread_mutex_unlock(&mutexBuffer);
        sem_post(&bufferFull);
        sleep(1);
    }
}
void* consumer1(void* args){
    while (readcount < 200) {
        int y;
        // Remove from the buffer
        sem_wait(&bufferFull);
        pthread_mutex_lock(&mutexBuffer);
        y = buffer[count - 1];
        readcount++;
        count--;
        pthread_mutex_unlock(&mutexBuffer);
        sem_post(&bufferEmpty);

        // Consume
        if(y % 2 == 0){
            printf("\nI am Consumer 1\n");
            printf("\nGot %d\n", y);
        }
        sleep(1);
    }
}
void* consumer2(void* args){
    while (readcount < 200) {
        int y;
        // Remove from the buffer
        sem_wait(&bufferFull);
        pthread_mutex_lock(&mutexBuffer);
        y = buffer[count - 1];
        readcount++;
        count--;
        pthread_mutex_unlock(&mutexBuffer);
        sem_post(&bufferEmpty);

        // Consume
        if(y % 2 != 0){
            printf("\nI am Consumer 2\n");
            printf("\nGot %d\n", y);
        }
    }
}

疑问解答

1. 条件判断的位置

生产者侧

  • 文件读取有效性判断:必须在写入缓冲区前判断读取是否成功,否则会把eof后的无效数据写入缓冲区。修改方式:将obj >> data的结果作为循环条件,读取失败直接退出,无需执行后续生产操作。
  • 奇数生成判断:当前逻辑可行,但可以简化为直接生成奇数(比如x = (rand() % 50) * 2 + 1),减少无效循环次数。

消费者侧

  • 数据类型匹配判断:必须在持有互斥锁且未移除数据时判断。当前代码先移除数据再判断,会导致不匹配的数据被丢弃(遗漏)。正确逻辑是:在锁内检查缓冲区顶部数据是否是自己需要的类型,是则移除并消费;不是则释放锁和信号量,重新等待。

2. 暂停线程的方式

线程暂停(阻塞)依赖以下同步原语:

  • 信号量阻塞:sem_wait(&bufferEmpty)会让生产者在缓冲区满时阻塞;sem_wait(&bufferFull)会让消费者在缓冲区空时阻塞,这是最基础的暂停方式。
  • 条件变量阻塞:针对“消费者需要等待特定类型数据”的场景,需新增condEven、condOdd两个条件变量,配合互斥锁使用。当消费者拿到非目标类型数据时,调用pthread_cond_wait()进入阻塞,等待对应类型数据的生产者唤醒。

3. 避免数据遗漏的核心优化

当前代码的遗漏风险:

  • 消费者移除数据后发现类型不匹配,直接跳过导致数据丢失。
  • 生产者1在eof后仍会执行一次写入(obj >> data失败后仍写入缓冲区)。
  • readcount < 200硬编码判断,若文件读取的偶数数量不是100,会导致多消费或未消费完。

优化方案:

  1. 消费者在锁内检查数据类型,不匹配则不修改缓冲区,释放锁后重新等待。
  2. 生产者1修改读取循环,以obj >> data作为循环条件,确保只有有效数据才进入生产流程。
  3. 新增全局变量记录已生产总数据量,消费者以此作为循环终止条件,避免硬编码。

修改后的代码

#include <iostream>
#include <fstream>
#include <pthread.h>
#include <semaphore.h>
#include <cstdlib>
#include <unistd.h>

using namespace std;

#define BUFFER_SIZE 10
int buffer[BUFFER_SIZE];
int count = 0;
int totalProduced = 0; // 记录总生产数据量
bool producersDone = false; // 标记所有生产者是否结束
sem_t bufferEmpty, bufferFull;
pthread_mutex_t mutexBuffer;
pthread_cond_t condEven, condOdd; // 新增条件变量,等待特定类型数据

// 生产者1:读取偶数
void* producer1(void* args){
    fstream obj;
    cout<<"生产者1:从OS_A2_Q1.txt读取偶数"<<endl;
    obj.open("OS_A2_Q1.txt",ios::in);
    if (!obj){
        cout<<"打开文件失败"<<endl;
        exit(0);
    }
    int data;
    // 以读取成功作为循环条件,避免无效数据写入
    while (obj >> data) {
        if (data % 2 != 0) { // 跳过文件中的奇数
            continue;
        }
        cout<<"生产者1:准备写入数据"<<data<<endl;
        sem_wait(&bufferEmpty);
        pthread_mutex_lock(&mutexBuffer);               
        buffer[count] = data;
        count++;
        totalProduced++;
        pthread_cond_signal(&condEven); // 唤醒等待偶数的消费者
        pthread_mutex_unlock(&mutexBuffer);
        sem_post(&bufferFull);
        sleep(1);
    }
    cout<<"生产者1:文件读取完成"<<endl;
    obj.close();
    pthread_exit(NULL);
}

// 生产者2:生成奇数
void* producer2(void* args){
    cout<<"生产者2:生成100个奇数写入缓冲区"<<endl;
    int x,i=0;
    while (i < 100){
        // 直接生成1-99的奇数
        x = (rand() % 50) * 2 + 1;
        i++;
        cout<<"生产者2:准备写入数据"<<x<<endl;
        sem_wait(&bufferEmpty);
        pthread_mutex_lock(&mutexBuffer);               
        buffer[count] = x;
        count++;
        totalProduced++;
        pthread_cond_signal(&condOdd); // 唤醒等待奇数的消费者
        pthread_mutex_unlock(&mutexBuffer);
        sem_post(&bufferFull);
        sleep(1);
    }
    pthread_exit(NULL);
}

// 消费者1:消费偶数
void* consumer1(void* args){
    while (true) {
        int y;
        sem_wait(&bufferFull);
        pthread_mutex_lock(&mutexBuffer);
        // 等待直到缓冲区有偶数,或所有生产者已结束
        while (count > 0 && buffer[count-1] % 2 != 0) {
            if (producersDone) break;
            pthread_cond_wait(&condEven, &mutexBuffer);
        }
        if (count == 0 && producersDone) { // 缓冲区空且生产者结束,退出
            pthread_mutex_unlock(&mutexBuffer);
            sem_post(&bufferFull);
            break;
        }
        if (count == 0) { // 缓冲区空但生产者未结束,继续等待
            pthread_mutex_unlock(&mutexBuffer);
            sem_post(&bufferFull);
            continue;
        }
        // 取出数据
        y = buffer[count - 1];
        count--;
        pthread_mutex_unlock(&mutexBuffer);
        sem_post(&bufferEmpty);

        // 消费数据
        cout<<"消费者1:获取到偶数"<<y<<endl;
        sleep(1);
    }
    pthread_exit(NULL);
}

// 消费者2:消费奇数
void* consumer2(void* args){
    while (true) {
        int y;
        sem_wait(&bufferFull);
        pthread_mutex_lock(&mutexBuffer);
        // 等待直到缓冲区有奇数,或所有生产者已结束
        while (count > 0 && buffer[count-1] % 2 == 0) {
            if (producersDone) break;
            pthread_cond_wait(&condOdd, &mutexBuffer);
        }
        if (count == 0 && producersDone) {
            pthread_mutex_unlock(&mutexBuffer);
            sem_post(&bufferFull);
            break;
        }
        if (count == 0) {
            pthread_mutex_unlock(&mutexBuffer);
            sem_post(&bufferFull);
            continue;
        }
        // 取出数据
        y = buffer[count - 1];
        count--;
        pthread_mutex_unlock(&mutexBuffer);
        sem_post(&bufferEmpty);

        // 消费数据
        cout<<"消费者2:获取到奇数"<<y<<endl;
        sleep(1);
    }
    pthread_exit(NULL);
}

// 主函数
int main() {
    sem_init(&bufferEmpty, 0, BUFFER_SIZE);
    sem_init(&bufferFull, 0, 0);
    pthread_mutex_init(&mutexBuffer, NULL);
    pthread_cond_init(&condEven, NULL);
    pthread_cond_init(&condOdd, NULL);

    pthread_t p1, p2, c1, c2;
    pthread_create(&p1, NULL, producer1, NULL);
    pthread_create(&p2, NULL, producer2, NULL);
    pthread_create(&c1, NULL, consumer1, NULL);
    pthread_create(&c2, NULL, consumer2, NULL);

    pthread_join(p1, NULL);
    pthread_join(p2, NULL);
    // 标记生产者全部结束,唤醒所有等待的消费者
    pthread_mutex_lock(&mutexBuffer);
    producersDone = true;
    pthread_cond_broadcast(&condEven);
    pthread_cond_broadcast(&condOdd);
    pthread_mutex_unlock(&mutexBuffer);

    pthread_join(c1, NULL);
    pthread_join(c2, NULL);

    sem_destroy(&bufferEmpty);
    sem_destroy(&bufferFull);
    pthread_mutex_destroy(&mutexBuffer);
    pthread_cond_destroy(&condEven);
    pthread_cond_destroy(&condOdd);

    return 0;
}

内容的提问来源于stack exchange,提问作者Syed Muhammad Ismail

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 19:46:59