如何修改C++中带Continuation的Zip实现以支持并行执行?
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
- State Management: The
Statestruct holds our optional values, a mutex for thread safety, and the user's final callback. We usestd::shared_ptrto ensure this state exists until both callbacks complete. - Parallel Execution: We start both
upstreamandotheroperations immediately (instead of nesting them like your original implementation). - Callback Trigger: Each time a value arrives, we update the corresponding
std::optionaland check if both values are ready. Only when both are available do we invoke the user's callback with the zipped result. - Thread Safety: The mutex ensures we don't have race conditions when updating the shared state, which is crucial if
upstreamorotherexecute asynchronously (e.g., on separate threads).
Notes
- If your
upstreamandotherare strictly synchronous (no async execution), you can omit the mutex andshared_ptrfor simplicity—just use stack-allocatedstd::optionalvariables. - This pattern extends naturally to zipping more than two sources (you'd add more
std::optionalmembers and adjust the readiness check).
内容的提问来源于stack exchange,提问作者Artiano

