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

如何统计生产者向缓冲区交付字符的轮次?

生产者-消费者程序轮次统计需求与实现

需求说明

现有一个生产者-消费者程序,功能为逐字符读取文件并将内容存入缓冲区。需要统计生产者向缓冲区交付字符的轮次,轮次定义为:一次或多次连续写入缓冲区且未因队列满而被wait中断的过程。

原程序代码

#include <pthread.h>
#include <semaphore.h>
#include <stdlib.h>
#include <stdio.h>

/*
This program provides a possible solution for producer-consumer problem using mutex and semaphore.
I have used 5 producers and 5 consumers to demonstrate the solution. You can always play with these values.
*/

#define MaxItems 5 // Maximum items a producer can produce or a consumer can consume
#define BufferSize 5 // Size of the buffer

sem_t empty;
sem_t full;
int in = 0;
int out = 0;
int buffer[BufferSize];
pthread_mutex_t mutex;

void *producer(void *pno)
{   
    int item;
    for(int i = 0; i < MaxItems; i++) {
        item = rand(); // Produce an random item
        sem_wait(&empty);
        pthread_mutex_lock(&mutex);
        buffer[in] = item;
        printf("Producer %d: Insert Item %d at %d\n", *((int *)pno),buffer[in],in);
        in = (in+1)%BufferSize;
        pthread_mutex_unlock(&mutex);
        sem_post(&full);
    }
}
void *consumer(void *cno)
{   
    for(int i = 0; i < MaxItems; i++) {
        sem_wait(&full);
        pthread_mutex_lock(&mutex);
        int item = buffer[out];
        printf("Consumer %d: Remove Item %d from %d\n",*((int *)cno),item, out);
        out = (out+1)%BufferSize;
        pthread_mutex_unlock(&mutex);
        sem_post(&empty);
    }
}

int main()
{   

    pthread_t pro[5],con[5];
    pthread_mutex_init(&mutex, NULL);
    sem_init(&empty,0,BufferSize);
    sem_init(&full,0,0);

     FILE *fp = fopen("file.txt", "r");
     if (fp != NULL) {

        if (fseek(fp, 0L, SEEK_END) == 0) {
            /* Get the size of the file. */
            p1.BUFFER_SIZE = ftell(fp);
            if (p1.BUFFER_SIZE == -1) { /* Error */ }

            /* Allocate our buffer to that size. */
            p1.item = malloc(sizeof(char) * (p1.BUFFER_SIZE + 1));

            /* Go back to the start of the file. */
            if (fseek(fp, 0L, SEEK_SET) != 0) { /* Error */ }

            /* Read the entire file into memory. */
            size_t newLen = fread(p1.item, sizeof(char), p1.BUFFER_SIZE, fp);
            if ( ferror( fp ) != 0 ) {
                fputs("Error reading file", stderr);
            } else {
                p1.item[newLen++] = '\0'; /* Just to be safe. */
            }
        }

    int a[5] = {1,2,3,4,5}; //Just used for numbering the producer and consumer

    for(int i = 0; i < 5; i++) {
        pthread_create(&pro[i], NULL, (void *)producer, (void *)&a[i]);
    }
    for(int i = 0; i < 5; i++) {
        pthread_create(&con[i], NULL, (void *)consumer, (void *)&a[i]);
    }

    for(int i = 0; i < 5; i++) {
        pthread_join(pro[i], NULL);
    }
    for(int i = 0; i < 5; i++) {
        pthread_join(con[i], NULL);
    }

    pthread_mutex_destroy(&mutex);
    sem_destroy(&empty);
    sem_destroy(&full);

    return 0;
    
}

轮次统计实现方案

1. 核心思路

轮次的关键判断依据:生产者每次写入前,若因缓冲区满被sem_wait阻塞,则唤醒后的写入属于新轮次;若未阻塞且已有连续写入,则属于当前轮次。通过sem_trywait提前判断是否需要阻塞,以此区分轮次边界。

2. 代码修改与实现

步骤1:添加必要变量与结构体

#include <pthread.h>
#include <semaphore.h>
#include <stdlib.h>
#include <stdio.h>
#include <string.h>
#include <errno.h>

#define BufferSize 5 // 缓冲区大小

sem_t empty;
sem_t full;
int in = 0;
int out = 0;
char buffer[BufferSize]; // 改为字符缓冲区适配需求
pthread_mutex_t mutex;
pthread_mutex_t round_mutex; // 保护全局轮次计数的锁

char *file_content = NULL;
long file_size = 0;
int global_round_count = 0; // 全局总轮次

// 生产者参数结构体,传递ID和自身轮次统计
typedef struct {
    int producer_id;
    int round_count;
    int first_write; // 标记是否为第一次写入
} ProducerData;

步骤2:修改生产者函数实现轮次统计

void *producer(void *args)
{   
    ProducerData *data = (ProducerData *)args;
    char item;

    // 分配每个生产者负责的文件字符范围(均分处理)
    long start = data->producer_id * (file_size / 5);
    long end = (data->producer_id + 1) * (file_size / 5);
    if (data->producer_id == 4) end = file_size; // 最后一个生产者处理剩余字符

    data->round_count = 0;
    data->first_write = 1;

    for(long i = start; i < end; i++) {
        item = file_content[i];
        int was_blocked = 0;

        // 尝试获取空槽,判断是否需要阻塞
        if (sem_trywait(&empty) != 0) {
            // 缓冲区已满,进入阻塞等待
            sem_wait(&empty);
            was_blocked = 1;
        }

        pthread_mutex_lock(&mutex);
        // 判断是否开启新轮次
        if (was_blocked || data->first_write) {
            data->round_count++;
            // 更新全局轮次,需加锁保护
            pthread_mutex_lock(&round_mutex);
            global_round_count++;
            pthread_mutex_unlock(&round_mutex);
            data->first_write = 0;
        }

        buffer[in] = item;
        printf("生产者 %d: 写入字符 '%c' 到位置 %d\n", data->producer_id, buffer[in], in);
        in = (in + 1) % BufferSize;
        pthread_mutex_unlock(&mutex);
        sem_post(&full);
    }

    printf("生产者 %d 完成任务,累计轮次:%d\n", data->producer_id, data->round_count);
    return NULL;
}

步骤3:调整消费者函数适配字符读取

void *consumer(void *cno)
{   
    int consumer_id = *((int *)cno);
    // 计算每个消费者需要处理的字符数
    long task_count = (file_size / 5) + (consumer_id == 5 ? file_size % 5 : 0);

    for(long i = 0; i < task_count; i++) {
        sem_wait(&full);
        pthread_mutex_lock(&mutex);
        char item = buffer[out];
        printf("消费者 %d: 读取字符 '%c' 从位置 %d\n", consumer_id, item, out);
        out = (out + 1) % BufferSize;
        pthread_mutex_unlock(&mutex);
        sem_post(&empty);
    }
    return NULL;
}

步骤4:修正main函数的文件读取与线程初始化

int main()
{   
    pthread_t pro[5], con[5];
    pthread_mutex_init(&mutex, NULL);
    pthread_mutex_init(&round_mutex, NULL);
    sem_init(&empty, 0, BufferSize);
    sem_init(&full, 0, 0);

    // 读取文件内容到内存
    FILE *fp = fopen("file.txt", "r");
    if (fp == NULL) {
        perror("打开文件失败");
        return 1;
    }

    if (fseek(fp, 0L, SEEK_END) == 0) {
        file_size = ftell(fp);
        if (file_size == -1) {
            perror("获取文件大小失败");
            fclose(fp);
            return 1;
        }

        file_content = malloc(sizeof(char) * (file_size + 1));
        if (!file_content) {
            perror("内存分配失败");
            fclose(fp);
            return 1;
        }

        if (fseek(fp, 0L, SEEK_SET) != 0) {
            perror("重置文件指针失败");
            free(file_content);
            fclose(fp);
            return 1;
        }

        size_t read_len = fread(file_content, sizeof(char), file_size, fp);
        if (ferror(fp) != 0) {
            fputs("读取文件错误", stderr);
            free(file_content);
            fclose(fp);
            return 1;
        } else {
            file_content[read_len] = '\0';
        }
    }
    fclose(fp);

    // 初始化生产者参数并创建线程
    ProducerData pro_data[5];
    for(int i = 0; i < 5; i++) {
        pro_data[i].producer_id = i + 1;
        pro_data[i].round_count = 0;
        pro_data[i].first_write = 1;
        pthread_create(&pro[i], NULL, producer, (void *)&pro_data[i]);
    }

    // 创建消费者线程
    int con_ids[5] = {1,2,3,4,5};
    for(int i = 0; i < 5; i++) {
        pthread_create(&con[i], NULL, consumer, (void *)&con_ids[i]);
    }

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

    printf("所有生产者完成,全局总轮次:%d\n", global_round_count);

    // 释放资源
    pthread_mutex_destroy(&mutex);
    pthread_mutex_destroy(&round_mutex);
    sem_destroy(&empty);
    sem_destroy(&full);
    free(file_content);

    return 0;
}

3. 逻辑说明

  • 通过sem_trywait非阻塞尝试获取空槽,若失败则说明缓冲区已满,生产者会进入阻塞,唤醒后写入的字符标记为新轮次。
  • 每个生产者的第一次写入无论是否阻塞,均算作第一个轮次。
  • 全局轮次计数使用互斥锁保护,避免多线程竞争导致计数错误。
  • 修正了原程序中p1未定义、缓冲区类型不匹配等问题,适配逐字符读取文件的需求。

内容的提问来源于stack exchange,提问作者beginner_101

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 03:10:44