D2300R11: `std::execution`

Lobsters Hottest 论文

摘要

本文提出了一个标准的C++框架,用于管理异步执行,引入了调度器、发送器和接收器,以增强C++编程中的异步性和并行性。

<p><a href="https://lobste.rs/s/l8evsw/d2300r11_std_execution">评论</a></p>
查看原文
查看缓存全文

缓存时间: 2026/09/21 20:29

# P2300R10:`std::execution` 来源:https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html ## 1. 引言 https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#intro 本文提出了一种用于在通用执行资源上管理异步执行的、自包含的标准C++框架设计方案。它基于《C++统一执行器提案》(https://wg21.link/p0443r14)及其配套论文中的思想。 ### 1.1. 动机 https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#motivation 如今,C++软件正日益异步化和并行化,这一趋势很可能将持续下去。异步性和并行性无处不在,从处理器硬件接口、网络、文件I/O、GUI到加速器。每个C++领域和每个平台都需要处理异步和并行,从科学计算、视频游戏、金融服务,到最小的移动设备、您的笔记本电脑,再到世界上最快超级计算机中的GPU。 尽管C++标准库提供了丰富的并发原语(`std::atomic`、`std::mutex`、`std::counting_semaphore`等)和更低层次的构建块(`std::thread`等),但我们仍然缺乏C++程序员迫切需要的、用于异步和并行的标准词汇表和框架。`std::async`/`std::future`/`std::promise`作为C++11为异步性设计的对外接口,效率低下、难以正确使用,并且严重缺乏泛型性,导致其在许多场景下无法使用。我们在C++17中将并行算法引入了C++标准库,虽然这是一个极好的开端,但它们本质上都是同步且不可组合的。 本文提出一个基于三个关键抽象——调度器、发送者和接收者——以及一系列可定制的异步算法的C++异步标准模型。 ### 1.2. 优先级 https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#priorities - **可组合与泛型化**:允许用户编写可用于多种不同执行资源的代码。 - **封装通用异步模式**:将通用异步模式封装在可定制和可复用的算法中,这样用户就不必自己重新发明轮子。 - **易于正确构建**:使代码易于正确构建。 - **支持执行资源和执行代理的多样性**:并非所有执行代理生而平等;有些能力较弱,但重要性不减。 - **允许执行资源定制一切**:包括切换到其他执行资源,但不要求执行资源定制所有内容。 - **关注所有合理的用例、领域和平台**。 - **错误必须传播**:错误处理不应成为负担。 - **支持取消**:取消不是一种错误。 - **清晰简洁的答案**:明确说明任务在哪里执行。 - **异步管理对象生命周期**:能够异步地管理和终止对象的生命周期。 ### 1.3. 示例:最终用户 https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#example-end-user 在本节中,我们将展示直接使用本文提出的发送者算法进行异步编程的终端用户体验。有关代码示例中使用的算法的简短说明,请参阅: - 第4.19节 面向用户的发送者工厂(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-factories) - 第4.20节 面向用户的发送者适配器(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptors) - 第4.21节 面向用户的发送者消费者(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-consumers) #### 1.3.1. Hello world https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#example-hello-world ```cpp using namespace std::execution; scheduler auto sch = thread_pool.scheduler(); // 1 sender auto begin = schedule(sch); // 2 sender auto hi = then(begin, []{ // 3 std::cout << "Hello world! Have an int."; // 3 return 13; // 3 }); // 3 sender auto add_42 = then(hi, [](int arg) { return arg + 42; }); // 4 auto [i] = this_thread::sync_wait(add_42).value(); // 5 ``` 此示例演示了调度器、发送者和接收者的基本概念: 1. 首先,我们需要从某个地方获取一个调度器,例如线程池。调度器是对执行资源的轻量级句柄。 2. 要在调度器上启动一连串工作,我们调用第4.19.1节 `execution::schedule`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-factory-schedule),它返回一个在指定调度器上完成的发送者。发送者描述异步工作,并在该工作完成时向某个接收者发送信号(值、错误或已停止)。 3. 我们使用发送者算法来生成发送者并组合异步工作。第4.20.2节 `execution::then`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptor-then)是一个发送者适配器,它接受一个输入发送者和一个 `std::invocable`(可调用对象),并在输入发送者发送的信号上调用该 `std::invocable`。`then` 返回的发送者会发送该调用的结果。在这个例子中,输入发送者来自 `schedule`,因此它是 `void` 类型,意味着它不会向我们发送值,所以我们的 `std::invocable` 不接受参数。但我们返回一个 `int`,它将被发送给下一个接收者。 4. 现在,我们使用第4.20.2节 `execution::then`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptor-then)向链中添加另一个操作。这一次,我们收到了一个值——前一步骤的 `int`。我们将其加上 `42`,然后返回结果。 5. 最后,我们准备提交整个异步管道并等待其完成。在此之前的一切都是完全异步的;工作甚至可能尚未开始。为确保工作已开始并阻塞直到其完成,我们使用第4.21.1节 `this_thread::sync_wait`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-consumer-sync_wait),它将返回一个包含最后一个发送者发送的值的 `std::optional`,或者如果最后一个发送者发送了已停止信号则返回一个空的 `std::optional`,或者如果最后一个发送者发送了错误则抛出异常。 #### 1.3.2. 异步包含扫描 https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#example-async-inclusive-scan ```cpp using namespace std::execution; sender auto async_inclusive_scan(scheduler auto sch, // 2 std::span<double> input, // 1 std::span<double> output, // 1 double init, // 1 std::size_t tile_count) // 3 { std::size_t const tile_size = (input.size() + tile_count - 1) / tile_count; std::vector<double> partials(tile_count + 1); // 4 partials[0] = init; // 4 return just(std::move(partials)) // 5 | continues_on(sch) // 6 | bulk(tile_count, // 6 [=](std::size_t i, std::vector<double>& partials) { // 7 auto start = i * tile_size; // 8 auto end = std::min(input.size(), (i + 1) * tile_size); // 8 partials[i + 1] = *std::inclusive_scan(begin(input) + start, // 9 begin(input) + end, // 9 begin(output) + start); // 9 }) // 10 | then( // 11 [](std::vector<double>&& partials) { std::inclusive_scan(begin(partials), end(partials), // 12 begin(partials)); // 12 return std::move(partials); // 13 }) | bulk(tile_count, // 14 [=](std::size_t i, std::vector<double>& partials) { // 14 auto start = i * tile_size; // 14 auto end = std::min(input.size(), (i + 1) * tile_size); // 14 std::for_each(begin(output) + start, begin(output) + end, // 14 [&](double& e) { e = partials[i] + e; } // 14 ); }) | then( // 15 [=](std::vector<double>&& partials) { // 15 return output; // 15 }); // 15 } ``` 此示例构建了一个异步的包含扫描计算: 1. 它扫描一个 `double` 序列(表示为 `std::span<double> input`),并将结果存储在另一个 `double` 序列中(表示为 `std::span<double> output`)。 2. 它接受一个调度器,指定应在哪个执行资源上启动扫描。 3. 它还接受一个 `tile_count` 参数,控制将生成的执行代理数量。 4. 首先,我们需要为算法分配所需的临时存储,这里我们使用 `std::vector<double> partials`。我们需要为创建的每个执行代理分配一个 `double` 的临时存储。 5. 接下来,我们使用第4.19.2节 `execution::just`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-factory-just)和第4.20.1节 `execution::continues_on`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptor-continues_on)创建初始发送者。这些发送者将发送我们已移入其中的临时存储。该发送者的完成调度器是 `sch`,这意味着链中的下一个任务将使用 `sch`。 6. 发送者和发送者适配器支持通过 `operator|` 进行组合,类似于C++ ranges。我们将使用 `operator|` 附加下一部分工作,该工作将使用第4.20.9节 `execution::bulk`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptor-bulk)生成 `tile_count` 个执行代理(详见第4.12节 大多数发送者适配器都是可管道化的(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-pipeable))。 7. 每个代理将调用一个 `std::invocable`,并向其传递两个参数。第一个参数是代理在第4.20.9节 `execution::bulk`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptor-bulk)操作中的索引(`i`),在此例中是一个在 `[0, tile_count)` 范围内的唯一整数。第二个参数是输入发送者发送的内容——临时存储。 8. 首先,我们根据代理索引计算该代理负责的输入和输出元素范围的起始和结束位置。 9. 然后,我们对我们的元素执行顺序 `std::inclusive_scan`。我们将最后一个元素的扫描结果(即我们所有元素的和)存储在我们的临时存储 `partials` 中。 10. 在该初始第4.20.9节 `execution::bulk`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptor-bulk)阶段的所有计算完成后,每个生成的执行代理都已将其元素的和写入其在 `partials` 中的槽位。 11. 现在我们需要扫描 `partials` 中的所有值。我们将在第4.20.9节 `execution::bulk`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptor-bulk)完成后,使用一个单独的执行代理来完成此操作。我们使用第4.20.2节 `execution::then`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptor-then)创建该执行代理。 12. 第4.20.2节 `execution::then`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptor-then)接受一个输入发送者和一个 `std::invocable`,并使用输入发送者发送的值调用该 `std::invocable`。在我们的 `std::invocable` 内部,我们对输入发送者将发送给我们的 `partials` 调用 `std::inclusive_scan`。 13. 然后我们返回 `partials`,下一阶段将需要它。 14. 最后,我们执行另一个与之前形状相同的第4.20.9节 `execution::bulk`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptor-bulk)。在这个 `bulk` 中,我们将使用 `partials` 中扫描后的值,将其他分片的和整合到我们的元素中,从而完成包含扫描。 15. `async_inclusive_scan` 返回一个发送输出 `std::span<double>` 的发送者。该算法的消费者可以链式附加使用扫描结果的更多工作。在 `async_inclusive_scan` 返回时,计算可能尚未完成。实际上,它甚至可能尚未开始。 #### 1.3.3. 异步动态大小读取 https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#example-async-dynamically-sized-read ```cpp using namespace std::execution; sender_of<void(std::span<std::byte>)> auto async_read( // 1 sender_of<void(std::span<std::byte>)> auto buffer, // 1 auto handle); // 1 struct dynamic_buffer { // 3 std::unique_ptr<std::byte[]> data; // 3 std::size_t size; // 3 }; // 3 sender_of<void(dynamic_buffer)> auto async_read_array(auto handle) { // 2 return just(dynamic_buffer{}) // 4 | let_value([handle] (dynamic_buffer& buf) { // 5 return just(std::as_writable_bytes(std::span(&buf.size, 1))) // 6 | async_read(handle) // 7 | then( // 8 [&buf] (std::size_t bytes_read) { // 9 assert(bytes_read == sizeof(buf.size)); // 10 buf.data = std::make_unique<std::byte[]>(buf.size); // 11 return std::span(buf.data.get(), buf.size); // 12 }) | async_read(handle) // 13 | then( [&buf] (std::size_t bytes_read) { assert(bytes_read == buf.size); // 14 return std::move(buf); // 15 }); }); } ``` 此示例演示了一种常见的异步I/O模式——通过先读取大小、然后读取由该大小指定的字节数来读取动态大小的有效负载: 1. `async_read` 是一个可管道化的发送者适配器。它是一个定制点对象,但其调用签名如下所示。它接受一个必须以 `std::span<std::byte>` 形式发送输入缓冲区的发送者参数,以及一个指向I/O上下文的句柄。它将异步读取数据到输入缓冲区,最大读取量为 `std::span` 的大小。它返回一个发送者,该发送者将在读取完成后发送读取的字节数。 2. `async_read_array` 接受一个I/O句柄,从中读取一个大小,然后读取该大小字节的缓冲区。它返回一个发送者,该发送者发送一个拥有已发送数据的 `dynamic_buffer` 对象。 3. `dynamic_buffer` 是一个聚合结构体,包含一个 `std::unique_ptr` 和一个大小。 4. 在 `async_read_array` 内部,我们首先使用第4.19.2节 `execution::just`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-factory-just)创建一个将发送新的空 `dynamic_buffer` 对象的发送者。我们可以使用 `operator|` 组合向管道附加更多工作(详见第4.12节 大多数发送者适配器都是可管道化的(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-pipeable))。 5. 我们需要这个 `dynamic_buffer` 对象的生命周期持续整个管道。因此,我们使用 `let_value`,它接受一个输入发送者和一个必须返回发送者本身的 `std::invocable`(详见第4.20.4节 `execution::let_*`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptor-let))。`let_value` 将输入发送者发送的值传递给 `std::invocable`。关键的是,发送对象的生命周期将持续到 `std::invocable` 返回的发送者完成为止。 6. 在 `let_value` 的 `std::invocable` 内部,我们有其余的逻辑。首先,我们想启动一个读取缓冲区大小的 `async_read`。为此,我们需要发送一个指向 `buf.size` 的 `std::span`。我们可以使用第4.19.2节 `execution::just`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-factory-just)来实现。 7. 我们使用 `operator|` 将 `async_read` 链接到 `just` 发送者上。 8. 接下来,我们使用第4.20.2节 `execution::then`(https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2024/p2300r10.html#design-sender-adaptor-then)管道一个在 `async_read` 完成后调用的 `std::invocable`。 9. 该 `std::invocable` 接收读取的字节数。 10. 我们需要检查读取的字节数是否符合预期。 11. 现在我们已经读取了数据的大小,可以为其分配存储空间。 12. 我们从 `std::invocable` 返回一个指向数据存储的 `std::span`。这将被发送给管道中的下一个接收者。 13. 该接收者将读取实际数据。

相似文章

std::call_once 与 std::async

The Old New Thing (Raymond Chen)

本文比较了 C++ 中用于延迟执行的 std::call_once 和 std::async,讨论了它们在实现、性能和异常行为方面的差异。

为3DS构建AsyncIO执行器

Lobsters Hottest

本文介绍了在Nintendo 3DS上进行异步编程的必要性,因为其采用协作式多任务处理,并开始解释如何为其构建一个asyncio执行器,重点讨论了Rust中的任务、未来、唤醒器和执行器这些概念。

Async/Await 的设计空间探索

Lobsters Hottest

本文介绍了编程语言中直线异步性的设计空间探索,审视了不同语言中async/await实现的差异及其语义后果。

面向分布式并行AI程序验证的定向神经符号随机执行

arXiv cs.AI

本文提出了DNSSE,一种混合框架,结合了基于LLM的调度预测、符号约束求解和覆盖引导的随机变异,用于验证分布式并行AI程序。它在现实基准上检测到的并发缺陷数量是基线的2.9倍,并将分支覆盖率从68.6%提高到91.6%。