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

如何用rangeless/STL实现含中间串行联合计算的并行任务流?

实现J与部分B并行的异步任务流程(STL/轻量库方案)

完全可以用STL的并发组件(std::future、std::shared_future、条件变量等)或rangeless库实现你想要的并行优化,核心思路是让J无需等待所有A任务完成再启动,而是边接收已完成的A结果边计算,同时让完成A的线程在J结果就绪后立即启动B,实现J与剩余A、部分B的并行执行。

核心实现逻辑

  1. 工作线程启动后先执行A任务,完成后将结果传递给J计算线程(用线程安全队列或std::future)
  2. J线程在单独线程中增量处理A结果,和未完成的A任务并行执行
  3. 每个工作线程完成A后,等待J的结果(用std::shared_future让多线程共享J的结果),一旦J就绪就立即启动B→C
  4. 主线程最后收集所有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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 23:45:54