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

如何修改C++中带Continuation的Zip实现以支持并行执行?

Solution for Parallel Zip with Continuations in C++

Great question! The core issue you're facing is safely storing the two values from your upstream and other sources while handling potential asynchronous execution, and ensuring we only invoke the final callback once both values are available. Let's walk through a practical solution step by step.

Step 1: Add Value Type Annotations

First, we need to explicitly declare the type of value each F-derived struct produces. This lets us know exactly what types we'll be storing in our zip operation. Update your F_id struct and any other F derivatives (like your fmap implementation) to include a value_type alias:

template <typename X> struct F_id : F<F_id<X>> {
    using value_type = X; // Add this line
    const X x;
    F_id(const X& x) : x(x) {}
    template <typename Callback>
    void operator()(Callback&& callback) const {
        callback(x);
    }
};

// Example for your existing fmap implementation (adjust as needed)
template <typename Upstream, typename Func>
struct F_fmap : F<F_fmap<Upstream, Func>> {
    using value_type = std::invoke_result_t<Func, typename Upstream::value_type>;
    Upstream upstream;
    Func func;
    F_fmap(const Upstream& up, Func f) : upstream(up), func(std::move(f)) {}
    template <typename Callback>
    void operator()(Callback&& callback) const {
        upstream([=](auto val) {
            callback(func(val));
        });
    }
};

Step 2: Rewrite F_zip for Parallel Execution

We'll use std::optional to track whether each value is available, a mutex to protect shared state (critical for async execution), and std::shared_ptr to manage the lifecycle of our state (so we don't have dangling references if callbacks run asynchronously). Here's the updated F_zip struct:

template <typename Upstream, typename Other, typename Zip> struct F_zip : F<F_zip<Upstream, Other, Zip>> {
    using value_type = std::invoke_result_t<Zip, typename Upstream::value_type, typename Other::value_type>;
    Upstream upstream;
    Other other;
    Zip zip;
    F_zip(const Upstream& upstream, const Other& other, const Zip& zip) : upstream(upstream), other(other), zip(zip) {}

    template <typename Callback>
    void operator()(Callback&& callback) const {
        using A = typename Upstream::value_type;
        using B = typename Other::value_type;

        // Shared state to track values and synchronization
        struct State {
            std::optional<A> a;
            std::optional<B> b;
            std::mutex mtx;
            Callback final_callback;

            State(Callback cb) : final_callback(std::move(cb)) {}
        };
        auto state = std::make_shared<State>(std::forward<Callback>(callback));

        // Helper to check if both values are ready and trigger the callback
        auto try_invoke_callback = [state, this]() {
            std::lock_guard<std::mutex> lock(state->mtx);
            if (state->a.has_value() && state->b.has_value()) {
                state->final_callback(zip(*state->a, *state->b));
            }
        };

        // Launch upstream execution
        upstream([state, try_invoke_callback](auto x) {
            std::lock_guard<std::mutex> lock(state->mtx);
            state->a = std::move(x);
            try_invoke_callback();
        });

        // Launch other execution
        other([state, try_invoke_callback](auto y) {
            std::lock_guard<std::mutex> lock(state->mtx);
            state->b = std::move(y);
            try_invoke_callback();
        });
    }
};

Step 3: Fix the Test Function

Your test had a small typo (atuo → auto). Here's the corrected version:

void test_zip() {
    F_id(10)
    .zip(F_id(20), [](int x, int y) { return std::make_tuple(x, y); })
    ([](auto x) {
        auto [a, b] = x;
        printf("(%d, %d)\n", a, b);
    });
}

How This Works

  1. State Management: The State struct holds our optional values, a mutex for thread safety, and the user's final callback. We use std::shared_ptr to ensure this state exists until both callbacks complete.
  2. Parallel Execution: We start both upstream and other operations immediately (instead of nesting them like your original implementation).
  3. Callback Trigger: Each time a value arrives, we update the corresponding std::optional and check if both values are ready. Only when both are available do we invoke the user's callback with the zipped result.
  4. Thread Safety: The mutex ensures we don't have race conditions when updating the shared state, which is crucial if upstream or other execute asynchronously (e.g., on separate threads).

Notes

  • If your upstream and other are strictly synchronous (no async execution), you can omit the mutex and shared_ptr for simplicity—just use stack-allocated std::optional variables.
  • This pattern extends naturally to zipping more than two sources (you'd add more std::optional members and adjust the readiness check).

内容的提问来源于stack exchange,提问作者Artiano

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 08:06:04