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

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锁/信号量实现同步:

方案一:共享任务队列(动态分配)

这是最灵活的方案,适合任意线程数与文件数的组合,核心是让线程从共享队列中取未处理的文件任务,直到所有任务完成。

实现步骤

  1. 定义共享任务管理结构:包含待处理文件的索引队列、互斥锁、任务完成计数(用于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;
  1. 初始化共享资源:在主线程中完成初始化,设置初始任务数为文件总数:
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);
  1. 线程函数逻辑:线程循环获取任务,处理完成后更新计数,直到所有任务处理完毕:
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;
}
  1. 主线程控制逻辑:创建指定数量的线程,等待所有线程完成,然后合并结果:
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);

方案二:静态分片分配(固定任务块)

如果不需要动态调整任务,可提前给每个线程分配固定数量的文件,适合文件数较多的场景。

实现步骤

  1. 计算每个线程的任务范围:
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;
  1. 定义线程参数结构:包含该线程负责的文件起始、结束索引,以及输入文件数组:
typedef struct ThreadArgs {
    InputFile* input_files;
    int start_index;
    int end_index;
    int thread_id;
} ThreadArgs;
  1. 线程函数逻辑:处理分配给自己的所有文件,完成后更新共享计数(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;
}
  1. 主线程创建线程:
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 16:14:57