如何统计生产者向缓冲区交付字符的轮次?
生产者-消费者程序轮次统计需求与实现
需求说明
现有一个生产者-消费者程序,功能为逐字符读取文件并将内容存入缓冲区。需要统计生产者向缓冲区交付字符的轮次,轮次定义为:一次或多次连续写入缓冲区且未因队列满而被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
相关产品推荐
相关产品推荐

