基于进程线程同步的C生产者消费者程序问题求助
进程与线程同步的生产者消费者模型C程序调试求助
程序设计说明
- 2个生产者进程,各自循环生产单一类型的产品(类型1和类型2)
- 1个消费者进程,包含5个线程:1个分发线程(distributor)、2个消费类型1产品的线程、2个消费类型2产品的线程
- 消费者进程内有两个不同容量的产品缓存队列,生产者与消费者通过单个共享管道(pipe)通信
当前问题
程序已完成大部分开发,但运行存在异常,需排查同步逻辑、队列操作、线程/进程管理等方面的问题
代码实现
#include <stdio.h> #include <stdlib.h> #include <string.h> #include <unistd.h> #include <dirent.h> #include <stdbool.h> #include <fcntl.h> #include <pthread.h> #define BUFFER_SIZE1 20 #define BUFFER_SIZE2 30 typedef struct { int count; int productType; } product; int count = 0; int fd[2]; pthread_mutex_t lock; pthread_cond_t cond; typedef struct { product *values; int head; int tail; int numEntries; int size; } queue; queue q1; queue q2; void producer(int args); void *consumer(void *args); void *distributor(void *args); void initQ(queue *q, int size); bool QEmpty(queue *q); bool QFull(queue *q); bool put(queue *q, product prod); product get(queue *q); int main(int argc, char const *argv[]) { pthread_t th[5]; if (pipe(fd) == -1) { perror("error creating pipe"); exit(1); } initQ(&q1, BUFFER_SIZE1); initQ(&q2, BUFFER_SIZE2); int pid1; int pid2; int consId1 = 1; int consId2 = 2; pthread_mutex_init(&lock, NULL); pthread_cond_init(&cond, NULL); if ((pid1 = fork()) == 0) { producer(1); } else if ((pid2 = fork()) == 0) { producer(2); } else { if (pthread_create(&th[4], NULL, &distributor, NULL) != 0) { perror("Error creating distributor thread"); } for (int i = 0; i < 4; i++) { if (i < 2) { if (pthread_create(&th[i], NULL, &consumer, &consId1) != 0) { perror("Error creating thread"); } } else { if (pthread_create(&th[i], NULL, &consumer, &consId2) != 0) { perror("Error creating thread"); } } } for (int i = 0; i < 5; i++) { if (pthread_join(th[i], NULL) != 0) { perror("Error joining thread"); } } close(fd[0]); close(fd[1]); } return 0; } void producer(int args) { int prodCount = 0; product prod; prod.productType = args; close(fd[0]); while (1) { prodCount++; prod.count = prodCount; write(fd[1], &prod, sizeof(prod)); int time = (rand() % (200000 - 10000 + 1)) + 10000; usleep(time); } close(fd[1]); } void *consumer(void *args) { int consCount = 0; FILE *fp = fopen("output.txt", "a"); if (fp == NULL) { perror("Failed to open output file"); pthread_exit(NULL); } product prod; int prodType = *(int *)args; while (1) { if (prodType == 1) { prod = get(&q1); consCount++; fprintf(fp, "Thread ID: %d\n", prodType); fprintf(fp, "Product Type: %d\n", prod.productType); fprintf(fp, "Production Sequence #: %d\n", prod.count); fprintf(fp, "Consumption Sequence #: %d\n\n", consCount); } else { prod = get(&q2); consCount++; fprintf(fp, "Thread ID: %d\n", prodType); fprintf(fp, "Product Type: %d\n", prod.productType); fprintf(fp, "Production Sequence #: %d\n", prod.count); fprintf(fp, "Consumption Sequence #: %d\n\n", consCount); } fflush(fp); } fclose(fp); pthread_exit(NULL); } void *distributor(void *args) { close(fd[1]); product prod; while (1) { if (read(fd[0], &prod, sizeof(prod)) <= 0) { perror("Read from pipe failed"); break; } pthread_mutex_lock(&lock); while ((prod.productType == 1 && QFull(&q1)) || (prod.productType == 2 && QFull(&q2))) { pthread_cond_wait(&cond, &lock); } if (prod.productType == 1) { put(&q1, prod); } else { put(&q2, prod); } pthread_cond_broadcast(&cond); pthread_mutex_unlock(&lock); } close(fd[0]); pthread_exit(NULL); } void initQ(queue *q, int size) { q->size = size; q->values = malloc(sizeof(product) * q->size); q->numEntries = 0; q->head = 0; q->tail = 0; } bool QEmpty(queue *q) { return (q->numEntries == 0); } bool QFull(queue *q) { return (q->numEntries == q->size); } bool put(queue *q, product prod) { if (QFull(q)) { return false; } q->values[q->tail] = prod; q->numEntries++; q->tail = (q->tail + 1) % q->size; return true; } product get(queue *q) { pthread_mutex_lock(&lock); while (QEmpty(q)) { pthread_cond_wait(&cond, &lock); } product result = q->values[q->head]; q->head = (q->head + 1) % q->size; q->numEntries--; pthread_cond_broadcast(&cond); pthread_mutex_unlock(&lock); return result; }
核心问题修复说明
1. 进程创建逻辑修正
原代码中pid1 = fork() == 0的运算符优先级错误,导致pid1被赋值为布尔值而非进程ID,修正为(pid1 = fork()) == 0,确保正确获取子进程ID。
2. 线程启动顺序与执行逻辑修正
- 原代码先等待消费者线程结束再执行分发逻辑,但消费者线程是无限循环,导致分发逻辑永远无法运行。调整为先创建分发线程,再创建消费者线程。
- 实现独立的
distributor函数,将管道读取、队列分发逻辑移入,确保主线程不被阻塞。
3. 队列初始化与循环计算修正
initQ中head和tail初始值从NULL改为0,符合整数类型的索引定义。get函数中循环队列索引计算从位运算&改为取余%,避免队列大小非2的幂时出现逻辑错误。
4. 同步机制完善
- 所有队列操作(
put、get)均通过互斥锁保护,避免竞态条件。 - 新增条件变量等待逻辑:分发线程等待队列有空闲空间,消费者线程等待队列有数据,确保生产者-消费者模型的同步正确性。
5. 文件写入逻辑修正
- 消费者线程以追加模式
"a"打开输出文件,避免多线程写入时的文件截断问题。 - 修正
fprintf缺少文件指针的错误,添加fflush确保写入内容实时落盘。
内容的提问来源于stack exchange,提问作者chami
相关产品推荐
相关产品推荐

