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

基于进程线程同步的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 19:25:20