生产者-消费者模型中条件判断与线程同步问题求助
生产者消费者问题优化解答
问题描述
我有两类生产者:一类从文本文件读取偶数数据,另一类通过rand()函数生成奇数数据;同时有两类消费者,分别消费偶数和奇数数据。已使用bufferEmpty、bufferFull两个信号量及mutexBuffer互斥锁实现同步,但仍有以下疑问:
- 应在何处应用条件判断?
- 如何暂停当前执行缓冲区操作的线程?
- 如何避免线程遗漏数据?
当前代码
#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修改读取循环,以
obj >> data作为循环条件,确保只有有效数据才进入生产流程。 - 新增全局变量记录已生产总数据量,消费者以此作为循环终止条件,避免硬编码。
修改后的代码
#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
相关产品推荐
相关产品推荐

