进阶 Steve Klabnik / Carol Nichols 著,KaiserY 中译 2026-09-13 14:34:37 · 0 阅读

第17章 异步编程基础:Async、Await、Future 与 Stream

很多我们要求计算机执行的操作都需要一段时间才能完成。如果在等待这些长时间运行的过程结束时,我们还能做点别的事情,就再好不过了。现代计算机提供了两种同时处理多个操作的技术:并行和并发。然而,我们的程序逻辑通常是以近乎线性的方式编写的。我们希望能够描述程序应执行的操作,以及函数在哪些点可以暂停并让程序的其他部分转而运行,而不需要预先精确指定每段代码究竟该以什么顺序、用什么方式运行。异步编程 正是一种抽象,它让我们能够用“可能暂停的点”和“最终结果”来表达代码,而协调执行的细节则由它来替我们处理。

本章以前一章通过线程实现并行和并发为基础,引入另一种编写代码的方式:Rust 的 futures、streams,以及 asyncawait 语法。它们让我们能够表达“操作可以是异步的”这一点;此外,还有由第三方 crate 提供的异步运行时,用来管理和协调这些异步操作的执行。

来看一个例子。假设你正在导出一个家庭聚会的视频,这个操作可能要花上几分钟甚至几小时。视频导出会尽可能多地占用 CPU 和 GPU 资源。如果你只有一个 CPU 核,而操作系统又不会在导出完成前暂停这个任务,也就是说它会以同步方式执行导出,那么在这个任务运行期间你就什么别的都做不了。这会是相当糟糕的体验。幸运的是,你的操作系统可以,而且确实会,足够频繁地打断导出任务,让你同时完成其他工作。

再假设你正在下载别人分享给你的视频。这个操作也可能耗时很长,但不会占用太多 CPU 时间。此时 CPU 主要是在等待网络数据到达。虽然数据一开始到达时你就可以开始读取,但全部数据到齐仍然可能需要一段时间。即便数据已经全部到达,如果视频文件很大,完整载入它也至少可能需要一两秒。听起来不算长,但对于每秒能执行数十亿次操作的现代处理器来说,这已经是很长的时间了。和前面的例子一样,操作系统也会在等待网络调用完成时悄悄打断程序,让 CPU 去处理其他工作。

视频导出属于 CPU 密集型CPU-bound)或 计算密集型compute-bound)操作:它受限于 CPU 或 GPU 处理数据的速度,以及操作能分配到多少计算能力。视频下载则属于 I/O-bound 操作,因为它受限于计算机的 输入输出input and output)速度;它的速度最多只能和网络传输数据的速度一样快。

在上述两个例子中,操作系统的隐式中断提供了一种形式的并发。不过这种并发仅限于整个程序的级别:操作系统中断一个程序并让其它程序得以执行。在很多场景中,由于我们能比操作系统在更细粒度上理解我们的程序,因此我们可以观察到很多操作系统无法察觉的并发机会。

例如,如果我们正在构建一个管理文件下载的工具,程序就应该能做到:启动一个下载任务不会让 UI 卡死,而且用户还能够同时启动多个下载任务。不过,许多操作系统中与网络交互的 API 都是 blocking 的;也就是说,在它们所处理的数据完全就绪之前,会阻止程序继续向前执行。

注意:如果仔细想想,这其实也是大多数函数调用的工作方式。不过,blocking 这个术语通常保留给那些与文件、网络或计算机上其他资源交互的函数调用,因为正是在这些场景中,单个程序才会从 non-blocking 操作中明显受益。

我们可以为每个文件单独创建一个线程来避免阻塞主线程。然而,这些线程所消耗的系统资源最终会成为问题。更理想的情况是:这些调用从一开始就不是阻塞的,并且我们只需定义程序想完成的一组任务,然后让运行时自行选择最佳的执行顺序和方式。

这正是 Rust 的 asyncasynchronous 的缩写)抽象所提供的能力。本章会介绍以下内容:

  • 如何使用 Rust 的 asyncawait 语法,并借助运行时执行异步函数
  • 如何用异步模型解决一些我们在第十六章中已经遇到过的挑战
  • 多线程和异步如何提供互补的解决方案,以及它们在许多场景下如何组合使用

不过在看到 async 的实际工作方式之前,我们需要先稍微绕个远路,讨论一下并行(parallelism)和并发(concurrency)的区别。

并行与并发

在上一章中,我们大致将并行和并发视为可以互换的概念。但现在我们需要更加精确地区分它们,因为它们的区别将在实际工作中显现出来。

思考一下不同的团队分割方法来开发一个软件项目。我们可以分配给一个个人多个任务,也可以每个团队成员各自负责一个任务,或者可以采用这两种方法的组合。

当一个人在任何一个任务都还没完成之前,就在多个不同任务之间切换工作,这就是 并发。一种实现并发的方式,很像你在电脑上同时 checkout 了两个不同项目;当你对其中一个项目感到厌倦,或是在上面卡住时,就切到另一个。你只有一个人,所以不可能在完全相同的时刻同时推进两个任务,但你可以通过在它们之间切换来多任务处理,一次推进一个任务。

并发工作流
图 17-1:一个并发工作流,在任务 A 和任务 B 之间切换

当团队把一组任务拆开,让每个成员各自负责一个任务并单独推进时,这就是 并行。团队中的每个人都可以在完全相同的时间取得进展。

并行工作流
图 17-2:一个并行流,其中任务 A 和任务 B 的工作同时独立进行

在这两种工作流中,你都可能需要在不同任务之间做协调。也许你原以为分配给某个人的任务和其他人的工作完全独立,但实际上它必须等另一个人先完成自己的任务。有些工作可以并行完成,但其中一些实际上是 串行 的:它们只能按顺序发生,一个接一个,如图 17-3 所示。

部分并行工作流
图 17-3:一个部分并行的工作流,其中任务 A 和任务 B 的工作相互独立,直到任务 A3 阻塞在等待任务 B3 的结果

同样,你也可能意识到自己的一个任务依赖于另一个任务。那么你原本的并发工作也变成了串行。

并行和并发之间也可能彼此交叉。如果你得知某位同事正在卡着,必须等你先完成某个任务,那你很可能会把全部精力都集中到那个任务上,好“解除阻塞”。这时你和同事就不能再并行工作了,而你自己也不能再并发地推进其他任务。

同样的基础动态也作用于软件与硬件。在一个单核的机器上,CPU 一次只能执行一个操作,不过它仍然可以并发工作。借助像线程、进程和异步(async)等工具,计算机可以暂停一个活动,并在最终切换回第一个活动之前切换到其它活动。在一个有多个 CPU 核心的机器上,它也可以并行工作。一个核心可以做一件工作的同时另一个核心可以做一些完全不相关的工作,而且这些工作实际上是同时发生的。

在 Rust 中运行 async 代码时,通常是在并发地执行。至于这种并发在底层是否也会利用并行,则取决于硬件、操作系统,以及所使用的异步运行时(稍后我们会进一步介绍异步运行时)。

现在,让我们深入看看 Rust 中的异步编程究竟是如何工作的。

Futures 和 async 语法

Future 与 async 语法

Rust 异步编程的关键元素是 futures 和 Rust 的 asyncawait 关键字。

future 是一个现在也许还没准备好,但会在将来某个时刻准备好的值。(这个概念在很多语言里都存在,只是有时会用 taskpromise 之类的名字。)Rust 提供了 Future trait 作为基础构件,让不同的异步操作可以用不同的数据结构来实现,同时又拥有统一的接口。在 Rust 中,future 就是那些实现了 Future trait 的类型。每个 future 都保存了自身的进度信息,以及“就绪”到底意味着什么。

async 关键字可以用于代码块和函数,表示它们可以被中断和恢复。在 async 块或 async 函数中,你可以使用 await 关键字来 await 一个 future,也就是等待它变为就绪。在 async 块或函数里,每个等待 future 的位置,都是这个块或函数可能暂停并随后恢复的点。检查 future、看看它的值是否已经可用,这个过程称为 polling(轮询)。

其他一些语言,例如 C# 和 JavaScript,也用 asyncawait 关键字进行异步编程。如果你熟悉这些语言,可能会注意到 Rust 在语法处理上存在一些明显差异。我们会看到,这样设计是有充分理由的。

编写异步 Rust 时,大多数时候我们直接使用 asyncawait 关键字。Rust 会把它们编译成等价的、基于 Future trait 的代码,就像它把 for 循环编译成基于 Iterator trait 的等价代码一样。不过,既然 Rust 提供了 Future trait,你在需要时也可以为自己的数据类型实现它。本章中我们会见到很多函数,它们都返回拥有各自 Future 实现的类型。我们会在本章结尾回到这个 trait 的定义,进一步深入理解它的工作原理;不过眼下这些细节已经足够让我们继续前进。

这些内容可能仍然有些抽象,所以我们来写第一个异步程序:一个小型网页抓取器。我们会从命令行传入两个 URL,并发地抓取它们,然后返回那个最先完成的结果。这个例子会带来不少新语法,不过不用担心,我们会一路把需要知道的内容都解释清楚。

第一个异步程序

为了让本章专注于学习 async,而不是在生态系统的各种组件之间来回切换,我们准备了一个 trpl crate(trpl 是 “The Rust Programming Language” 的缩写)。它重新导出了本章需要的所有类型、trait 和函数,主要来自 futurestokio crate。futures crate 是 Rust 异步代码实验的官方阵地,Future trait 最初就是在那里设计出来的。Tokio 则是目前 Rust 中使用最广泛的异步运行时(async runtime),尤其常见于 Web 应用。生态中也还有其他很优秀的运行时,而且它们可能更适合你的实际用途。我们在 trpl 的底层使用 tokio,是因为它经过了充分测试,也足够常用。

在某些场景下,trpl 还会对原始 API 进行重命名或包装,好让你把注意力集中在本章相关的细节上。如果你想了解这个 crate 实际做了什么,我们建议你看看它的源码。你可以从中看到每个重导出项究竟来自哪个 crate,我们也留下了很多注释来解释这个 crate 的行为。

创建一个名为 hello-async 的二进制项目并将 trpl crate 作为一个依赖添加:

$ cargo new hello-async
$ cd hello-async
$ cargo add trpl

现在我们可以利用 trpl 提供的各种组件来编写第一个异步程序。我们要构建一个小型命令行工具:抓取两个网页,从各自页面中提取 <title> 元素,然后打印出那个最先完成整套流程的页面标题。

定义 page_title 函数

让我们开始编写一个函数,它获取一个网页 URL 作为参数,请求该 URL 并返回标题元素的文本(见示例 17-1)。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

fn main() {
    // TODO: we'll add this next!
}

use trpl::Html;

async fn page_title(url: &str) -> Option<String> {
    let response = trpl::get(url).await;
    let response_text = response.text().await;
    Html::parse(&response_text)
        .select_first("title")
        .map(|title| title.inner_html())
}
示例 17-1:定义一个 async 函数来获取一个 HTML 页面的标题元素

首先,我们定义了一个名为 page_title 的函数,并用 async 关键字标记它。然后使用 trpl::get 函数抓取传入的 URL,再用 await 关键字等待响应。为了得到 response 的文本,我们调用它的 text 方法,并再次使用 await 进行等待。这两个步骤都是异步的。对于 get 函数来说,我们必须等待服务器先把响应的第一部分发回来,其中包括 HTTP headers、cookies 等,这些内容可以和响应体分开发送。尤其当响应体很大时,全部数据到达可能要花上一些时间。由于我们必须等待响应完整到达,text 方法自然也是 async 的。

我们必须显式地等待这两个 future,因为 Rust 中的 future 是 lazy 的:在你用 await 请求它之前,它什么都不会做。(实际上,如果你创建了 future 却不使用它,Rust 还会给出编译器警告。)这大概会让你想起第十三章“使用迭代器处理元素序列”中的讨论。迭代器只有在你调用 next 方法时才会工作,无论是直接调用,还是通过 for 循环,或者借助像 map 这样底层会调用 next 的方法。future 也是一样,只有你显式要求它运行时,它才会开始工作。这种惰性让 Rust 能够避免在真正需要之前就运行异步代码。

注意:这和我们在第十六章“使用 spawn 创建新线程”里看到的 thread::spawn 的行为不同,在那里我们传给新线程的闭包会立刻开始执行。它也和许多其他语言处理 async 的方式不同。但这对于 Rust 提供它一贯的性能保证很重要,正如迭代器也是如此。

有了 response_text 之后,我们就可以用 Html::parse 把它解析成 Html 类型的实例。这样一来,我们得到的就不再是原始字符串,而是一个可以把 HTML 当作更丰富数据结构来操作的类型。特别是,我们可以用 select_first 方法找到给定 CSS selector 的第一个匹配项。传入字符串 "title" 后,我们就能拿到文档中的第一个 <title> 元素,如果它存在的话。因为也可能根本没有匹配项,所以 select_first 返回的是 Option<ElementRef>。最后,我们使用 Option::map 方法:如果 Option 中有值,它就会对其中的值进行处理;如果没有,就什么都不做。(这里当然也可以使用 match 表达式,不过 map 更符合惯用写法。)在我们传给 map 的闭包里,会对 title 调用 inner_html 来获取其中的内容,它是一个 String。到这里,我们最终得到的就是一个 Option<String>

注意,Rust 的 await 关键字放在要等待的表达式后面,而不是前面。也就是说,它是一个 postfix keyword(后缀关键字)。如果你在其他语言里用过 async,这一点可能和你的习惯不同;但在 Rust 中,这种设计会让链式方法调用更易读。因此,我们可以把 page_title 的函数体改写成在 trpl::gettext 调用之间插入 await 的链式写法,如示例 17-2 所示:

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use trpl::Html;

fn main() {
    // TODO: we'll add this next!
}

async fn page_title(url: &str) -> Option<String> {
    let response_text = trpl::get(url).await.text().await;
    Html::parse(&response_text)
        .select_first("title")
        .map(|title| title.inner_html())
}
示例 17-2:使用 `await` 关键字的链式调用

这样我们就成功编写了第一个异步函数!在我们向 main 加入一些代码调用它之前,让我们再多了解下我们写了什么以及它的意义。

当 Rust 遇到一个 async 关键字标记的代码块时,会将其编译为一个实现了 Future trait 的唯一的、匿名的数据类型。当 Rust 遇到一个被标记为 async 的函数时,会将其编译成一个函数体是异步代码块的非异步函数。异步函数的返回值类型是编译器为异步代码块所创建的匿名数据类型。

因此,编写 async fn 就等同于编写一个返回类型为 future 的函数。当编译器遇到类似示例 17-1 中 async fn page_title 的函数定义时,它等价于以下定义的非异步函数:

#![allow(unused)]
fn main() {
extern crate trpl; // required for mdbook test
use std::future::Future;
use trpl::Html;

fn page_title(url: &str) -> impl Future<Output = Option<String>> {
    async move {
        let text = trpl::get(url).await.text().await;
        Html::parse(&text)
            .select_first("title")
            .map(|title| title.inner_html())
    }
}
}

让我们挨个看一下转换后版本的每一个部分:

  • 它使用了之前第十章 “trait 作为参数” 部分讨论过的 impl Trait 语法。
  • 它返回的值实现了 Future trait,并且这个 trait 有一个关联类型 Output。注意 Output 的类型是 Option<String>,这和 async fn 版本的 page_title 的原始返回类型一致。
  • 原始函数体中的所有代码都被包进了一个 async move 块。回忆一下,代码块本身就是表达式。整个块就是函数返回的那个表达式。
  • 如上所述,这个异步代码块产生一个 Option<String> 类型的值。这个值与返回类型中的 Output 类型一致。这正类似于你已经见过的其它代码块。
  • 这个新函数体之所以是 async move 块,是由它使用 url 参数的方式决定的。(本章后面会更详细地讨论 asyncasync move 的区别。)

现在我们可以在 main 中调用 page_title

使用运行时执行异步函数

首先,我们只获取单个页面的标题,如示例 17-3 所示。不幸的是,这段代码还不能编译。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use trpl::Html;

async fn main() {
    let args: Vec<String> = std::env::args().collect();
    let url = &args[1];
    match page_title(url).await {
        Some(title) => println!("The title for {url} was {title}"),
        None => println!("{url} had no title"),
    }
}

async fn page_title(url: &str) -> Option<String> {
    let response_text = trpl::get(url).await.text().await;
    Html::parse(&response_text)
        .select_first("title")
        .map(|title| title.inner_html())
}
示例 17-3:在 `main` 中通过一个用户提供的参数调用 `page_title` 函数

我们沿用了第十二章“接受命令行参数”一节中获取命令行参数的模式。然后把 URL 参数传给 page_title,再等待它的结果。由于 future 产出的值是 Option<String>,我们使用 match 表达式来根据页面是否含有 <title> 打印不同的信息。

唯一能使用 await 关键字的地方,是 async 函数或 async 代码块中,而 Rust 又不允许我们把特殊的 main 函数标记为 async

error[E0752]: `main` function is not allowed to be `async`
 --> src/main.rs:6:1
  |
6 | async fn main() {
  | ^^^^^^^^^^^^^^^ `main` function is not allowed to be `async`

main 不能标记为 async 的原因是异步代码需要一个 运行时:即一个管理执行异步代码细节的 Rust crate。一个程序的 main 函数可以 初始化 一个运行时,但是其 自身 并不是一个运行时。(稍后我们会进一步解释原因。)每一个执行异步代码的 Rust 程序必须至少有一个设置运行时并执行 futures 的地方。

大多数支持 async 的语言都会自带运行时,但 Rust 不会。相反,Rust 有很多不同的异步运行时可供选择,每一种都针对自己的目标用例做了不同权衡。比如,一个拥有许多 CPU 核心和大量 RAM 的高吞吐 Web 服务器,和一个单核、RAM 很小、甚至不能进行堆分配的微控制器,需求就截然不同。提供这些运行时的 crate 往往也会一并提供文件或网络 I/O 等常见功能的异步版本。

在这里,以及本章余下的部分,我们会使用 trpl crate 提供的 block_on 函数。它接受一个 future 作为参数,并阻塞当前线程,直到这个 future 运行完成为止。在内部,调用 block_on 会借助 tokio crate 设置一个运行时,用来执行传入的 future(trplblock_on 和其他运行时 crate 提供的同名函数行为类似)。一旦 future 完成,block_on 就会返回 future 产生的值。

我们当然可以把 page_title 返回的 future 直接传给 block_on,并在它完成后对得到的 Option<String> 进行匹配,就像我们在示例 17-3 中本来打算做的那样。不过,本章的大部分例子里(以及现实中的大多数 async 代码里),我们都不止会进行一次异步函数调用,因此我们改为传入一个 async 块,并在其中显式等待 page_title 的结果,如示例 17-4 所示。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use trpl::Html;

fn main() {
    let args: Vec<String> = std::env::args().collect();

    trpl::block_on(async {
        let url = &args[1];
        match page_title(url).await {
            Some(title) => println!("The title for {url} was {title}"),
            None => println!("{url} had no title"),
        }
    })
}

async fn page_title(url: &str) -> Option<String> {
    let response_text = trpl::get(url).await.text().await;
    Html::parse(&response_text)
        .select_first("title")
        .map(|title| title.inner_html())
}
示例 17-4:使用 `trpl::block_on` 等待一个 async 代码块

当我们运行这段代码时,就会得到一开始期待的行为:

$ cargo run -- https://www.rust-lang.org
    Finished `dev` profile [unoptimized + debuginfo] target(s) in 0.05s
     Running `target/debug/async_await 'https://www.rust-lang.org'`
The title for https://www.rust-lang.org was
            Rust Programming Language

我们终于有了一些可以正常工作的异步代码!不过在我们添加代码让两个网址进行竞争之前,让我们简要地回顾一下 future 是如何工作的。

每一个 await point,也就是代码使用 await 关键字的地方,代表将控制权交还给运行时的地方。为此 Rust 需要记录异步代码块中涉及的状态,这样运行时可以去执行其他工作,并在准备好时回来继续推进当前的任务。这就像你通过编写一个枚举来保存每一个 await point 的状态一样:

#![allow(unused)]
fn main() {
extern crate trpl; // required for mdbook test

enum PageTitleFuture<'a> {
    Initial { url: &'a str },
    GetAwaitPoint { url: &'a str },
    TextAwaitPoint { response: trpl::Response },
}
}

编写代码来手动控制不同状态之间的转换是非常乏味且容易出错的,特别是之后增加了更多功能和状态的时候。相反,Rust 编译器自动创建并管理异步代码的状态机数据结构。如果你感兴趣的话:是的,正常的借用和所有权也全部适用于这些数据结构。幸运的是,编译器也会为我们处理这些检查,并提供友好的错误信息。本章稍后会讲解一些相关内容!

最终,总得有某个组件来执行这个状态机,而那个组件就是运行时。(这也是为什么在了解运行时时,你可能会看到 executor 这个词:executor 是运行时中负责执行异步代码的那一部分。)

现在你就能理解,为什么编译器会在示例 17-3 中阻止我们把 main 本身写成异步函数了。如果 main 是 async 函数,那么就必须有别的东西来管理 main 返回的 future 对应的状态机;可 main 本身就是程序的入口点!因此,我们改为在 main 中调用 trpl::block_on,让它设置好运行时,并运行 async 块返回的 future,直到执行完成。

注意:有些运行时会提供宏,因此你确实可以写异步版的 main 函数。这些宏会把 async fn main() { ... } 重写成普通的 fn main,其逻辑和我们在示例 17-4 中手动做的事情一样:调用一个像 trpl::block_on 这样的函数,把 future 跑到完成为止。

现在让我们把这些部分组合起来,看看如何编写并发代码。

让两个 URL 并发竞争

在示例 17-5 中,我们会对从命令行传入的两个不同 URL 分别调用 page_title,并选出最先完成的那个 future。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use trpl::{Either, Html};

fn main() {
    let args: Vec<String> = std::env::args().collect();

    trpl::block_on(async {
        let title_fut_1 = page_title(&args[1]);
        let title_fut_2 = page_title(&args[2]);

        let (url, maybe_title) =
            match trpl::select(title_fut_1, title_fut_2).await {
                Either::Left(left) => left,
                Either::Right(right) => right,
            };

        println!("{url} returned first");
        match maybe_title {
            Some(title) => println!("Its page title was: '{title}'"),
            None => println!("It had no title."),
        }
    })
}

async fn page_title(url: &str) -> (&str, Option<String>) {
    let response_text = trpl::get(url).await.text().await;
    let title = Html::parse(&response_text)
        .select_first("title")
        .map(|title| title.inner_html());
    (url, title)
}
示例 17-5:对两个 URL 调用 `page_title`,看谁先返回

我们首先分别对用户提供的两个 URL 调用 page_title。随后把得到的 future 保存到 title_fut_1title_fut_2 中。记住,它们此时还什么都没做,因为 future 是惰性的,而我们也还没有等待它们。接着我们把这些 future 传给 trpl::select,它会返回一个值,用来表明传入的 future 中哪一个最先完成。

注意:在底层,trpl::select 建立在 futures crate 中更通用的 select 函数之上。futures crate 的 select 函数能做很多 trpl::select 做不到的事,不过它也带来了一些额外复杂性,所以我们暂时先跳过。

任意一个 future 都有可能“获胜”,因此这里返回 Result 并不合理。相反,trpl::select 返回的是一个我们之前还没见过的类型:trpl::EitherEither 在某种程度上有点像 Result,也有两个分支;但不同的是,它并没有内建“成功”或“失败”的语义,而是用 LeftRight 来表示“这个或那个”。

#![allow(unused)]
fn main() {
enum Either<A, B> {
    Left(A),
    Right(B),
}
}

如果第一个参数先完成,select 就返回 Left,其中包含该 future 的输出;如果第二个 future 先完成,则返回 Right,其中包含第二个 future 的输出。这正好对应函数调用时参数的顺序:第一个参数位于第二个参数的左边。

我们还更新了 page_title,让它把传入的 URL 一并返回。这样一来,即使最先返回的页面无法解析出 <title>,我们仍然可以打印出一条有意义的信息。有了这些数据之后,我们最后再调整 println! 的输出,让它既能显示哪个 URL 最先完成,也能在页面存在 <title> 时打印出标题内容。

至此,你已经构建出了一个可以工作的迷你网页抓取器!随便选两个 URL 运行一下这个命令行工具吧。你会发现有些站点总是比另一些更快,而另一些情况下则每次运行谁快谁慢都不一定。更重要的是,你已经掌握了使用 future 的基础知识,所以现在我们可以继续深入,看看 async 还能做些什么。

并发与 async

使用 async 实现并发

在这一部分,我们将使用异步来应对一些与第十六章中通过线程解决的相同的并发问题。因为之前我们已经讨论了很多关键理念了,这一部分我们会专注于线程与 future 的区别。

在很多情况下,使用异步处理并发的 API 与使用线程的非常相似。在其它的一些情况,它们则非常不同。即便线程与异步的 API 看起来 很类似,通常它们有着不同的行为,同时它们几乎总是有着不同的性能特点。

使用 spawn_task 创建新任务

第十六章中我们应付的第一个任务是在两个不同的线程中计数。让我们用异步来完成相同的任务。trpl crate 提供了一个 spawn_task 函数,它看起来非常像 thread::spawn API,和一个 sleep 函数,这是 thread::sleep API 的异步版本。我们可以将它们结合使用,实现与线程示例相同的计数功能,如示例 17-6 所示。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        trpl::spawn_task(async {
            for i in 1..10 {
                println!("hi number {i} from the first task!");
                trpl::sleep(Duration::from_millis(500)).await;
            }
        });

        for i in 1..5 {
            println!("hi number {i} from the second task!");
            trpl::sleep(Duration::from_millis(500)).await;
        }
    });
}
示例 17-6:创建一个新任务,在主任务打印内容的同时打印另一组内容

作为起点,我们在 main 函数中使用 trpl::block_on,这样顶层函数就可以写成 async 风格。

注意:从这里开始,本章中的每个示例在 main 中都会包含这段几乎完全一样的 trpl::block_on 包装代码,所以之后我们通常会像省略 main 一样把它省掉。记得在你自己的代码里补上它!

然后我们在这个代码块里写了两个循环,每个循环中都调用了 trpl::sleep,在输出下一条消息之前等待半秒(500 毫秒)。其中一个循环放在 trpl::spawn_task 的函数体里,另一个则放在顶层的 for 循环中。我们还在 sleep 调用后加上了 await

这段代码的行为和线程版实现很像,包括当你亲自运行时,终端中的消息顺序可能和这里不完全一样:

hi number 1 from the second task!
hi number 1 from the first task!
hi number 2 from the first task!
hi number 2 from the second task!
hi number 3 from the first task!
hi number 3 from the second task!
hi number 4 from the first task!
hi number 4 from the second task!
hi number 5 from the first task!

这个版本会在主 async 块中的 for 循环一结束就停止,因为当 main 函数结束时,由 spawn_task 生成的任务也会被关闭。如果你想让它一直运行到任务自身完成,就需要使用 join handle 来等待第一个任务结束。在线程的版本中,我们使用 join 方法“阻塞”等待线程运行结束。在示例 17-7 中,我们可以使用 await 做同样的事,因为任务句柄本身就是一个 future。它的 Output 类型是 Result,所以在等待之后还要再 unwrap 一次。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let handle = trpl::spawn_task(async {
            for i in 1..10 {
                println!("hi number {i} from the first task!");
                trpl::sleep(Duration::from_millis(500)).await;
            }
        });

        for i in 1..5 {
            println!("hi number {i} from the second task!");
            trpl::sleep(Duration::from_millis(500)).await;
        }

        handle.await.unwrap();
    });
}
示例 17-7:在 join 句柄上使用 `await`,让任务运行到完成

更新后的版本会一直运行到两个循环都完成为止:

hi number 1 from the second task!
hi number 1 from the first task!
hi number 2 from the first task!
hi number 2 from the second task!
hi number 3 from the first task!
hi number 3 from the second task!
hi number 4 from the first task!
hi number 4 from the second task!
hi number 5 from the first task!
hi number 6 from the first task!
hi number 7 from the first task!
hi number 8 from the first task!
hi number 9 from the first task!

到目前为止,看起来 async 和线程只是用不同语法实现了相似效果:在 join handle 上使用 await,而不是调用 join;同时对 sleep 调用也使用 await

更大的不同在于,我们根本不需要再创建另一个操作系统线程来做这件事。实际上,这里甚至连任务都不一定要创建。因为 async 代码块会被编译成匿名 future,我们可以把每个循环都放进一个 async 代码块里,然后让运行时使用 trpl::join 让它们都执行到完成。

在第十六章“等待所有线程完成”一节中,我们展示了如何对 std::thread::spawn 返回的 JoinHandle 调用 join 方法。trpl::join 与之类似,不过它面向的是 future。当你把两个 future 传给它时,它会生成一个新的 future;等到两个传入的 future 都完成时,这个新 future 的输出就是一个包含它们各自输出值的元组。因此,在示例 17-8 中,我们用 trpl::join 来等待 fut1fut2 完成。我们不会分别等待 fut1fut2,而是等待 trpl::join 生成的那个新 future。这里我们忽略它的输出,因为那不过是一个包含两个 unit 值的元组。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let fut1 = async {
            for i in 1..10 {
                println!("hi number {i} from the first task!");
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        let fut2 = async {
            for i in 1..5 {
                println!("hi number {i} from the second task!");
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        trpl::join(fut1, fut2).await;
    });
}
示例 17-8:使用 `trpl::join` 等待两个匿名 future

运行后,我们会看到两个 future 都执行到了结束:

hi number 1 from the first task!
hi number 1 from the second task!
hi number 2 from the first task!
hi number 2 from the second task!
hi number 3 from the first task!
hi number 3 from the second task!
hi number 4 from the first task!
hi number 4 from the second task!
hi number 5 from the first task!
hi number 6 from the first task!
hi number 7 from the first task!
hi number 8 from the first task!
hi number 9 from the first task!

现在你会发现,每次运行时顺序都完全一样,这和线程版本以及示例 17-7 中使用 trpl::spawn_task 的情况非常不同。这是因为 trpl::joinfair 的,也就是它会以同样的频率检查每一个 future,在它们之间交替进行;只要另一个 future 已经就绪,它就不会让其中一个一路领先。在线程模型下,由操作系统决定先检查哪个线程、让它运行多久。对于 async Rust,则由运行时决定先检查哪个任务。(在实践中,细节会复杂得多,因为异步运行时可能会在底层借助操作系统线程来实现并发,因此要保证公平性,对运行时来说可能意味着更多工作,但这仍然是可能做到的。)运行时并不一定会为任何给定操作都保证公平性,而且它们通常会提供不同的 API,让你自行决定是否需要公平性。

尝试这些不同的 await future 的变体来观察它们的效果:

  • 去掉一个或者两个循环外的异步代码块。
  • 在定义两个异步代码块后立刻 await 它们。
  • 只将第一个循环封装进异步代码块,并在第二个循环体之后 await 作为结果的 future。

作为额外的挑战,看看你能否在运行代码 之前 想出每个情况下的输出!

通过消息传递在两个任务之间发送数据

在 future 之间共享数据的方式也会让你感到熟悉:我们再次使用消息传递,只不过这次使用的是异步版本的类型和函数。为了展示基于线程的并发和基于 future 的并发之间的一些关键差别,我们会和第十六章“通过消息传递在线程间传送数据”一节稍微走一条不一样的路线。在示例 17-9 中,我们先只使用一个 async 代码块,而像之前那样显式地创建一个独立任务。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let val = String::from("hi");
        tx.send(val).unwrap();

        let received = rx.recv().await.unwrap();
        println!("received '{received}'");
    });
}
示例 17-9:创建一个异步信道(async channel)并赋值其两端为 `tx` 和 `rx`

这里我们使用了 trpl::channel,一个第十六章用于线程的多生产者、单消费者信道 API 的异步版本。异步版本的 API 与基于线程的版本只有一点微小的区别:它使用一个可变的而不是不可变的 rx,并且它的 recv 方法产生一个需要 await 的 future 而不是直接返回值。现在我们可以发送端向接收端发送消息了。注意我们无需产生一个独立的线程或者任务;只需等待(await) rx.recv 调用。

std::mpsc::channel 中的同步 Receiver::recv 方法阻塞执行直到它接收一个消息。trpl::Receiver::recv 则不会阻塞,因为它是异步的。不同于阻塞,它将控制权交还给运行时,直到接收到一个消息或者信道的发送端关闭。相比之下,我们不用 await send,因为它不会阻塞。也无需阻塞,因为信道的发送端的数量是没有限制的。

注意:因为这些 async 代码都运行在传给 trpl::block_on 的 async 代码块里,所以块中的所有内容都可以避免阻塞。不过,块外部的代码则会阻塞,直到 block_on 返回为止。这正是 trpl::block_on 的意义所在:它让你可以选择在哪一处对一组 async 代码进行阻塞,从而也就决定了在什么地方切换同步和异步代码。

请注意这个示例中的两个地方:首先,消息立刻就会到达!其次,虽然我们使用了 future,但是这里还没有并发。示例中的所有事情都是顺序发生的,就像没涉及到 future 时一样。

让我们通过发送一系列消息并在之间休眠来解决第一个问题,如示例 17-10 所示:

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let vals = vec![
            String::from("hi"),
            String::from("from"),
            String::from("the"),
            String::from("future"),
        ];

        for val in vals {
            tx.send(val).unwrap();
            trpl::sleep(Duration::from_millis(500)).await;
        }

        while let Some(value) = rx.recv().await {
            println!("received '{value}'");
        }
    });
}
示例 17-10:通过异步信道发送和接收多个消息并在每个消息之间通过 `await` 休眠

除了发送消息之外,我们还需要接收它们。在这个例子中我们可以手动接收,就是调用四次 rx.recv().await,因为我们知道进来了多少条消息。然而,在现实世界中,我们通常会等待 未知 数量的消息。这时我们需要一直等待直到可以确认没有更多消息了为止。

在示例 16-10 中,我们使用 for 循环处理从同步信道接收到的所有条目。不过,Rust 目前还没有办法对异步产生的一系列条目使用 for 循环。因此,我们需要一种前面还没见过的循环:while let 条件循环。它正是我们在第六章“使用 if letlet...else 实现简洁控制流”中见过的 if let 结构的循环版本。只要它指定的模式还在持续匹配,循环就会继续执行。

rx.recv 调用产生一个 Future,我们会 await 它。运行时会暂停 Future 直到它就绪。一旦消息到达,future 会解析为 Some(message),每次消息到达时都会如此。当信道关闭时,不管是否有 任何 消息到达,future 都会解析为 None 来表明没有更多的值了,我们也就应该停止轮询,也就是停止等待。

while let 循环将上述逻辑整合在一起。如果 rx.recv().await 调用的结果是 Some(message),我们会得到消息并可以在循环体中使用它,就像使用 if let 一样。如果结果是 None,则循环停止。每次循环执行完毕,它会再次触发 await point,如此运行时会再次暂停直到另一条消息到达。

现在代码可以成功发送和接收所有的消息了。不幸的是,这里还有一些问题。首先,消息并不是按照半秒的间隔到达的。它们在程序启动后两秒(2000 毫秒)后立刻一起到达。其次,程序永远也不会退出!相反它会永远等待新消息。你会需要使用 ctrl-c 来关闭它。

一个 async 代码块中的代码会线性执行

先来看为什么这些消息会在完整延迟之后一起到达,而不是在每次延迟之后逐条到达。在一个给定的 async 代码块里,代码中 await 出现的顺序,也就是程序运行时它们执行的顺序。

示例 17-10 中只有一个 async 代码块,所以里面的一切都按线性顺序执行。这里依然没有并发。所有 tx.send 调用,连同 trpl::sleep 调用及其相应的 await 点,都会先全部依次发生。只有在那之后,while let 循环才有机会开始执行 recv 调用上的那些 await 点。

为了得到我们真正想要的行为,也就是在每条消息之间都出现休眠间隔,我们需要把 txrx 的操作分别放进各自的 async 代码块中,如示例 17-11 所示。这样运行时就可以像示例 17-8 那样,使用 trpl::join 分别执行它们。我们再次等待的是 trpl::join 调用的结果,而不是分别等待每个 future。要是依次等待它们,我们就又回到了顺序执行的流程,这正是我们想要的。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let tx_fut = async {
            let vals = vec![
                String::from("hi"),
                String::from("from"),
                String::from("the"),
                String::from("future"),
            ];

            for val in vals {
                tx.send(val).unwrap();
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        let rx_fut = async {
            while let Some(value) = rx.recv().await {
                println!("received '{value}'");
            }
        };

        trpl::join(tx_fut, rx_fut).await;
    });
}
示例 17-11:将 `send` 和 `recv` 分隔到其各自的 `async` 代码块中并 await 这些代码块的 future

使用示例 17-11 中更新后的代码后,消息就会以 500 毫秒的间隔输出,而不是在 2 秒之后一次性全部打印出来。

将所有权移入 async 代码块

但是程序仍然永远也不会退出,这是由于 while let 循环与 trpl::join 的交互方式所致:

  • trpl::join 返回的 future 只会完成一次,即传递的 两个 future 都完成的时候。
  • tx_fut future 会在发送完 vals 中最后一条消息后,再完成最后一次休眠之后结束。
  • rx_fut future 则要等到 while let 循环结束时才会结束。
  • 只有当等待 rx.recv 的结果变成 None 时,while let 循环才会结束。
  • 只有在信道另一端关闭后,等待 rx.recv 才会返回 None
  • 只有在我们调用 rx.close,或者发送端 tx 被 drop 时,信道才会关闭。
  • 我们根本没有调用 rx.close,而 tx 也要等到传给 trpl::block_on 的最外层 async 代码块结束后才会被 drop。
  • 但那个最外层 async 代码块又必须等 trpl::join 完成才能结束,于是我们就又回到了这个列表的起点。

目前,发送消息的那个 async 代码块只是借用tx,因为发送消息并不需要取得它的所有权。但如果我们能把 tx move 进那个 async 代码块里,那么一旦该代码块结束,tx 就会被 drop。在第十三章“捕获引用或移动所有权”中,你学过如何在闭包上使用 move 关键字;而正如第十六章“将 move 闭包与线程一同使用”一节提到的那样,在线程场景下我们也经常需要把数据 move 进闭包。相同的基本原理也适用于 async 代码块,因此 move 关键字同样可以和 async 代码块一起使用。

在示例 17-12 中,我们把发送消息用的代码块从 async 改为 async move

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let tx_fut = async move {
            // --snip--
            let vals = vec![
                String::from("hi"),
                String::from("from"),
                String::from("the"),
                String::from("future"),
            ];

            for val in vals {
                tx.send(val).unwrap();
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        let rx_fut = async {
            while let Some(value) = rx.recv().await {
                println!("received '{value}'");
            }
        };

        trpl::join(tx_fut, rx_fut).await;
    });
}
示例 17-12:对示例 17-11 的修改版本,它会在完成后正确关闭

运行这个版本的代码后,它就会在最后一条消息发送并接收完之后正常退出。接下来,我们来看看,如果要从多个 future 发送数据,又需要做哪些变化。

使用 join! 宏合并多个 future

这个异步信道同样也是多生产者信道,因此如果我们希望从多个 future 发送消息,就可以对 tx 调用 clone,如示例 17-13 所示。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let tx1 = tx.clone();
        let tx1_fut = async move {
            let vals = vec![
                String::from("hi"),
                String::from("from"),
                String::from("the"),
                String::from("future"),
            ];

            for val in vals {
                tx1.send(val).unwrap();
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        let rx_fut = async {
            while let Some(value) = rx.recv().await {
                println!("received '{value}'");
            }
        };

        let tx_fut = async move {
            let vals = vec![
                String::from("more"),
                String::from("messages"),
                String::from("for"),
                String::from("you"),
            ];

            for val in vals {
                tx.send(val).unwrap();
                trpl::sleep(Duration::from_millis(1500)).await;
            }
        };

        trpl::join!(tx1_fut, tx_fut, rx_fut);
    });
}
示例 17-13:在 async 代码块中使用多个生产者

首先,我们克隆 tx,在第一个 async 代码块外创建出 tx1。然后像之前处理 tx 那样,把 tx1 move 进这个代码块里。随后,我们再把原始的 tx move 进一个新的 async 代码块,在那里以稍慢一点的节奏继续发送更多消息。这里我们把这个新 async 代码块放在接收消息的 async 代码块后面,不过放在前面也同样可以。关键在于 future 被等待的顺序,而不是它们被创建的顺序。

两个负责发送消息的 async 代码块都必须写成 async move,这样当代码块结束时,txtx1 都会被 drop。否则,我们又会回到一开始那个无限循环的问题。

最后,我们从 trpl::join 切换为 trpl::join! 来处理新增的 future。join! 宏可以在 future 数量已知于编译期的情况下,等待任意数量的 future。本章稍后我们还会讨论,如何等待一个数量事先未知的 future 集合。

现在我们就能看到来自两个发送 future 的所有消息了。由于这两个发送 future 在发送后使用了略微不同的延迟,接收到这些消息的时间间隔也会相应不同:

received 'hi'
received 'more'
received 'from'
received 'the'
received 'messages'
received 'future'
received 'for'
received 'you'

我们已经探索了如何用消息传递在 future 之间发送数据、一个 async 代码块中的代码如何按顺序执行、如何将所有权 move 进 async 代码块,以及如何合并多个 future。接下来,我们来讨论一下,为什么以及如何告诉运行时:它现在可以切换去执行别的任务了。

使用任意数量的 futures

将控制权交还给运行时

回忆一下我们在“第一个异步程序”一节中提到的内容:在每个 await 点,如果被等待的 future 还没准备好,Rust 就会给运行时一个机会来暂停当前任务并切换到其他任务。反过来也成立:Rust 只会 在 await 点暂停 async 代码块,并把控制权交还给运行时。await 点之间的所有内容都是同步执行的。

这意味着,如果你在一个 async 代码块中做了大量工作,却没有任何 await 点,那么这个 future 就会阻止其他 future 取得进展。有时你会听到人们把这称为一个 future 让其他 future starve(饥饿)。在某些场景下,这也许不是什么大问题;但如果你在做昂贵的初始化、长时间运行的工作,或者你有一个会无限执行某项任务的 future,就需要认真考虑该在何时、何地把控制权交还给运行时。

让我们通过模拟一个长时间运行的操作来展示这种“饥饿”问题,再看看该如何解决。示例 17-14 引入了一个 slow 函数。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::{thread, time::Duration};

fn main() {
    trpl::block_on(async {
        // We will call `slow` here later
    });
}

fn slow(name: &str, ms: u64) {
    thread::sleep(Duration::from_millis(ms));
    println!("'{name}' ran for {ms}ms");
}
示例 17-14:使用 `thread::sleep` 来模拟缓慢的操作

这段代码使用的是 std::thread::sleep,而不是 trpl::sleep,因此调用 slow 会让当前线程阻塞若干毫秒。我们可以把 slow 看作现实世界中那些既耗时又会阻塞的操作的替身。

在示例 17-15 中,我们用 slow 来模拟在一对 future 中执行这类 CPU 密集型工作。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::{thread, time::Duration};

fn main() {
    trpl::block_on(async {
        let a = async {
            println!("'a' started.");
            slow("a", 30);
            slow("a", 10);
            slow("a", 20);
            trpl::sleep(Duration::from_millis(50)).await;
            println!("'a' finished.");
        };

        let b = async {
            println!("'b' started.");
            slow("b", 75);
            slow("b", 10);
            slow("b", 15);
            slow("b", 350);
            trpl::sleep(Duration::from_millis(50)).await;
            println!("'b' finished.");
        };

        trpl::select(a, b).await;
    });
}

fn slow(name: &str, ms: u64) {
    thread::sleep(Duration::from_millis(ms));
    println!("'{name}' ran for {ms}ms");
}
示例 17-15:调用 `slow` 函数来模拟缓慢操作

每个 future 都会在完成一大串缓慢操作之后,才把控制权交还给运行时。如果你运行这段代码,就会看到如下输出:

'a' started.
'a' ran for 30ms
'a' ran for 10ms
'a' ran for 20ms
'b' started.
'b' ran for 75ms
'b' ran for 10ms
'b' ran for 15ms
'b' ran for 350ms
'a' finished.

和示例 17-5 中用 trpl::select 让两个 URL 获取任务竞争时一样,select 仍然会在 a 完成时立刻结束。不过,这两个 future 里的 slow 调用之间完全没有交错执行。a future 会一路把自己的工作做完,直到等待 trpl::sleep 调用;接着 b future 又一路做完自己的工作,直到它自己的 trpl::sleep 被等待;最后 a future 才完成。要想让两个 future 在这些缓慢任务之间都取得进展,我们就需要一些 await 点,好把控制权交还给运行时。也就是说,我们得有某种可以被 await 的东西!

实际上,我们已经能在示例 17-15 中看到这种“交接”是如何发生的:如果去掉 a future 末尾的 trpl::sleep,那么它会直接完成,而 b future 根本不会运行。让我们先从 trpl::sleep 入手,试着让这些操作能够轮流取得进展,如示例 17-16 所示。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::{thread, time::Duration};

fn main() {
    trpl::block_on(async {
        let one_ms = Duration::from_millis(1);

        let a = async {
            println!("'a' started.");
            slow("a", 30);
            trpl::sleep(one_ms).await;
            slow("a", 10);
            trpl::sleep(one_ms).await;
            slow("a", 20);
            trpl::sleep(one_ms).await;
            println!("'a' finished.");
        };

        let b = async {
            println!("'b' started.");
            slow("b", 75);
            trpl::sleep(one_ms).await;
            slow("b", 10);
            trpl::sleep(one_ms).await;
            slow("b", 15);
            trpl::sleep(one_ms).await;
            slow("b", 350);
            trpl::sleep(one_ms).await;
            println!("'b' finished.");
        };

        trpl::select(a, b).await;
    });
}

fn slow(name: &str, ms: u64) {
    thread::sleep(Duration::from_millis(ms));
    println!("'{name}' ran for {ms}ms");
}
示例 17-16:使用 `trpl::sleep` 让操作轮流推进

我们在每次调用 slow 之间都插入了 trpl::sleep 调用和 await 点。现在,两个 future 的工作就交错在一起了:

'a' started.
'a' ran for 30ms
'b' started.
'b' ran for 75ms
'a' ran for 10ms
'b' ran for 10ms
'a' ran for 20ms
'b' ran for 15ms
'a' finished.

a future 仍然会在第一次把控制权交给 b 之前先运行一阵,因为它是在第一次调用 trpl::sleep 之前先执行了 slow;但在那之后,每当其中一个 future 命中 await 点,它们就会来回切换。在这个例子中,我们是在每次 slow 之后这么做的,不过实际上也可以按任何对你最合理的方式来拆分工作。

但我们其实并不是真的想在这里休眠sleep);我们只是希望程序尽可能快地前进。我们真正需要的只是把控制权交还给运行时。可以直接通过 trpl::yield_now 来做到这一点。在示例 17-17 中,我们把前面的 trpl::sleep 全部替换成 trpl::yield_now

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::{thread, time::Duration};

fn main() {
    trpl::block_on(async {
        let a = async {
            println!("'a' started.");
            slow("a", 30);
            trpl::yield_now().await;
            slow("a", 10);
            trpl::yield_now().await;
            slow("a", 20);
            trpl::yield_now().await;
            println!("'a' finished.");
        };

        let b = async {
            println!("'b' started.");
            slow("b", 75);
            trpl::yield_now().await;
            slow("b", 10);
            trpl::yield_now().await;
            slow("b", 15);
            trpl::yield_now().await;
            slow("b", 350);
            trpl::yield_now().await;
            println!("'b' finished.");
        };

        trpl::select(a, b).await;
    });
}

fn slow(name: &str, ms: u64) {
    thread::sleep(Duration::from_millis(ms));
    println!("'{name}' ran for {ms}ms");
}
示例 17-17:使用 `yield_now` 让操作轮流推进

这段代码不仅更清楚地表达了真实意图,而且通常也会比使用 sleep 快得多,因为像 sleep 使用的那类计时器,往往都会受到最小粒度限制。比如我们这里使用的 sleep,即使你传入的是 1 纳秒的 Duration,它也至少会睡眠 1 毫秒。再说一次,现代计算机是非常的:1 毫秒里已经能完成大量工作了!

这说明:即使面对计算密集型任务,async 依然可能是有用的,具体取决于你的程序还在做什么,因为它提供了一种很实用的手段,用来组织程序不同部分之间的关系(当然代价是 async 状态机本身也有一定开销)。这是一种 协作式多任务cooperative multitasking):每个 future 都可以通过 await 点来决定何时交出控制权,因此每个 future 也都负有避免长时间阻塞的责任。在某些基于 Rust 的嵌入式操作系统中,这甚至是唯一的多任务形式!

当然,在真实代码里,你通常不会在每一行之间都交替插入函数调用和 await 点。像这样主动交出控制权虽然相对便宜,但并不是没有代价。在很多情况下,试图把一个计算密集型任务切得太碎,反而可能显著拖慢它的执行速度,所以有时为了整体性能,让某个操作短暂阻塞一下反而更好。还是那句老话:一定要靠测量来确认代码真正的性能瓶颈。不过,如果你发现原本预期会并发执行的工作,实际上却大量串行发生,那么就要记住这里的底层机制。

构建我们自己的异步抽象

我们还可以把多个 future 组合在一起,创造出新的模式。比如,完全可以用手头已有的 async 构件来写一个 timeout 函数。等我们做完,它本身就会成为另一个可以继续拿来构建更多 async 抽象的基础模块。

示例 17-18 展示了我们希望这个 timeout 在面对一个缓慢 future 时应该如何工作。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let slow = async {
            trpl::sleep(Duration::from_secs(5)).await;
            "Finally finished"
        };

        match timeout(slow, Duration::from_secs(2)).await {
            Ok(message) => println!("Succeeded with '{message}'"),
            Err(duration) => {
                println!("Failed after {} seconds", duration.as_secs())
            }
        }
    });
}
示例 17-18:使用我们设想的 `timeout` 为一个缓慢操作设置时间限制

让我们来实现它。首先先想想 timeout 的 API:

  • 它本身需要是一个 async 函数,这样我们才能等待它。
  • 它的第一个参数应该是一个要执行的 future。我们可以把它设计成泛型,从而支持任意 future。
  • 它的第二个参数应该是最大等待时间。如果用 Duration,就能很方便地直接传给 trpl::sleep
  • 它应该返回一个 Result。如果传入的 future 成功完成,结果应当是 Ok,内部带上该 future 产生的值;如果先发生超时,结果就应当是 Err,内部带上等待的时长。

示例 17-19 展示了这个声明。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let slow = async {
            trpl::sleep(Duration::from_secs(5)).await;
            "Finally finished"
        };

        match timeout(slow, Duration::from_secs(2)).await {
            Ok(message) => println!("Succeeded with '{message}'"),
            Err(duration) => {
                println!("Failed after {} seconds", duration.as_secs())
            }
        }
    });
}

async fn timeout<F: Future>(
    future_to_try: F,
    max_time: Duration,
) -> Result<F::Output, Duration> {
    // Here is where our implementation will go!
}
示例 17-19:定义 `timeout` 的签名

这样类型层面的目标就满足了。接下来想想我们需要的行为:我们希望让传入的 future 和这个时长“竞争”。可以用 trpl::sleep 根据这个时长构造一个计时器 future,再用 trpl::select 让计时器和调用者传入的 future 一起运行。

在示例 17-20 中,我们通过匹配 trpl::select 的等待结果来实现 timeout

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::time::Duration;

use trpl::Either;

// --snip--

fn main() {
    trpl::block_on(async {
        let slow = async {
            trpl::sleep(Duration::from_secs(5)).await;
            "Finally finished"
        };

        match timeout(slow, Duration::from_secs(2)).await {
            Ok(message) => println!("Succeeded with '{message}'"),
            Err(duration) => {
                println!("Failed after {} seconds", duration.as_secs())
            }
        }
    });
}

async fn timeout<F: Future>(
    future_to_try: F,
    max_time: Duration,
) -> Result<F::Output, Duration> {
    match trpl::select(future_to_try, trpl::sleep(max_time)).await {
        Either::Left(output) => Ok(output),
        Either::Right(_) => Err(max_time),
    }
}
示例 17-20:使用 `select` 和 `sleep` 定义 `timeout`

trpl::select 的实现并不是公平的:它总是按参数传入的顺序进行轮询(其他一些 select 实现会随机选择先轮询哪个参数)。因此,我们把 future_to_try 作为第一个参数传给 select,好让它即使在 max_time 很短的情况下,也仍然有机会先完成。如果 future_to_try 先完成,select 会返回 Left,其中包含 future_to_try 的输出;如果 timer 先完成,select 就会返回 Right,其中包含计时器的输出 ()

如果 future_to_try 成功完成,并且我们得到了 Left(output),那么就返回 Ok(output)。如果相反是睡眠计时器先结束,我们得到 Right(()),那就用 _ 忽略这个 (),并返回 Err(max_time)

这样一来,我们就用另外两个 async 小工具拼出了一个可工作的 timeout。如果运行代码,它会在超时后打印出失败信息:

Failed after 2 seconds

因为 future 可以和其他 future 组合,你就能利用更小的 async 构件构建出非常强大的工具。比如,完全可以用同样的方法把 timeout 和 retry 组合起来,再进一步把它们用在网络请求一类的操作上(例如示例 17-5 里的那些)。

在实践中,你通常会主要直接使用 asyncawait,其次才是像 select 这样的函数,以及像 join! 这样的宏,来控制最外层的 future 应该如何执行。

到这里,我们已经看过好几种同时处理多个 future 的方式了。接下来,我们将看看如何借助 stream,按照时间顺序处理一串 future。

Stream:按顺序出现的 Future

Stream:按顺序出现的 Future

回忆一下本章前面在“通过消息传递在两个任务之间发送数据”一节中,我们是如何使用异步信道的接收端的。异步版的 recv 方法会随着时间推移产出一系列条目。这正是一种更普遍模式的实例,通常称为 stream)。很多概念都很自然地适合表示成 stream:队列中逐步变得可用的项、当完整数据集太大而无法一次装入内存时从文件系统中逐块拉取的数据、或者随着时间逐渐从网络到达的数据。由于 stream 本身也和 future 密切相关,我们可以把它和其他类型的 future 一起使用,并以有趣的方式进行组合。比如,我们可以把事件分批处理,以避免触发过多网络调用;可以为一串长时间运行的操作设置超时;也可以对 UI 事件进行节流,避免做无谓的工作。

我们在第十三章“Iterator trait 和 next 方法”一节中已经见过“按顺序产生一系列项”这回事,但迭代器和异步信道接收端之间有两个区别。第一个区别是时间:迭代器是同步的,而信道接收端是异步的。第二个区别是 API。直接处理 Iterator 时,我们会调用同步的 next 方法;而对于 trpl::Receiver 这个具体的 stream 来说,我们调用的是异步的 recv 方法。除此之外,这些 API 给人的感觉非常相似,而这种相似并非巧合。stream 就像迭代的一种异步形式。不过,trpl::Receiver 专门用于等待接收消息,而更通用的 stream API 则宽泛得多:它像 Iterator 一样提供“下一个条目”,只不过是以异步方式来做。

Rust 中迭代器和 stream 的这种相似性意味着,我们实际上可以从任意迭代器创建一个 stream。和使用迭代器一样,我们也可以通过调用 stream 的 next 方法,再等待其输出,来处理它,如示例 17-21 所示。不过这段代码暂时还编译不过。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

fn main() {
    trpl::block_on(async {
        let values = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
        let iter = values.iter().map(|n| n * 2);
        let mut stream = trpl::stream_from_iter(iter);

        while let Some(value) = stream.next().await {
            println!("The value was: {value}");
        }
    });
}
示例 17-21:从迭代器创建一个 stream,并打印出其中的值

我们从一个数字数组开始,把它转换为迭代器,然后调用 map 把其中所有值都翻倍。接着再使用 trpl::stream_from_iter 函数,把这个迭代器转换成一个 stream。最后,我们用 while let 循环,在 stream 中的值陆续到达时逐个处理它们。

遗憾的是,当我们尝试运行这段代码时,编译器并不会通过,而是报告说没有可用的 next 方法:

error[E0599]: no method named `next` found for struct `tokio_stream::iter::Iter` in the current scope
  --> src/main.rs:10:40
   |
10 |         while let Some(value) = stream.next().await {
   |                                        ^^^^
   |
   = help: items from traits can only be used if the trait is in scope
help: the following traits which provide `next` are implemented but not in scope; perhaps you want to import one of them
   |
1  + use crate::trpl::StreamExt;
   |
1  + use futures_util::stream::stream::StreamExt;
   |
1  + use std::iter::Iterator;
   |
1  + use std::str::pattern::Searcher;
   |
help: there is a method `try_next` with a similar name
   |
10 |         while let Some(value) = stream.try_next().await {
   |                                        ~~~~~~~~

正如这段输出解释的那样,编译错误的原因是:我们需要把正确的 trait 放进作用域,才能使用 next 方法。根据前面的讨论,你很可能会合理地猜测这个 trait 应该是 Stream,但实际上它是 StreamExt。这里的 Extextension 的缩写;在 Rust 社区里,用一个 trait 去扩展另一个 trait,是非常常见的模式。

Stream trait 定义的是一个底层接口,它实际上把 IteratorFuture trait 的特征结合在了一起。StreamExt 则在 Stream 之上提供了一组更高层的 API,其中包括 next 方法,以及其他一些和 Iterator trait 提供的工具方法相似的辅助方法。StreamStreamExt 目前都还不是 Rust 标准库的一部分,不过生态系统中的大多数 crate 都使用相似的定义。

修复这个编译错误的方式,就是像示例 17-22 那样,添加一条 use trpl::StreamExt 语句。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use trpl::StreamExt;

fn main() {
    trpl::block_on(async {
        let values = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
        // --snip--
        let iter = values.iter().map(|n| n * 2);
        let mut stream = trpl::stream_from_iter(iter);

        while let Some(value) = stream.next().await {
            println!("The value was: {value}");
        }
    });
}
示例 17-22:成功把迭代器作为 stream 的基础来使用

把这些部分拼起来之后,这段代码就会按我们想要的方式工作!更重要的是,既然我们已经把 StreamExt 引入作用域,就也能像使用迭代器时那样,使用它提供的整套工具方法。

深入理解 async 相关的 traits

深入理解 async 相关的 trait

贯穿本章,我们以各种方式使用了 FutureStreamStreamExt trait。不过到目前为止,我们一直刻意没有太深入它们究竟是如何工作的、又是如何彼此配合的。对日常 Rust 编程来说,这通常完全没问题。不过有时你会遇到一些场景,在那里你需要额外理解这些 trait 的更多细节,以及 Pin 类型和 Unpin trait。在这一节里,我们会适度深入,足够帮助你应对这些情况,但把真正深入的内容留给其他文档。

Future trait

让我们先更仔细地看看 Future trait 是如何工作的。Rust 中它的定义如下:

#![allow(unused)]
fn main() {
use std::pin::Pin;
use std::task::{Context, Poll};

pub trait Future {
    type Output;

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>;
}
}

这个 trait 定义里包含了不少新类型,也有一些我们之前还没见过的语法,所以我们逐部分来看。

首先,Future 的关联类型 Output 指明了这个 future 最终会解析成什么值。这和 Iterator trait 里的关联类型 Item 是类似的。其次,Future 提供了一个 poll 方法。它接收一个特殊的 Pin 包裹的 self 引用、一个指向 Context 类型的可变引用,并返回 Poll<Self::Output>。稍后我们会再讲 PinContext。现在,先聚焦到这个方法的返回值 Poll

#![allow(unused)]
fn main() {
pub enum Poll<T> {
    Ready(T),
    Pending,
}
}

这个 Poll 类型有点像 Option。它也有一个带值的变体 Ready(T),以及一个不带值的变体 Pending。但 Poll 的语义和 Option 完全不同。Pending 表示这个 future 还有工作没做完,因此调用方稍后还需要再次检查。Ready 则表示这个 Future 已经完成,其结果值 T 现在已经可用。

注意:直接调用 poll 的场景很少,但如果你真的需要这么做,请记住:对于大多数 future 来说,一旦它已经返回过 Ready,调用方就不应再对它调用 poll。很多 future 在 ready 之后再次被轮询时会 panic。那些可以安全重复轮询的 future,会在文档里明确说明。这和 Iterator::next 的行为有些相似。

当你看到使用 await 的代码时,Rust 在底层会把它编译成调用 poll 的代码。如果你回头看示例 17-4,也就是在单个 URL 的标题解析完成后把它打印出来的那个例子,Rust 编译出来的代码大致会像下面这样(虽然并不完全一致):

match page_title(url).poll() {
    Ready(page_title) => match page_title {
        Some(title) => println!("The title for {url} was {title}"),
        None => println!("{url} had no title"),
    }
    Pending => {
        // 这里该怎么办?
    }
}

如果 future 仍然是 Pending,那我们该怎么办?我们需要一种办法不断重试,直到 future 最终准备好。换句话说,我们需要一个循环:

let mut page_title_fut = page_title(url);
loop {
    match page_title_fut.poll() {
        Ready(value) => match page_title {
            Some(title) => println!("The title for {url} was {title}"),
            None => println!("{url} had no title"),
        }
        Pending => {
            // continue
        }
    }
}

但如果 Rust 真按这段代码精确地编译,那么每个 await 就都会变成阻塞式的,这恰恰和我们想要的效果相反!Rust 实际上会保证:这个循环能够把控制权交给某个东西,由它暂停当前 future 的工作,去处理别的 future,然后稍后再回来重新检查当前这个。正如我们已经见过的,这个“某个东西”就是异步运行时,而调度和协调这些工作,正是运行时的核心职责之一。

“通过消息传递在两个任务之间发送数据”一节中,我们描述过等待 rx.recv 的过程。recv 调用会返回一个 future,而等待这个 future 本质上就是在轮询它。我们之前提到,运行时会暂停这个 future,直到它准备好,最终要么得到 Some(message),要么在信道关闭时得到 None。现在,借助对 Future trait,尤其是 Future::poll 的更深入理解,我们就能看清它的工作方式了:当返回 Poll::Pending 时,运行时知道这个 future 还没准备好;反过来,当 poll 返回 Poll::Ready(Some(message))Poll::Ready(None) 时,运行时就知道这个 future 已经准备好,可以继续推进它。

至于运行时具体是怎么做到这一点的,已经超出了本书的范围。不过关键是看清 future 的基本机制:运行时会去轮询它所负责的每个 future,而当 future 还没准备好时,就让它重新休眠。

Pin 类型与 Unpin trait

回到示例 17-13,我们使用过 trpl::join! 宏来等待三个 future。不过,更常见的情况是你会有一个集合,比如一个向量,其中包含若干个 future,而这些 future 的个数要到运行时才知道。让我们把示例 17-13 改成示例 17-23 中的代码:把这三个 future 放进一个向量里,再调用 trpl::join_all。不过,这段代码暂时还编译不过。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let tx1 = tx.clone();
        let tx1_fut = async move {
            let vals = vec![
                String::from("hi"),
                String::from("from"),
                String::from("the"),
                String::from("future"),
            ];

            for val in vals {
                tx1.send(val).unwrap();
                trpl::sleep(Duration::from_secs(1)).await;
            }
        };

        let rx_fut = async {
            while let Some(value) = rx.recv().await {
                println!("received '{value}'");
            }
        };

        let tx_fut = async move {
            // --snip--
            let vals = vec![
                String::from("more"),
                String::from("messages"),
                String::from("for"),
                String::from("you"),
            ];

            for val in vals {
                tx.send(val).unwrap();
                trpl::sleep(Duration::from_secs(1)).await;
            }
        };

        let futures: Vec<Box<dyn Future<Output = ()>>> =
            vec![Box::new(tx1_fut), Box::new(rx_fut), Box::new(tx_fut)];

        trpl::join_all(futures).await;
    });
}
示例 17-23:等待一个集合中的多个 future

我们把每个 future 都放进了一个 Box 中,好把它们变成 trait object,就像我们在第十二章“从 run 返回错误”那一节做的那样。(我们会在第十八章详细讨论 trait object。)使用 trait object 后,我们就能把这些类型各不相同的匿名 future 当成同一种类型来对待,因为它们全都实现了 Future trait。

这也许会让人意外。毕竟,这些 async 代码块都没有返回任何值,所以它们每一个产生的都是 Future<Output = ()>。但别忘了:Future 是个 trait,而编译器会为每个 async 代码块生成一个独一无二的 enum,即使它们的输出类型完全相同。就像你不能把两个不同的手写 struct 放进同一个 Vec,你也同样不能把这些编译器生成的不同 enum 混在一起。

然后,我们把这组 future 传给 trpl::join_all,再等待结果。然而,这段代码仍然无法编译。下面是报错中最关键的一部分:

error[E0277]: `dyn Future<Output = ()>` cannot be unpinned
  --> src/main.rs:48:33
   |
48 |         trpl::join_all(futures).await;
   |                                 ^^^^^ the trait `Unpin` is not implemented for `dyn Future<Output = ()>`
   |
   = note: consider using the `pin!` macro
           consider using `Box::pin` if you need to access the pinned value outside of the current scope
   = note: required for `Box<dyn Future<Output = ()>>` to implement `Future`
note: required by a bound in `futures_util::future::join_all::JoinAll`
  --> file:///home/.cargo/registry/src/index.crates.io-1949cf8c6b5b557f/futures-util-0.3.30/src/future/join_all.rs:29:8
   |
27 | pub struct JoinAll<F>
   |            ------- required by a bound in this struct
28 | where
29 |     F: Future,
   |        ^^^^^^ required by this bound in `JoinAll`

这段错误信息里的 note 告诉我们,应该使用 pin! 宏来 pin 这些值,也就是把它们放进 Pin 类型中,以保证这些值不会在内存中移动。报错之所以说需要 pin,是因为 dyn Future<Output = ()> 需要实现 Unpin trait,而它当前并没有实现。

trpl::join_all 返回的是一个名为 JoinAll 的结构体。这个结构体在类型参数 F 上是泛型的,而 F 又被约束必须实现 Future trait。直接通过 await 去等待一个 future 时,Rust 会隐式地把它 pin 住。这也正是为什么我们平常不需要在每个想等待 future 的地方都显式写 pin!

但这里,我们并不是直接在等待某个 future。相反,我们是通过把一组 future 传给 join_all,构造出了一个新的 future:JoinAll。而 join_all 的签名要求集合中的元素类型都必须实现 Future trait。另一方面,Box<T> 只有在它包裹的 T 本身是 future 且实现了 Unpin trait 时,才会实现 Future

这一下信息量很大!为了真正理解它,我们得再更深入一点,看清 Future trait 尤其是 pinning 这一部分到底是如何运作的。再看一遍 Future trait 的定义:

#![allow(unused)]
fn main() {
use std::pin::Pin;
use std::task::{Context, Poll};

pub trait Future {
    type Output;

    // Required method
    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>;
}
}

这里的 cx 参数以及它的 Context 类型,是运行时在保持 lazy 的同时,真正知道该在什么时候重新检查某个 future 的关键。和前面一样,这部分具体机制超出了本章范围,而且通常也只有在你自己实现 Future 时才需要关注。我们这里聚焦的是 self 的类型,因为这是我们第一次见到一个方法里的 self 带有类型注解。对 self 进行类型注解,和给其他函数参数写类型注解类似,但有两个关键区别:

  • 它告诉 Rust:要调用这个方法,self 必须是什么类型。
  • 它不能随便写成任意类型。它必须是方法所实现类型本身、该类型的引用或智能指针,或者是一个包裹了该类型引用的 Pin

我们会在第十八章里看到更多相关语法。眼下,只要知道:如果我们想通过轮询 future 来检查它到底是 Pending 还是 Ready(Output),那么就需要一个 Pin 包裹的、指向该类型的可变引用。

Pin 是一种针对指针类类型的包装器,比如 &&mutBoxRc。(严格来说,Pin 作用于实现了 DerefDerefMut 的类型,但实际效果基本等同于“引用和智能指针”。)Pin 本身并不是指针,也不像 RcArc 那样自带引用计数之类的行为;它纯粹是一个让编译器能够对指针使用方式施加约束的工具。

回忆一下:await 是通过调用 poll 实现的。理解这一点以后,前面的错误信息就已经开始变得容易理解了,不过那个报错说的是 Unpin,不是 Pin。那么,PinUnpin 究竟是什么关系?为什么 Future 又要求 self 必须放在 Pin 里才能调用 poll 呢?

记住,我们在本章前面提过,一个 future 里的多个 await 点会被编译成一个状态机,而编译器会确保这个状态机遵守 Rust 关于安全性的全部常规规则,包括借用和所有权。为了做到这一点,Rust 会分析:在某个 await 点和下一个 await 点之间,或者直到 async 代码块结束之前,哪些数据是需要保留的。然后,它会在编译出来的状态机里生成对应的变体。每个变体都会得到其对应源代码片段所需的数据访问权限,这种访问可能是获得所有权,也可能是获得可变或不可变引用。

到这里为止,一切都很好:如果你在某个 async 代码块里把所有权或引用关系写错了,借用检查器会告诉你。但当我们想要移动这个代码块对应的 future 时,比如把它放进 Vec 然后传给 join_all,事情就开始变复杂了。

当我们移动一个 future 时,无论是把它放进数据结构,以便通过 join_all 这种方式迭代处理,还是从函数里返回它,本质上都是在移动 Rust 为我们生成的那个状态机。与 Rust 中大多数其他类型不同的是,Rust 为 async 代码块生成的 future,可能会在某个状态变体的字段里保存指向它自身其他字段的引用,就像图 17-4 里的简化示意图那样。

A single-column, three-row table representing a future, fut1, which has data values 0 and 1 in the first two rows and an arrow pointing from the third row back to the second row, representing an internal reference within the future.
图 17-4:一个自引用的数据类型

但默认情况下,任何包含自引用的对象,一旦移动就是不安全的,因为引用始终指向它们所引用对象的真实内存地址(见图 17-5)。如果我们移动了这个数据结构本身,那么这些内部引用仍然会指向旧位置。然而那个内存地址现在已经失效了。一方面,你之后对数据结构做的修改不会再反映到那些旧引用上;另一方面,更严重的是,计算机此时已经可以把那块内存拿去做别的用途了。最后你很可能会读到完全无关的数据。

Two tables, depicting two futures, fut1 and fut2, each of which has one column and three rows, representing the result of having moved a future out of fut1 into fut2. The first, fut1, is grayed out, with a question mark in each index, representing unknown memory. The second, fut2, has 0 and 1 in the first and second rows and an arrow pointing from its third row back to the second row of fut1, representing a pointer that is referencing the old location in memory of the future before it was moved.
图 17-5:移动自引用数据类型后产生的不安全结果

理论上,Rust 编译器也可以尝试在对象被移动时更新所有引用,但这样很可能带来大量性能开销,尤其在需要更新的是一整张引用网络的时候。如果我们反过来,确保这个数据结构根本不在内存中移动,那就完全不需要更新任何引用。这正是 Rust 借用检查器要做的事:在安全代码里,它会阻止你移动任何仍然存在活动引用的值。

Pin 正是在这个基础上,进一步提供了我们需要的精确保证。当我们把一个指向某值的指针包进 Pin 里,也就是对这个值进行 pin 之后,它就不能再被移动了。因此,如果你有的是 Pin<Box<SomeType>>,那么真正被 pin 住的是 SomeType 这个值,而不是 Box 指针本身。图 17-6 展示了这个过程。

Three boxes laid out side by side. The first is labeled “Pin”, the second “b1”, and the third “pinned”. Within “pinned” is a table labeled “fut”, with a single column; it represents a future with cells for each part of the data structure. Its first cell has the value “0”, its second cell has an arrow coming out of it and pointing to the fourth and final cell, which has the value “1” in it, and the third cell has dashed lines and an ellipsis to indicate there may be other parts to the data structure. All together, the “fut” table represents a future which is self-referential. An arrow leaves the box labeled “Pin”, goes through the box labeled “b1” and terminates inside the “pinned” box at the “fut” table.
图 17-6:把一个指向自引用 future 类型的 `Box` pin 住

实际上,Box 指针本身仍然可以自由移动。请记住:我们真正关心的是最终被引用的数据必须固定不动。如果指针移动了,但它指向的数据仍然留在原地,就像图 17-7 那样,那么就不会产生问题。(你可以把这当作一个独立练习:去查阅相关类型以及 std::pin 模块的文档,试着想清楚如果是 Pin 包着 Box,到底如何做到这一点。)关键在于:那个自引用的类型本身不能移动,因为它仍然是被 pin 住的。

Four boxes laid out in three rough columns, identical to the previous diagram with a change to the second column. Now there are two boxes in the second column, labeled “b1” and “b2”, “b1” is grayed out, and the arrow from “Pin” goes through “b2” instead of “b1”, indicating that the pointer has moved from “b1” to “b2”, but the data in “pinned” has not moved.
图 17-7:移动一个指向自引用 future 类型的 `Box`

不过,大多数类型即使碰巧放在 Pin 指针后面,也完全可以安全移动。只有当某个值内部真的包含引用时,我们才需要关心 pin。比如数字和布尔值这类基本类型显然没有内部引用,所以当然是安全的。你平时在 Rust 里处理的大多数类型也都是这样。比如一个 Vec 就可以自由移动而不用担心。考虑到目前为止我们看到的内容,如果你有一个 Pin<Vec<String>>,那理论上你必须通过 Pin 提供的那套安全但受限的 API 来操作它,哪怕 Vec<String> 在没有其他引用存在时始终都是可以安全移动的。因此,我们需要一种机制来告诉编译器:像这种情况,移动它完全没问题。这正是 Unpin 的用途。

Unpin 是一个标记 trait(marker trait),就像我们在第十六章见过的 SendSync 一样,它本身没有任何功能。marker trait 的存在,只是为了告诉编译器:实现了该 trait 的类型,在某种特定上下文里可以被安全使用。Unpin 告诉编译器,某个类型不需要维护“这个值是否可以安全移动”方面的额外保证。

就像 SendSync 一样,只要编译器能证明某个类型这样做是安全的,它就会自动为其实现 Unpin。同样也存在一个特殊情况:某个类型不会实现 Unpin。这种写法是 impl !Unpin for SomeType,其中 SomeType 表示的是:为了在被 Pin 指针引用时保持安全,该类型必须保证自身不会被移动。

换句话说,关于 PinUnpin 的关系,有两件事要记住。第一,Unpin 才是“正常情况”,!Unpin 才是特殊情况。第二,一个类型到底实现的是 Unpin 还是 !Unpin只有在你使用像 Pin<&mut SomeType> 这样指向该类型的 pin 过的指针时,才真正有意义。

为了更具体一点,想想 String。它内部保存的是长度以及组成它的 Unicode 字符。我们完全可以把一个 String 包进 Pin,如图 17-8 所示。不过,String 会自动实现 Unpin,Rust 中绝大多数其他类型也一样。

A box labeled “Pin” on the left with an arrow going from it to a box labeled “String” on the right. The “String” box contains the data 5usize, representing the length of the string, and the letters “h”, “e”, “l”, “l”, and “o” representing the characters of the string “hello” stored in this String instance. A dotted rectangle surrounds the “String” box and its label, but not the “Pin” box.
图 17-8:把一个 `String` pin 起来;虚线表示 `String` 实现了 `Unpin` trait,因此它实际上并没有被固定住

结果就是,我们可以做一些如果 String 实现的是 !Unpin 就会非法的事情,比如像图 17-9 那样,在同一块内存位置上把一个字符串直接替换成另一个完全不同的字符串。这并没有违反 Pin 的约定,因为 String 内部没有那种会让它在移动时变得不安全的自引用。也正因为如此,它实现的是 Unpin,而不是 !Unpin

The same “hello” string data from the previous example, now labeled “s1” and grayed out. The “Pin” box from the previous example now points to a different String instance, one that is labeled “s2”, is valid, has a length of 7usize, and contains the characters of the string “goodbye”. s2 is surrounded by a dotted rectangle because it, too, implements the Unpin trait.
图 17-9:在内存中用一个完全不同的 `String` 替换原来的 `String`

到这里,我们已经知道得足够多,可以理解前面示例 17-23 中那个 join_all 调用为什么会报错了。我们最初试图把 async 代码块生成的 future 移动进 Vec<Box<dyn Future<Output = ()>>> 中,但正如我们刚刚看到的,那些 future 可能带有内部引用,因此它们不会自动实现 Unpin。一旦把它们 pin 住,我们就可以放心地把得到的 Pin 类型放进 Vec,因为此时这些 future 底层的数据就不会再被移动。示例 17-24 展示了修复这段代码的方法:在定义每个 future 的地方调用 pin! 宏,并相应调整 trait object 的类型。

文件名:src/main.rs

extern crate trpl; // required for mdbook test

use std::pin::{Pin, pin};

// --snip--

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let tx1 = tx.clone();
        let tx1_fut = pin!(async move {
            // --snip--
            let vals = vec![
                String::from("hi"),
                String::from("from"),
                String::from("the"),
                String::from("future"),
            ];

            for val in vals {
                tx1.send(val).unwrap();
                trpl::sleep(Duration::from_secs(1)).await;
            }
        });

        let rx_fut = pin!(async {
            // --snip--
            while let Some(value) = rx.recv().await {
                println!("received '{value}'");
            }
        });

        let tx_fut = pin!(async move {
            // --snip--
            let vals = vec![
                String::from("more"),
                String::from("messages"),
                String::from("for"),
                String::from("you"),
            ];

            for val in vals {
                tx.send(val).unwrap();
                trpl::sleep(Duration::from_secs(1)).await;
            }
        });

        let futures: Vec<Pin<&mut dyn Future<Output = ()>>> =
            vec![tx1_fut, rx_fut, tx_fut];

        trpl::join_all(futures).await;
    });
}
示例 17-24:将 future pin 住,以便把它们移动进向量中

这段代码现在已经可以编译和运行了,而且我们还能在运行时动态地从向量里增加或删除 future,再把它们全部 join 在一起。

Stream trait

现在你已经对 FuturePinUnpin 有了更深入的理解,我们可以把注意力转向 Stream trait 了。正如你在本章前面学到的,stream 很像异步迭代器。不过和 Iterator 以及 Future 不同的是,到本书写作时,标准库里还没有 Stream 的定义;但 futures crate 提供了一个在整个生态系统中被广泛采用的通用定义。

在看 Stream 如何把 IteratorFuture 的特征结合起来之前,我们先回顾一下这两个 trait 的定义。从 Iterator 我们得到了“序列”这个概念:它的 next 方法返回 Option<Self::Item>。从 Future 我们得到了“值会随着时间变得就绪”这个概念:它的 poll 方法返回 Poll<Self::Output>。为了表示“一串会随着时间逐渐就绪的项”,我们就可以定义这样一个 Stream trait,把两者的特征合并起来:

#![allow(unused)]
fn main() {
use std::pin::Pin;
use std::task::{Context, Poll};

trait Stream {
    type Item;

    fn poll_next(
        self: Pin<&mut Self>,
        cx: &mut Context<'_>
    ) -> Poll<Option<Self::Item>>;
}
}

Stream trait 定义了一个名为 Item 的关联类型,用来表示 stream 产生的条目类型。这和 Iterator 很像,因为它可以有零个到多个条目;而和 Future 不同,后者始终只有一个 Output,哪怕这个输出只是 unit 类型 ()

Stream 还定义了一个获取这些条目的方法。它叫 poll_next,这个名字清楚地表明:它既像 Future::poll 那样进行轮询,又像 Iterator::next 那样生成一个接一个的条目。它的返回类型把 PollOption 组合了起来。最外层是 Poll,因为和 future 一样,它需要先检查是否就绪;里面那层是 Option,因为和迭代器一样,它还得表示“后面是否还有更多条目”。

和这个定义非常相似的版本,将来很可能会进入 Rust 标准库。在此之前,它已经是大多数运行时工具箱的一部分,因此你完全可以依赖它,而我们接下来讲的内容通常也都会成立。

不过,在我们前面[“Stream:按顺序出现的 Future”][streams]一节中见到的那些例子里,我们并没有直接用 poll_nextStream,而是用了 nextStreamExt。当然,我们可以像直接操作 future 的 poll 方法那样,手写自己的 Stream 状态机,直接基于 poll_next 来工作。不过,用 await 显然舒服得多,而 StreamExt trait 则为此提供了 next 方法:

#![allow(unused)]
fn main() {
use std::pin::Pin;
use std::task::{Context, Poll};

trait Stream {
    type Item;
    fn poll_next(
        self: Pin<&mut Self>,
        cx: &mut Context<'_>,
    ) -> Poll<Option<Self::Item>>;
}

trait StreamExt: Stream {
    async fn next(&mut self) -> Option<Self::Item>
    where
        Self: Unpin;

    // other methods...
}
}

注意:我们在本章前面实际使用到的定义,看起来会和这个稍微有点不同,因为它需要兼容那些还不支持“在 trait 中使用 async 函数”的 Rust 版本。所以它实际上更像这样:

fn next(&mut self) -> Next<'_, Self> where Self: Unpin;

这里的 Next 类型是一个实现了 Futurestruct,它通过 Next<'_, Self> 的形式,把对 self 的引用生命周期显式命名出来,这样 await 才能和这个方法一起工作。

StreamExt trait 还是所有那些“用于 stream 的有趣方法”的所在地。任何实现了 Stream 的类型,都会自动获得 StreamExt 的实现;不过这两个 trait 之所以分开定义,是为了让社区能够在不影响底层基础 trait 的前提下,不断迭代那些更方便的高层 API。

trpl crate 使用的这个 StreamExt 版本里,这个 trait 不仅定义了 next 方法,还给 next 提供了一个默认实现,这个实现会正确处理 Stream::poll_next 的各种细节。这意味着,即便你将来需要自己写一种流式数据类型,也只需要实现 Stream;然后,任何使用你这个数据类型的人,都会自动获得 StreamExt 及其方法。

关于这些 trait 的底层细节,我们就讲到这里。最后,让我们来想一想:future(包括 stream)、任务和线程到底是如何一起协作的。

future、任务和线程

结合起来看:Future、任务与线程

正如我们在第十六章中看到的那样,线程提供了一种实现并发的方式。而在本章里,我们又看到了另一种方式:使用基于 future 和 stream 的 async。如果你想知道应该在什么时候选哪一种,答案是:这取决于具体情况!而且在很多情况下,真正的选择并不是线程 async,而是线程 async 一起用。

许多操作系统几十年来一直都提供基于线程的并发模型,因此许多编程语言也都支持它们。但这些模型并非没有代价。在许多操作系统中,每个线程都要占用相当可观的内存。线程也只有在你的操作系统和硬件支持它们时才可用。与主流桌面和移动设备不同,一些嵌入式系统甚至根本没有操作系统,因此也就根本没有线程。

async 模型提供了一组不同的,而且归根结底是互补的权衡。在 async 模型中,并发操作不需要各自独占一个线程。相反,它们可以运行在任务上,就像我们在 stream 那一节里使用 trpl::spawn_task 从同步函数中启动工作时那样。任务和线程有点像,但它不是由操作系统管理,而是由库层面的代码,也就是运行时来管理。

线程 API 和任务 API 长得这么像并不是偶然。线程是“一组同步操作”的边界;并发发生在线程之间。任务则是“一组异步操作”的边界;并发既可能发生在任务之间,也可能发生在任务内部,因为任务的函数体里可以在多个 future 之间切换。最后,future 是 Rust 中最细粒度的并发单位,而每个 future 又可能代表一棵由其他 future 组成的树。运行时,准确地说是它的 executor,管理任务;任务再去管理 future。从这个角度说,任务有点像“轻量级的、由运行时管理的线程”,只不过由于它们是由运行时而不是操作系统管理,所以又拥有了一些额外能力。

这并不意味着 async 任务就一定总比线程更好,反过来也一样。在线程基础上的并发,在某些方面其实是比 async 并发更简单的编程模型。这一点既可能是优点,也可能是缺点。线程有点像 “射后不理”(“fire and forget”):它们没有原生对应 future 的机制,因此通常只是一路运行到结束,除非被操作系统本身打断。

而线程和任务常常又能很好地协同工作,因为任务(至少在某些运行时里)可以在线程之间来回移动。事实上,我们这一章一直使用的运行时,在底层默认就是多线程的,包括 spawn_blockingspawn_task 在内。许多运行时会采用一种叫作 work stealing(工作窃取)的方法,根据各个线程当前的利用情况,把任务透明地在线程之间迁移,以提高系统整体性能。而这种方法本身就同时需要线程、任务,也就同时需要 future。

在思考什么时候该用哪种方式时,可以先记住这些经验法则:

  • 如果工作是非常适合并行化的,也就是典型的 CPU 密集型任务,比如有一大批数据而且每一部分都能单独处理,那么线程通常是更好的选择。
  • 如果工作是高度并发的,也就是典型的 I/O 密集型任务,比如要同时处理来自很多不同来源、且到达间隔和频率都各不相同的消息,那么 async 通常是更好的选择。

如果你同时既需要并行,又需要并发,那也完全不必在“线程”和“async”之间二选一。你可以自由地把它们组合起来,让每一种都去做自己最擅长的那一部分。比如,示例 17-25 展示的就是一种在现实 Rust 代码里相当常见的组合方式。

文件名:src/main.rs

extern crate trpl; // for mdbook test

use std::{thread, time::Duration};

fn main() {
    let (tx, mut rx) = trpl::channel();

    thread::spawn(move || {
        for i in 1..11 {
            tx.send(i).unwrap();
            thread::sleep(Duration::from_secs(1));
        }
    });

    trpl::block_on(async {
        while let Some(message) = rx.recv().await {
            println!("{message}");
        }
    });
}
示例 17-25:在线程中用阻塞代码发送消息,并在 async 代码块中等待这些消息

我们先创建一个异步信道,然后启动一个线程,用 move 关键字把信道发送端的所有权移入线程中。在那个线程里,我们发送 1 到 10 这些数字,并在每次发送之间睡眠 1 秒。最后,我们像本章前面一直做的那样,把一个由 async 代码块构造出来的 future 交给 trpl::block_on 去执行。在这个 future 里,我们等待那些消息,就像前面那些消息传递示例里做的一样。

回到本章一开始举过的场景:想象一下,你用一个专门的线程来执行一组视频编码任务,因为视频编码是计算密集型工作;但当这些操作完成时,再通过一个异步信道通知 UI。现实世界里这类组合的例子多得数不过来。

总结

这不会是你在本书里最后一次见到并发。第二十一章中的项目,会把这些概念放到一个比这里那些小例子更真实的场景里来运用,并且更直接地比较“线程”和“任务 / future”这两种解决问题的方式。

无论你最终选择哪一种方法,Rust 都为你提供了编写安全、快速并发代码所需的工具,不管你的目标是高吞吐的 Web 服务器,还是嵌入式操作系统。

接下来,我们将讨论:随着 Rust 程序不断变大,应该如何用符合 Rust 惯例的方式来建模问题并组织解决方案。此外,我们也会谈谈 Rust 的这些惯用法,和你可能已经熟悉的面向对象编程风格之间是什么关系。

评论 (0)