如何用rangeless/STL实现含中间串行联合计算的并行任务流?
实现J与部分B并行的异步任务流程(STL/轻量库方案)
完全可以用STL的并发组件(std::future、std::shared_future、条件变量等)或rangeless库实现你想要的并行优化,核心思路是让J无需等待所有A任务完成再启动,而是边接收已完成的A结果边计算,同时让完成A的线程在J结果就绪后立即启动B,实现J与剩余A、部分B的并行执行。
核心实现逻辑
- 工作线程启动后先执行A任务,完成后将结果传递给J计算线程(用线程安全队列或
std::future) - J线程在单独线程中增量处理A结果,和未完成的A任务并行执行
- 每个工作线程完成A后,等待J的结果(用
std::shared_future让多线程共享J的结果),一旦J就绪就立即启动B→C - 主线程最后收集所有C结果,执行L任务
STL代码示例
#include <vector> #include <future> #include <mutex> #include <condition_variable> #include <iostream> #include <chrono> // 定义任务结果类型 struct AResult { int id; }; struct JResult { int total; }; struct BResult { int product; }; struct CResult { int square; }; // 模拟耗时的A任务 AResult taskA(int worker_id) { std::this_thread::sleep_for(std::chrono::milliseconds(100 * worker_id)); return {worker_id}; } // J任务:增量计算所有A结果的总和,和A并行执行 JResult taskJ(const int num_workers, std::vector<AResult>& a_queue, std::mutex& mtx, std::condition_variable& cv, bool& all_a_done) { JResult res{0}; int processed = 0; while (processed < num_workers) { std::unique_lock<std::mutex> lock(mtx); cv.wait(lock, [&]() { return !a_queue.empty() || all_a_done; }); if (!a_queue.empty()) { auto a_res = a_queue.back(); a_queue.pop_back(); res.total += a_res.id; processed++; std::cout << "J processed A from worker " << a_res.id << ", current total: " << res.total << "\n"; } else if (all_a_done) { break; } } return res; } // B任务:依赖J的结果 BResult taskB(const AResult& a_res, const JResult& j_res) { std::this_thread::sleep_for(std::chrono::milliseconds(50)); return {a_res.id * j_res.total}; } // C任务:依赖B的结果 CResult taskC(const BResult& b_res) { std::this_thread::sleep_for(std::chrono::milliseconds(30)); return {b_res.product * b_res.product}; } // L任务:汇总所有C结果 void taskL(const std::vector<CResult>& c_results) { int sum = 0; for (const auto& c : c_results) sum += c.square; std::cout << "L task completed, total sum: " << sum << "\n"; } int main() { const int worker_count = 5; std::vector<AResult> a_result_queue; std::mutex queue_mtx; std::condition_variable cv; bool all_a_finished = false; // 启动J计算线程 auto j_future = std::async(std::launch::async, taskJ, worker_count, std::ref(a_result_queue), std::ref(queue_mtx), std::ref(cv), std::ref(all_a_finished)); std::shared_future<JResult> shared_j = j_future.share(); // 让所有工作线程共享J的结果 // 启动所有工作线程 std::vector<std::future<CResult>> c_futures; for (int i = 0; i < worker_count; ++i) { c_futures.emplace_back(std::async(std::launch::async, [i, &queue_mtx, &cv, &a_result_queue, shared_j]() { // 执行A任务 AResult a_res = taskA(i); std::cout << "Worker " << i << " finished A\n"; // 将A结果传递给J { std::lock_guard<std::mutex> lock(queue_mtx); a_result_queue.push_back(a_res); } cv.notify_one(); // 等待J结果就绪,然后执行B→C JResult j_res = shared_j.get(); std::cout << "Worker " << i << " got J result, starting B\n"; BResult b_res = taskB(a_res, j_res); std::cout << "Worker " << i << " finished B\n"; CResult c_res = taskC(b_res); std::cout << "Worker " << i << " finished C\n"; return c_res; })); } // 通知J所有A任务已启动完成,无更多结果可接收 { std::lock_guard<std::mutex> lock(queue_mtx); all_a_finished = true; } cv.notify_one(); // 收集所有C结果,执行L std::vector<CResult> c_results; for (auto& fut : c_futures) { c_results.push_back(fut.get()); } taskL(c_results); return 0; }
优化点说明
- 如果J仅依赖部分A结果(比如前2个A),可以修改
taskJ中的processed < num_workers为processed < required_a_count,这样J提前完成,已完成A的线程可以立即启动B,和剩余的A任务并行,进一步提升资源利用率 - 使用
std::shared_future避免J结果被多次拷贝,同时支持多线程等待 - 线程安全队列+条件变量保证J能及时接收A结果,避免忙等
Rangeless库简化实现
Rangeless可以用range语法简化异步数据流的处理,核心逻辑和STL方案一致,只是代码更简洁:
#include <rangeless.hpp> #include <future> #include <vector> // 复用前面定义的taskA/taskB/taskC/taskL和结果类型 int main() { const int worker_count = 5; // 生成A任务的异步流 auto a_stream = rl::generate([id = 0, worker_count]() mutable -> std::optional<std::future<AResult>> { if (id >= worker_count) return std::nullopt; return std::async(std::launch::async, taskA, id++); }); // J任务消费A流,增量计算总和 JResult j_res = rl::fold( a_stream | rl::transform([](auto&& fut) { return fut.get(); }), JResult{0}, [](JResult acc, AResult a) { acc.total += a.id; return acc; } ); // 启动工作线程执行A→等待J→B→C std::vector<std::future<CResult>> c_futures; for (int i = 0; i < worker_count; ++i) { c_futures.emplace_back(std::async(std::launch::async, [i, j_res]() { AResult a_res = taskA(i); BResult b_res = taskB(a_res, j_res); return taskC(b_res); })); } // 收集结果执行L std::vector<CResult> c_results; rl::for_each(c_futures, [&](auto&& fut) { c_results.push_back(fut.get()); }); taskL(c_results); return 0; }
内容的提问来源于stack exchange,提问作者Sir Nate
相关产品推荐
相关产品推荐

