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

多线程并发文件读写:维持固定线程数处理任务的实现问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 03:25:19