C++中多线程拆分数据集计算:1500行2D向量的并行实现
多线程优化计算密集型Loss任务的实现方案
核心思路
将1500行数据集平均拆分给6个线程是最优方案——你的CPU是6核,每个线程对应一个核心,任务量均匀分配能避免负载不均,最大化利用硬件资源。1500刚好能被6整除,每个线程处理250行;如果数据集行数无法整除,最后一个线程多处理剩余行数即可。
实现方案(C++)
下面是基于std::thread的具体实现,重点用局部总和+最后汇总的方式避免线程锁开销(计算密集型任务中,锁会显著降低效率):
#include <vector> #include <thread> #include <numeric> // 建议修改Loss函数参数为const引用,避免不必要的拷贝 double Loss(const vector<double>& input, const vector<double>& expectedOutput); // 单线程实现(参考) double TotalLoss(const vector<vector<double>>& inputs, const vector<vector<double>>& expectedOutputs) { double total = 0.0; for (size_t i = 0; i < inputs.size(); ++i) { total += Loss(inputs[i], expectedOutputs[i]); } return total / inputs.size(); } // 多线程实现 double MultiThreadedTotalLoss(const vector<vector<double>>& inputs, const vector<vector<double>>& expectedOutputs) { const size_t num_threads = 6; const size_t total_rows = inputs.size(); vector<double> thread_sums(num_threads, 0.0); // 存储每个线程计算的局部总和 // 定义每个线程的任务函数 auto compute_range = [&](size_t thread_id) { size_t rows_per_thread = total_rows / num_threads; size_t start = thread_id * rows_per_thread; // 最后一个线程处理剩余所有行(如果不能整除) size_t end = (thread_id == num_threads - 1) ? total_rows : (thread_id + 1) * rows_per_thread; double local_sum = 0.0; for (size_t i = start; i < end; ++i) { local_sum += Loss(inputs[i], expectedOutputs[i]); } thread_sums[thread_id] = local_sum; }; // 创建并启动线程 vector<thread> threads; for (size_t i = 0; i < num_threads; ++i) { threads.emplace_back(compute_range, i); } // 等待所有线程完成 for (auto& t : threads) { t.join(); } // 汇总所有局部总和,计算平均Loss double total_sum = accumulate(thread_sums.begin(), thread_sums.end(), 0.0); return total_sum / total_rows; }
关键优化点
- 参数传递优化:把所有函数的vector参数改为
const引用,避免大内存的拷贝操作——原函数的传值方式会导致每次调用都复制整个vector,对性能影响极大。 - 无锁汇总:每个线程只写入自己对应的
thread_sums元素,不存在线程竞争,完全不需要互斥锁,避免了锁带来的性能损耗。 - 任务均匀分配:通过计算起始/结束索引拆分任务,保证每个线程的工作量尽可能一致,避免某几个线程提前空闲导致CPU资源浪费。
替代方案(更简洁)
如果不想手动管理线程,可以用std::async自动调度任务,代码更简洁:
#include <future> double MultiThreadedTotalLoss(const vector<vector<double>>& inputs, const vector<vector<double>>& expectedOutputs) { const size_t num_threads = 6; const size_t total_rows = inputs.size(); vector<future<double>> futures; for (size_t i = 0; i < num_threads; ++i) { size_t start = i * (total_rows / num_threads); size_t end = (i == num_threads - 1) ? total_rows : (i + 1) * (total_rows / num_threads); futures.emplace_back(std::async(std::launch::async, [&](size_t s, size_t e) { double sum = 0.0; for (size_t j = s; j < e; ++j) { sum += Loss(inputs[j], expectedOutputs[j]); } return sum; }, start, end)); } double total_sum = 0.0; for (auto& f : futures) { total_sum += f.get(); } return total_sum / total_rows; }
内容的提问来源于stack exchange,提问作者Soumish Das
相关产品推荐
相关产品推荐

