POSIX多线程编程作业求助:动态任务分配与线程同步
问题描述
我有一项学校作业需要完成,核心要求如下:
- 作业目标:通过创建动态数量的线程并合理分配任务,理解多线程编程;需使用checkpointing和锁来协调线程进度,必须用C语言和POSIX线程库实现。
- 程序参数:从控制台接收元数据文件路径和最大可用线程数。
- 元数据文件格式:包含各输入文件信息及计算所需数值,示例:
./input_files/file_1.txt 0.5 1 ./input_files/file_2.txt 0.5 0.5
- 数据结构:需读取元数据文件,创建
InputFile结构体数组,每个结构体对应一个输入文件,定义如下:
typedef struct InputFile { char file_path[1024]; float number_1; float number_2; } InputFile;
- 处理逻辑:每个输入文件包含若干数值,需为每个文件开启线程处理;处理完成后合并所有文件的输出为最终结果并打印,示例输出:
File_1 10 20 30 File_2 50 70 90 The output should be like this: 60 >> (50 + 10) 90 >> (70 + 20) 120 >> (30 + 90)
- 特殊场景:必须处理线程数与文件数不匹配的情况(如5文件2线程或5线程2文件)。
我已完成的部分:成功获取元数据文件和线程数,创建了InputFile_array数组存储所有InputFile结构体。当线程数与文件数相等时,我采用的方案如下:
// 存储线程数为"num_threads",从命令行参数获取 InputFile InputFiles[num_threads]; // 假设文件数和线程数相等,以此作为数组大小 // ... // 程序填充InputFiles数组,包含每个文件的信息 // ... pthread_t threads[num_threads]; // 创建每个线程并传递对应的结构体(1线程处理1个文件) for (size_t index = 0; index < num_threads; index++) { if (pthread_create(&threads[index], NULL, &thread_function, InputFiles[index]) != 0) { perror("Failed to create thread"); exit(1); // 出错则终止程序 } } // 等待所有线程完成 for (size_t index = 0; index < num_threads; index++) { if (pthread_join(threads[index], NULL) != 0) { exit(2); } } printf("Entire process has finished.\n");
现在需要解决的问题:在可使用mutex锁和信号量的前提下,如何针对线程数与文件数不匹配的场景,设计合理的任务分配与线程同步方案?
解决方案
针对线程数与文件数不匹配的场景,推荐两种实用的任务分配方案,结合POSIX锁/信号量实现同步:
方案一:共享任务队列(动态分配)
这是最灵活的方案,适合任意线程数与文件数的组合,核心是让线程从共享队列中取未处理的文件任务,直到所有任务完成。
实现步骤
- 定义共享任务管理结构:包含待处理文件的索引队列、互斥锁、任务完成计数(用于checkpointing):
#include <pthread.h> #include <semaphore.h> // 共享任务管理器 typedef struct TaskManager { InputFile* input_files; // 所有输入文件数组 int total_files; // 文件总数 int current_index; // 下一个待处理的文件索引 pthread_mutex_t mutex; // 保护共享队列的互斥锁 sem_t task_sem; // 信号量:标识可处理的任务数 int completed_count; // 已完成的任务数(checkpoint用) pthread_mutex_t count_mutex; // 保护completed_count的锁 } TaskManager;
- 初始化共享资源:在主线程中完成初始化,设置初始任务数为文件总数:
TaskManager task_manager; task_manager.input_files = InputFile_array; // 已创建的InputFile数组 task_manager.total_files = total_files; // 实际文件总数 task_manager.current_index = 0; pthread_mutex_init(&task_manager.mutex, NULL); sem_init(&task_manager.task_sem, 0, total_files); // 初始信号量值为文件数 task_manager.completed_count = 0; pthread_mutex_init(&task_manager.count_mutex, NULL);
- 线程函数逻辑:线程循环获取任务,处理完成后更新计数,直到所有任务处理完毕:
void* thread_function(void* arg) { TaskManager* manager = (TaskManager*)arg; int file_index; while (1) { // 等待可用任务 sem_wait(&manager->task_sem); // 获取下一个待处理文件的索引 pthread_mutex_lock(&manager->mutex); if (manager->current_index >= manager->total_files) { pthread_mutex_unlock(&manager->mutex); break; // 所有任务已处理,退出线程 } file_index = manager->current_index++; pthread_mutex_unlock(&manager->mutex); // 处理当前文件:读取文件数值、计算(示例逻辑) InputFile* file = &manager->input_files[file_index]; // ... 这里写你的文件读取和计算逻辑,比如生成该文件的输出数组 ... // 更新完成任务计数(checkpointing) pthread_mutex_lock(&manager->count_mutex); manager->completed_count++; // 可选:打印checkpoint信息,比如"已完成第X个文件,共Y个" printf("Checkpoint: %d/%d files processed\n", manager->completed_count, manager->total_files); pthread_mutex_unlock(&manager->count_mutex); } return NULL; }
- 主线程控制逻辑:创建指定数量的线程,等待所有线程完成,然后合并结果:
int num_threads = ...; // 命令行传入的最大线程数 // 实际创建的线程数取线程数和文件数的较小值(如果线程数多于文件数,没必要创建多余线程) int actual_threads = num_threads > total_files ? total_files : num_threads; pthread_t threads[actual_threads]; // 创建线程 for (int i = 0; i < actual_threads; i++) { if (pthread_create(&threads[i], NULL, thread_function, &task_manager) != 0) { perror("Failed to create thread"); exit(1); } } // 等待所有线程完成 for (int i = 0; i < actual_threads; i++) { pthread_join(threads[i], NULL); } // 合并所有文件的输出结果并打印 // ... 这里写你的结果合并逻辑,按照示例输出格式处理 ... // 销毁同步资源 pthread_mutex_destroy(&task_manager.mutex); sem_destroy(&task_manager.task_sem); pthread_mutex_destroy(&task_manager.count_mutex);
方案二:静态分片分配(固定任务块)
如果不需要动态调整任务,可提前给每个线程分配固定数量的文件,适合文件数较多的场景。
实现步骤
- 计算每个线程的任务范围:
int total_files = ...; // 文件总数 int num_threads = ...; // 传入的线程数 int actual_threads = num_threads > total_files ? total_files : num_threads; // 计算每个线程处理的文件数,剩余文件分配给前几个线程 int files_per_thread = total_files / actual_threads; int remainder = total_files % actual_threads;
- 定义线程参数结构:包含该线程负责的文件起始、结束索引,以及输入文件数组:
typedef struct ThreadArgs { InputFile* input_files; int start_index; int end_index; int thread_id; } ThreadArgs;
- 线程函数逻辑:处理分配给自己的所有文件,完成后更新共享计数(checkpointing):
// 全局或共享的完成计数和锁(也可放入ThreadArgs) int completed_count = 0; pthread_mutex_t count_mutex; void* thread_function(void* arg) { ThreadArgs* args = (ThreadArgs*)arg; for (int i = args->start_index; i < args->end_index; i++) { InputFile* file = &args->input_files[i]; // ... 处理当前文件 ... // 更新完成计数 pthread_mutex_lock(&count_mutex); completed_count++; printf("Checkpoint: %d/%d files processed\n", completed_count, total_files); pthread_mutex_unlock(&count_mutex); } free(args); // 释放参数内存 return NULL; }
- 主线程创建线程:
pthread_mutex_init(&count_mutex, NULL); pthread_t threads[actual_threads]; for (int i = 0; i < actual_threads; i++) { ThreadArgs* args = malloc(sizeof(ThreadArgs)); args->input_files = InputFile_array; args->start_index = i * files_per_thread + (i < remainder ? i : remainder); args->end_index = args->start_index + files_per_thread + (i < remainder ? 1 : 0); args->thread_id = i; if (pthread_create(&threads[i], NULL, thread_function, args) != 0) { perror("Failed to create thread"); exit(1); } } // 等待所有线程完成 for (int i = 0; i < actual_threads; i++) { pthread_join(threads[i], NULL); } // 合并结果、销毁资源 // ... pthread_mutex_destroy(&count_mutex);
关键同步与Checkpointing说明
- 互斥锁的作用:保护共享变量(如任务索引、完成计数),避免多线程同时修改导致的数据竞争。
- 信号量的作用:在动态队列方案中,控制线程获取任务的节奏,避免线程空转或重复处理任务。
- Checkpointing实现:通过保护的
completed_count变量,每次任务完成后更新并打印进度,确保进度信息的准确性(不会出现计数混乱)。 - 线程数优化:当线程数多于文件数时,只创建与文件数相等的线程,避免不必要的线程开销。
内容的提问来源于stack exchange,提问作者Default
相关产品推荐
相关产品推荐

