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

多线程C生产消费程序求助:需完善锁与条件变量实现

C多线程生产消费程序问题求助

需求说明

  • 产品类型由创建它的进程决定,只能为1或2
  • 产品ID每生产一个递增1
  • 消费者ID固定为1、2、3或4,取决于消费产品的线程
  • 消费计数每消费一个递增1

我已经完成了大部分C多线程生产消费程序的编写,但程序存在问题。我知道需要通过pthread_mutex_t锁和pthread_cond_t条件变量来解决同步问题,但目前陷入了困境,希望能得到帮助。

原代码

#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[])
{
    // Creating 5 threads 4 consumer and 1 distributor
    pthread_t th[5];
    // Creating our pipe, fd[0] is read end, fd[1] is write end
    if (pipe(fd) == -1)
    {
        perror("error creating pipe");
        exit(1);
    }

    // Initializing both buffers
    initQ(&q1, BUFFER_SIZE1);
    initQ(&q2, BUFFER_SIZE2);

    int pid1;
    int pid2;
    int consId1 = 1;
    int consId2 = 2;

    // Initializing lock
    pthread_mutex_init(&lock, NULL);

    // Initialziing condition variables
    pthread_cond_init(&cond, NULL);

    // Create first producer process using fork(), child process 1
    pid1 = fork();
    if (pid1 == 0)
    {
        producer(1);
    }
    // Create second prodcuer process using fork(), child process 2
    pid2 = fork();
     if ( pid2== 0)
    {
        producer(2);
    }
    // Create distrubtor and consumer threads, parent process
    else
    {
        // Creating 4 threads using for loop and pthread_create
        for (int i = 0; i < 4; i++)
        {
            // 2 consumer threads for product type 1
            if (i == 1 || i == 2)
            {
                if (pthread_create(&th[i], NULL, &consumer, &consId1) != 0)
                {
                    perror("Error creating thread");
                }
            }
            // 2 consumer threads for product type 2
            else
            {
                if (pthread_create(&th[i], NULL, &consumer, &consId2) != 0)
                {
                    perror("Error creating thread");
                }
            }
        }
        // use pthread_join to wait for preivous thread to terminate
        for (int i = 0; i < 4; i++)
        {
            if (pthread_join(th[i], NULL) != 0)
            {
                perror("Error joining thread");
            }
        }
        // Distributor thread
        close(fd[1]);

        while (1)
        {
            product prod;

            // Using lock and condition variable around crit section to avoid race condition
            // pthread_mutex_lock(&lock);
            // pthread_cond_wait(&cond, &lock);
            // Read from the pipe
            read(fd[0], &prod, sizeof(prod));
            if (prod.productType == 1)
            {
                put(&q1, prod);
            }
            else
            {
                put(&q2, prod);
            }
        }
        // pthread_cond_signal(&amp;cond);
        // pthread_mutex_unlock(&amp;lock);
        // Close read end of the pipe
        close(fd[0]);
    }
    return 0;
}

// Creating the producers
void producer(int args)
{
    int prodCount = 0;
    product prod;
    prod.productType = args;

    // Close read end of the pipe
    close(fd[0]);

    while (1)
    {
        prodCount++;
        prod.count = prodCount;
        // Send product to the pipe so the consumer can use
        write(fd[1], &prod, sizeof(prod));
        // Sleep for 0.01 - 0.2 seconds after each loop
        int time = (rand() % (200000 - 10000 + 1)) + 10000;
        usleep(time);
    }

    // Close write end of the pipe
    close(fd[1]);
}

void *consumer(void *args)
{
    int consCount1;
    int consCount2;

    FILE *fp;
    fp = fopen("output.txt", "w");

    product prod;
    int prodType = *(int *)args;

    while (1)
    {
        if (prodType == 1)
        {
            get(&q1);
            consCount1++;
            fputs("Thread ID: \n", fp);
            fprintf(fp, "Product Type: %d\n", prod.productType);
            fprintf(fp, "Production Sequence #: %d\n", prod.count);
            fprintf(fp, "Consumption Sequence #: %d\n", consCount1);
        }
        else
        {
            get(&q2);
            consCount2++;
            fputs("Thread ID: 2\n", fp);
            fprintf(fp, "Product Type: %d\n", prod.productType);
            fprintf(fp, "Production Sequence #: %d\n", prod.count);
            fprintf(fp, "Consumption Sequence #: %d\n", consCount2);
        }
    }
    fclose(fp);
}

// Circular buffer
void initQ(queue *q, int size)
{
    q->size = size;
    q->values = malloc(sizeof(product) * q->size);
    q->numEntries = 0;
    q->head = NULL;
    q->tail = NULL;
}

// Checks if the queue is empty
bool QEmpty(queue *q)
{
    return (q->numEntries == 0);
}

// Checks if the queue is full
bool QFull(queue *q)
{
    return (q->numEntries == q->size);
}

// Used for adding products to the queue
bool put(queue *q, product prod)
{
    // If the queue is full we can not add to it
    if (QFull(q))
    {
        return false;
    }
    // Add product to the end of the queue
    q->values[q->tail] = prod;
    q->numEntries++;
    // Move index of the tail
    q->tail = (q->tail + 1);
    // If index goes out of bounds set back to 0
    if (q->tail >= q->size)
    {
        q->tail = 0;
    }
    return true;
}

// Used for removing products for the queue
product get(queue *q)
{
    product result;

    // If the queue is empty we can not dequeue anymore
    if (QEmpty(q))
    {
        perror("Error on dequeue");
    }

    // Remove from the head of the queue
    result = q->values[q->head];
    q->head = (q->head + 1) & q->size;
    q->numEntries--;

    return result;
}

问题分析与修复方案

你的代码存在同步逻辑错误、线程ID分配错误、文件写入冲突等多个问题,以下是修复后的完整代码及关键修改说明:

修复后的代码

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <stdbool.h>
#include <fcntl.h>
#include <pthread.h>
#include <sys/wait.h>

#define BUFFER_SIZE1 20
#define BUFFER_SIZE2 30

typedef struct
{
    int count;          // 产品生产序列号
    int productType;    // 产品类型1/2
} product;

int fd[2];

// 每个队列独立配备锁和条件变量,避免跨队列同步干扰
typedef struct
{
    product *values;
    int head;
    int tail;
    int numEntries;
    int size;
    pthread_mutex_t lock;
    pthread_cond_t not_empty;
    pthread_cond_t not_full;
} 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]; // 4个消费者线程 + 1个分发线程
    if (pipe(fd) == -1)
    {
        perror("error creating pipe");
        exit(1);
    }

    // 初始化两个产品队列
    initQ(&q1, BUFFER_SIZE1);
    initQ(&q2, BUFFER_SIZE2);

    // 创建两个生产者进程
    int pid1 = fork();
    if (pid1 == 0)
    {
        producer(1);
        exit(0);
    }

    int pid2 = fork();
    if (pid2 == 0)
    {
        producer(2);
        exit(0);
    }

    // 创建分发线程
    if (pthread_create(&th[0], NULL, &distributor, NULL) != 0)
    {
        perror("Error creating distributor thread");
    }

    // 创建4个消费者线程,分配唯一ID(1-4)
    int consIds[4] = {1,2,3,4};
    for (int i = 0; i < 4; i++)
    {
        // 线程1、2消费类型1,线程3、4消费类型2
        int prodType = (i < 2) ? 1 : 2;
        // 动态分配参数,传递线程ID和产品类型
        int *threadArgs = malloc(sizeof(int)*2);
        threadArgs[0] = consIds[i];
        threadArgs[1] = prodType;
        
        if (pthread_create(&th[i+1], NULL, &consumer, threadArgs) != 0)
        {
            perror("Error creating consumer thread");
            free(threadArgs);
        }
    }

    // 等待所有线程结束
    for (int i = 0; i < 5; i++)
    {
        pthread_join(th[i], NULL);
    }

    // 等待生产者子进程结束,避免僵尸进程
    waitpid(pid1, NULL, 0);
    waitpid(pid2, NULL, 0);

    close(fd[0]);
    close(fd[1]);
    return 0;
}

void producer(int args)
{
    int prodCount = 0;
    product prod;
    prod.productType = args;

    close(fd[0]); // 关闭管道读端

    srand(getpid()); // 每个进程使用不同随机种子
    while (1)
    {
        prodCount++;
        prod.count = prodCount;
        // 管道满时write会自动阻塞,无需额外处理
        write(fd[1], &prod, sizeof(prod));
        // 随机休眠0.01-0.2秒
        int time = (rand() % (200000 - 10000 + 1)) + 10000;
        usleep(time);
    }

    close(fd[1]);
}

void *distributor(void *args)
{
    close(fd[1]); // 关闭管道写端
    product prod;

    while (1)
    {
        read(fd[0], &prod, sizeof(prod));
        // 根据产品类型放入对应队列
        if (prod.productType == 1)
        {
            put(&q1, prod);
        }
        else
        {
            put(&q2, prod);
        }
    }
    return NULL;
}

void *consumer(void *args)
{
    int threadId = ((int*)args)[0];
    int prodType = ((int*)args)[1];
    free(args); // 释放动态分配的参数

    int consCount = 0;
    // 以追加模式打开文件,避免覆盖已有内容
    FILE *fp = fopen("output.txt", "a");
    if (!fp)
    {
        perror("Failed to open output file");
        pthread_exit(NULL);
    }

    while (1)
    {
        product prod;
        // 从对应队列获取产品
        if (prodType == 1)
        {
            prod = get(&q1);
        }
        else
        {
            prod = get(&q2);
        }
        consCount++;

        // 加锁保证多线程写入文件时内容不混乱
        pthread_mutex_lock(&q1.lock);
        fprintf(fp, "Thread ID: %d\n", threadId);
        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); // 立即刷新缓冲区,确保内容写入文件
        pthread_mutex_unlock(&q1.lock);

        // 模拟消费耗时
        usleep(rand()%100000);
    }

    fclose(fp);
    return 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;
    // 初始化队列专属的锁和条件变量
    pthread_mutex_init(&q->lock, NULL);
    pthread_cond_init(&q->not_empty, NULL);
    pthread_cond_init(&q->not_full, NULL);
}

bool QEmpty(queue *q)
{
    return (q->numEntries == 0);
}

bool QFull(queue *q)
{
    return (q->numEntries == q->size);
}

bool put(queue *q, product prod)
{
    pthread_mutex_lock(&q->lock);
    // 队列满时等待消费者取走产品
    while (QFull(q))
    {
        pthread_cond_wait(&q->not_full, &q->lock);
    }

    q->values[q->tail] = prod;
    q->numEntries++;
    q->tail = (q->tail + 1) % q->size; // 修正循环队列索引逻辑

    // 通知消费者队列非空
    pthread_cond_signal(&q->not_empty);
    pthread_mutex_unlock(&q->lock);
    return true;
}

product get(queue *q)
{
    pthread_mutex_lock(&q->lock);
    // 队列空时等待生产者放入产品
    while (QEmpty(q))
    {
        pthread_cond_wait(&q->not_empty, &q->lock);
    }

    product result = q->values[q->head];
    q->numEntries--;
    q->head = (q->head + 1) % q->size; // 修正循环队列索引逻辑

    // 通知生产者队列非满
    pthread_cond_signal(&q->not_full);
    pthread_mutex_unlock(&q->lock);
    return result;
}

关键修改说明

  1. 队列同步机制:每个队列独立配置锁和条件变量,put操作在队列满时等待,get操作在队列空时等待,彻底解决多线程竞争和忙等问题。
  2. 主线程逻辑调整:先创建所有线程再调用pthread_join,确保分发线程和消费者线程能同时运行,避免原代码中分发逻辑永远无法执行的问题。
  3. 线程ID分配:通过动态参数传递每个线程的唯一ID(1-4),满足需求中消费者ID固定的要求。
  4. 文件写入优化:用追加模式打开输出文件,加锁保证多线程写入时内容不混乱,同时强制刷新缓冲区确保内容及时落地。
  5. 循环队列逻辑修复:将原代码中的位运算错误改为取模运算,修正head和tail的初始化值,符合循环队列的
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 23:55:19