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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 12:17:45