基于pthreads的多生产者单消费者问题实现受阻求助
嘿,看你在用pthreads实现多生产者-消费者模型时遇到阻碍了,我来给你梳理下核心要点,再整个能跑的完整实现示例吧!
核心问题拆解
你的场景里,多生产者读文件行到固定大小的共享缓冲区,消费者取数据写文件,关键要搞定这几个同步问题:
- 缓冲区满的时候,生产者得停下来等消费者腾出空间
- 缓冲区空的时候,消费者得等生产者塞数据
- 多个线程抢缓冲区时不能乱(竞态条件)
- 所有生产者读完文件后,要通知消费者把剩下的数据处理完再退出
必备的同步工具
在pthreads里,咱们得靠这俩家伙来管同步:
pthread_mutex_t:互斥锁,保证同一时间只有一个线程能碰共享缓冲区,避免数据乱掉pthread_cond_t:两个条件变量,一个告诉生产者“缓冲区有空位啦”,一个告诉消费者“缓冲区有数据啦”
完整实现代码
先把共享数据和参数结构定义好,方便线程传参:
#include <stdio.h> #include <stdlib.h> #include <pthread.h> #include <string.h> #include <unistd.h> #define BUFFER_SIZE 20 // 你定义的缓冲区大小 #define MAX_LINE_LENGTH 256 // 每行最大长度 // 共享数据结构,所有线程都要用到 typedef struct { char *buffer[BUFFER_SIZE]; int in; // 生产者写入的位置 int out; // 消费者读取的位置 int count; // 缓冲区里当前有多少条数据 pthread_mutex_t mutex; pthread_cond_t not_full; // 缓冲区不满的条件 pthread_cond_t not_empty; // 缓冲区不空的条件 int producers_done; // 已经完成任务的生产者数量 int total_producers; // 生产者总数 } SharedData; // 生产者的参数结构,包含共享数据和要读的文件名 typedef struct { SharedData *shared; const char *filename; } ProducerArgs; // 消费者的参数结构,包含共享数据和要写的输出文件名 typedef struct { SharedData *shared; const char *output_filename; } ConsumerArgs;
然后写生产者函数,负责读文件、往缓冲区塞数据:
void* Producer(void* args) { ProducerArgs *prod_args = (ProducerArgs*)args; SharedData *shared = prod_args->shared; FILE *fp = fopen(prod_args->filename, "r"); if (!fp) { perror("打不开输入文件"); pthread_exit(NULL); } char line[MAX_LINE_LENGTH]; while (fgets(line, MAX_LINE_LENGTH, fp) != NULL) { // 注意:必须给每行内容分配堆内存,不能用栈上的line! // 不然线程退出后栈内存会被回收,消费者拿到的就是垃圾数据 char *line_copy = malloc(strlen(line) + 1); if (!line_copy) { perror("内存分配失败"); break; } strcpy(line_copy, line); // 加锁,开始操作共享缓冲区 pthread_mutex_lock(&shared->mutex); // 如果缓冲区满了,就等着消费者取数据 // 这里一定要用while循环,不能用if!因为线程可能被虚假唤醒 while (shared->count == BUFFER_SIZE) { pthread_cond_wait(&shared->not_full, &shared->mutex); } // 把数据塞进缓冲区,更新位置和计数 shared->buffer[shared->in] = line_copy; shared->in = (shared->in + 1) % BUFFER_SIZE; shared->count++; // 通知消费者:缓冲区有数据啦,可以来取了 pthread_cond_signal(&shared->not_empty); pthread_mutex_unlock(&shared->mutex); // 模拟一下生产延迟(可选,去掉也没事) usleep(10000); } fclose(fp); // 当前生产者完成任务,更新计数,还要告诉消费者们 pthread_mutex_lock(&shared->mutex); shared->producers_done++; // 用broadcast而不是signal,确保所有等待的消费者都收到通知 pthread_cond_broadcast(&shared->not_empty); pthread_mutex_unlock(&shared->mutex); free(prod_args); pthread_exit(NULL); }
接下来是消费者函数,负责从缓冲区拿数据、写入输出文件:
void* Consumer(void* args) { ConsumerArgs *cons_args = (ConsumerArgs*)args; SharedData *shared = cons_args->shared; FILE *fp = fopen(cons_args->output_filename, "a"); if (!fp) { perror("打不开输出文件"); pthread_exit(NULL); } while (1) { pthread_mutex_lock(&shared->mutex); // 如果缓冲区空,而且还有生产者在干活,就等着 while (shared->count == 0 && shared->producers_done < shared->total_producers) { pthread_cond_wait(&shared->not_empty, &shared->mutex); } // 如果缓冲区空了,所有生产者也都干完了,那就退出循环 if (shared->count == 0 && shared->producers_done == shared->total_producers) { pthread_mutex_unlock(&shared->mutex); break; } // 从缓冲区取出数据,更新位置和计数 char *line = shared->buffer[shared->out]; shared->out = (shared->out + 1) % BUFFER_SIZE; shared->count--; // 通知生产者:缓冲区有空位啦,可以塞数据了 pthread_cond_signal(&shared->not_full); pthread_mutex_unlock(&shared->mutex); // 把数据写入文件,记得释放生产者分配的内存 fputs(line, fp); free(line); // 模拟消费延迟(可选) usleep(15000); } fclose(fp); free(cons_args); pthread_exit(NULL); }
最后是主函数,负责初始化线程和共享数据:
int main() { SharedData shared = {0}; // 初始化互斥锁和条件变量 pthread_mutex_init(&shared.mutex, NULL); pthread_cond_init(&shared.not_full, NULL); pthread_cond_init(&shared.not_empty, NULL); // 这里可以自己改输入文件列表和消费者数量 const char *input_files[] = {"input1.txt", "input2.txt", "input3.txt"}; shared.total_producers = sizeof(input_files) / sizeof(input_files[0]); int num_consumers = 2; // 创建生产者线程 pthread_t producers[shared.total_producers]; for (int i = 0; i < shared.total_producers; i++) { ProducerArgs *args = malloc(sizeof(ProducerArgs)); args->shared = &shared; args->filename = input_files[i]; pthread_create(&producers[i], NULL, Producer, args); } // 创建消费者线程 pthread_t consumers[num_consumers]; for (int i = 0; i < num_consumers; i++) { ConsumerArgs *args = malloc(sizeof(ConsumerArgs)); args->shared = &shared; args->output_filename = "output.txt"; pthread_create(&consumers[i], NULL, Consumer, args); } // 等待所有生产者线程完成 for (int i = 0; i < shared.total_producers; i++) { pthread_join(producers[i], NULL); } // 等待所有消费者线程完成 for (int i = 0; i < num_consumers; i++) { pthread_join(consumers[i], NULL); } // 清理资源 pthread_mutex_destroy(&shared.mutex); pthread_cond_destroy(&shared.not_full); pthread_cond_destroy(&shared.not_empty); return 0; }
重点提醒
- 内存安全:生产者一定要用
malloc给每行分配内存,消费者记得free,不然会内存泄漏 - 条件变量的while循环:绝对不能用
if代替,因为线程可能被操作系统虚假唤醒,必须重新检查条件 - 生产者完成通知:要用
pthread_cond_broadcast,确保所有等待的消费者都收到消息,避免有的消费者卡死 - 锁的范围:尽量缩短持有锁的时间,只在操作共享数据时加锁,不然会影响程序性能
内容的提问来源于stack exchange,提问作者Magna
相关产品推荐
相关产品推荐

