多线程并发文件读写:维持固定线程数处理任务的实现问题
C++固定并发数的文件批量处理线程实现
我是多线程技术新手,正在编写一个C++程序:输入为包含大文本文件路径的OBJ对象向量,以及指定并发线程数n。要求程序始终维持n个线程同时运行以最大化效率,直到所有文件处理完成。
已实现文件处理函数proc_file,需要完善multi_file_proc的线程管理逻辑。
现有代码实现
// 读取文件并生成过滤后的新版本 void proc_file(OBJ obj) { std::string inFileStr(obj.get_path().c_str()); std::string outFileStr(std::string(obj.get_path().replace_extension("new.txt").c_str())); std::ifstream inFile(inFileStr); std::ofstream outFile(outFileStr); std::string currLine; while (getline(inFile, currLine)) { if (currLine.size() == 1 || currLine.compare(currLine.length()-5, 5, "thing") != 0) { outFile << currLine << '\n'; } else { for (int i = 0; i < 3; i++) { getline(inFile, currLine); } } } inFile.close(); outFile.close(); } // 需要完善的多文件处理函数:维持n个并发线程处理所有OBJ void multi_file_proc(std::vector<OBJ> objs, int n) { std::vector<std::thread> procVec; for (int i = 0; i < objs.size(); i++) { /* 需实现逻辑:始终保持n个线程运行, 一个线程完成后立即启动新的,直到所有文件处理完毕 */ } }
尝试过的错误方案
错误方案1:假设头部线程最先完成
代码实现:
if (i >= n) { procVec.front().join(); procVec.erase(procVec.begin()); } procVec.push_back(std::thread(proc_file, objs[i]));
问题:错误假设向量头部的线程会最先完成,实际线程完成顺序不确定,可能导致主线程等待未完成的线程,浪费并发资源;同时erase操作可能引发迭代器失效问题。
错误方案2:全局变量+线程自清理
代码实现:
// 全局变量 std::vector<std::thread> procVec; std::mutex threadMutex; void remThread(std::thread::id id) { std::lock_guard<std::mutex> lock(threadMutex); auto iter = std::find_if(procVec.begin(), procVec.end(), [=](std::thread &t) {return (t.get_id() == id); }); if (iter != procVec.end()) { iter->join(); procVec.erase(iter); } } void lamb(OBJ obj, std::thread::id id) { proc_file(obj); remThread(id); } // multi_file_proc中的循环代码 std::lock_guard<std::mutex> lock(threadMutex); procVec.push_back(std::thread([&objs, i]() { std::thread(lamb, objs[i], std::this_thread::get_id()).detach(); }));
问题:存在未被join的线程,程序退出时会触发异常;嵌套detach的线程无法被正确跟踪,导致资源泄漏。
正确实现方案
方案1:任务队列+固定工作线程池(推荐)
这种方式复用线程,避免频繁创建销毁线程,且能严格维持n个并发数:
#include <queue> #include <mutex> #include <condition_variable> #include <atomic> void multi_file_proc(std::vector<OBJ> objs, int n) { // 任务队列:存储待处理的OBJ std::queue<OBJ> task_queue; for (const auto& obj : objs) { task_queue.push(obj); } std::mutex queue_mutex; std::condition_variable cv; std::atomic<bool> stop_flag{false}; // 标记是否停止工作线程 std::vector<std::thread> workers; workers.reserve(n); // 创建n个工作线程 for (int i = 0; i < n; ++i) { workers.emplace_back([&]() { while (true) { OBJ current_obj; { std::unique_lock<std::mutex> lock(queue_mutex); // 等待任务或停止信号 cv.wait(lock, [&]() { return stop_flag || !task_queue.empty(); }); if (stop_flag && task_queue.empty()) { break; // 无任务且收到停止信号,退出线程 } current_obj = std::move(task_queue.front()); task_queue.pop(); } // 处理文件 proc_file(current_obj); } }); } // 等待所有任务被取走 { std::unique_lock<std::mutex> lock(queue_mutex); cv.wait(lock, [&]() { return task_queue.empty(); }); } // 发送停止信号并唤醒所有工作线程 stop_flag = true; cv.notify_all(); // 等待所有工作线程结束 for (auto& worker : workers) { if (worker.joinable()) { worker.join(); } } }
方案说明
- 任务队列:集中管理所有待处理的文件任务,工作线程从队列中取任务执行。
- 条件变量+互斥锁:安全同步任务队列的访问,工作线程空闲时等待新任务或停止信号,避免空轮询浪费CPU。
- 原子停止标记:确保线程安全接收停止信号,所有任务处理完成后,主线程通知工作线程退出。
- 线程复用:固定创建n个工作线程,直到所有任务完成后才销毁,避免频繁创建销毁线程的开销。
方案2:基于std::future的动态任务管理(简化版)
如果不需要线程复用,可使用std::future跟踪任务状态,动态替换完成的任务:
#include <future> #include <algorithm> void multi_file_proc(std::vector<OBJ> objs, int n) { std::vector<std::future<void>> futures; int task_idx = 0; // 先启动n个初始任务 for (; task_idx < n && task_idx < objs.size(); ++task_idx) { futures.emplace_back(std::async(std::launch::async, proc_file, objs[task_idx])); } // 循环替换完成的任务,直到所有任务启动 while (task_idx < objs.size()) { // 找到第一个完成的任务 auto it = std::find_if(futures.begin(), futures.end(), [](std::future<void>& f) { return f.wait_for(std::chrono::milliseconds(0)) == std::future_status::ready; }); if (it != futures.end()) { it->get(); // 等待任务完成(非阻塞,因为已ready) // 替换为新任务 *it = std::async(std::launch::async, proc_file, objs[task_idx++]); } } // 等待剩余所有任务完成 for (auto& f : futures) { f.get(); } }
方案说明
- 使用
std::async创建异步任务,返回std::future跟踪任务状态。 - 主线程循环检查并替换完成的任务,维持n个并发任务。
- 最后等待所有剩余任务完成,确保无残留线程。
内容的提问来源于stack exchange,提问作者gladshire
相关产品推荐
相关产品推荐

