mirror of
https://github.com/gnu4cn/rust-lang-zh_CN.git
synced 2026-08-19 12:43:28 +08:00
Updated 'src/async/concurrency_n_async.md'.
This commit is contained in:
7
projects/message-passing_n_tasks/Cargo.toml
Normal file
7
projects/message-passing_n_tasks/Cargo.toml
Normal file
@@ -0,0 +1,7 @@
|
||||
[package]
|
||||
name = "message-passing_n_tasks"
|
||||
version = "0.1.0"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
trpl = "0.3.0"
|
||||
27
projects/message-passing_n_tasks/src/main.rs
Normal file
27
projects/message-passing_n_tasks/src/main.rs
Normal file
@@ -0,0 +1,27 @@
|
||||
use std::time::Duration;
|
||||
|
||||
fn main() {
|
||||
let fut = 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!("收到 '{value}'");
|
||||
}
|
||||
};
|
||||
|
||||
trpl::block_on(fut);
|
||||
|
||||
println! ("非并发部分");
|
||||
}
|
||||
@@ -147,29 +147,24 @@ hi 来自第一个任务的数字 8 !
|
||||
hi 来自第一个任务的数字 9 !
|
||||
```
|
||||
|
||||
现在,咱们将看到每次都完全相同的顺序,这与我们在线程下看到的情况截然不同。这是因为 `trpl::join` 函数是 *公平的*,意味着他会以相同频率检查各个未来值,在二者之间交替进行,并绝不会在一个已就绪时,让另一个超前。在线程下,是由操作系统决定要检查哪个线程,以及让他运行多久。而在异步的 Rust 下,是由运行时决定要检查哪个任务。(在实践中,细节会变得复杂,因为异步运行时可能会在表象之下,使用操作系统的线程,作为其管理并发性的一部分,所以对运行时来说,保证公平性可能会更费事 -- 但这仍然是可行的!)运行时不必保证对任何给定操作的公平性,他们通常提供不同的 API,让咱们选择是否需要公平性。
|
||||
现在,咱们每次都将看到完全相同的顺序,这与我们在线程下以及清单 17-7 中 `trpl::spawn_task` 下看到的情况大不不同。这是因为 `trpl::join` 函数是 *公平的*,这意味着他会以相等频率检查各个未来值,在二者之间交替进行,并且绝不会让其中一个抢先,即使另一个已经就绪。在线程下,操作系统决定要检查哪个线程以及让他运行多长时间。在异步 Rust 下,运行时决定检查哪个任务。(实际上,细节会变得复杂,因为异步运行时可能在底层使用操作系统的线程,作为其管理并发的一部分,因此保证公平性对于运行时来说可能需要更多工作 -- 但这仍然是可行的!)运行时不必保证任何给定操作的公平性,他们通常提供不同的 API,来让咱们选择是否需要公平性。
|
||||
|
||||
请在等待未来值上,尝试以下这些变化,看看他们能做些什么:
|
||||
请对等待未来值尝试以下这些变化,并看看他们有什么作用:
|
||||
|
||||
- 移除任一循环,或同时两个循环中的异步代码块;
|
||||
- 在定义出每个异步代码块后,立即等待他们;
|
||||
- 仅将第一个循环封装在异步代码块中,并在第二个循环的主体之后等待生成的未来值。
|
||||
|
||||
作为额外挑战,看看咱们是否能在运行代码 *前*,得出每种情况下的输出结果!
|
||||
|
||||
|
||||
- 移除任一循环,或同时两个循环的异步代码块;
|
||||
- 在定义出各个异步代码块后,立即等待他们;
|
||||
- 只将第一个循环封装在异步代码块中,并在第二个循环的主体后,等待所得到的未来值。
|
||||
|
||||
|
||||
作为额外挑战,请在运行代码 *前*,看看咱们能否得出每种情况下的输出结果!
|
||||
|
||||
|
||||
|
||||
## 使用消息传递在两个任务上计数
|
||||
|
||||
|
||||
在未来值间共用数据也将很常见:再次我们将使用消息传递,但这次将使用异步版本的类型与函数。我们将采用与咱们在 [使用消息传递在线程间传输数据](../concurrency/message_passing.md) 小节中略微不同的路径,说明基于线程与基于未来值的并发间的一些关键区别。在清单 17-9 中,我们将从单个的异步代码块开始 -- 而 *不是* 像咱们曾生成单个线程时,生成单个任务。
|
||||
## 使用消息传递在两个任务之间发送
|
||||
|
||||
在未来值之间共用数据也将很常见:我们将再次使用消息传递,但这次是在异步版本的类型和函数下。我们将采取与第 16 章中 [通过消息传递在线程间传输数据](../concurrency/message_passing.md) 小节中的略有不同的路径,以演示基于线程和基于未来值的并发之间的一些关键区别。在下面清单 17-9 中,我们将仅从单个异步代码块开始 -- 而 *不是* 像我们生成单个线程那样生成单个任务。
|
||||
|
||||
<a name="listing_17-9"></a>
|
||||
文件名:`src/main.rs`
|
||||
|
||||
|
||||
```rust
|
||||
let (tx, mut rx) = trpl::channel();
|
||||
|
||||
@@ -177,16 +172,15 @@ hi 来自第一个任务的数字 9 !
|
||||
tx.send(val).unwrap();
|
||||
|
||||
let received = rx.recv().await.unwrap();
|
||||
println!("Got: {received}");
|
||||
println!("收到 {received}");
|
||||
```
|
||||
|
||||
*请单 17-9:创建出一个异步通道并将其中两半分别赋值给 `tx` 与 `rx`*
|
||||
**请单 17-9**:创建异步信道,并指派两端给 `tx` 与 `rx`
|
||||
|
||||
这里,我们使用了我们在第 16 章中,与线程一起使用的多生产者、单消费者通道 API 的一个异步版本 `trpl::channel`。该 API 的异步版本,与基于线程的版本只有一点不同:他使用了可变的接收器 `rx`,而不是不可变接收器 `rx`,而且他的 `recv` 方法会产生一个我们需要等待的未来值,而不是直接产生值。现在,我们可以从发送方往接收方发送消息了。请注意,我们不必生成单独线程,甚至不需要生成任务;我们只需等待这个 `rx.recv` 调用。
|
||||
在这里,我们使用 `trpl::channel`,这是我们在第 16 章中与线程一起使用的多生产者、单消费者信道 API 的异步版本。这一 API 的异步版本与基于线程的版本只有细微差别:他使用可变接收器 `rx` 而不是不可变的,并且他的 `recv` 方法产生一个我们需要等待的未来值,而不是直接生成值。现在,我们可以从发送方发送消息到接受方。请注意,我们不必生成单独的线程甚至任务;我们只需等待 `rx.recv` 调用。
|
||||
|
||||
|
||||
`std::mpsc::channel` 中的同步 `Receiver::recv` 方法,会在收到消息前一直阻塞。而 `trpl::Receiver::recv` 这个方法不会,因为他是异步的。他不会阻塞,而是在消息被接收或通道的发送侧关闭前,会将控制权交还给运行时。相比之下,我们不会等待 `send` 调用,因为他不会阻塞。他之所以无需等待,是因为我们将消息发入的通道,是不受限的 <sup>1</sup>。
|
||||
|
||||
`std::mpsc::channel` 中的同步 `Receiver::recv` 方法会在接收到消息前一直阻塞。而 `trpl::Receiver::recv` 方法则不会,因为他属于异步的。他不会阻塞,而是交还控制权给运行时,直到收到消息或信道的发送侧关闭。相比之下,我们不会等待 `send` 调用,因为他不会阻塞。他之所以无需阻塞,因为我们发入消息的信道是不受限的 <sup>1</sup>。
|
||||
|
||||
> **译注**:
|
||||
>
|
||||
@@ -195,15 +189,15 @@ hi 来自第一个任务的数字 9 !
|
||||
> 参考:[Unbounded nondeterminism](https://en.wikipedia.org/wiki/Unbounded_nondeterminism)
|
||||
|
||||
|
||||
> **注意**:由于所有这些异步代码,都在一个 `trpl::run` 调用的异步代码块中运行,因此其中的所有代码,都可避免阻塞。但是,当 `run` 函数返回时,异步代码块 *外部* 的代码会阻塞。这正是 `trpl::run` 函数的意义所在:他可以让咱们 *选择*,在哪些地方阻塞某些异步代码,从而在哪些地方切换同步代码和异步代码。在大多数异步运行时中,`run` 其实被命名为 `block_on`,正是出于这个原因。
|
||||
|
||||
> **注意**:由于所有这些异步代码都在 `trpl::block_on` 调用的异步代码块中运行,因此其中的所有代码都可避免阻塞。但是,异步代码块 *外部* 的代码将在 `block_on` 函数运行时阻塞。这正是 `trpl::block_on` 函数的核心意义:他让咱们可以 *选择* 于何处阻塞某段异步代码,从而选择了于何处在同步代码和异步代码之间切换。
|
||||
|
||||
|
||||
请注意这个例子的两点。首先,消息将立即到达。其次,虽然我们在这里使用了一个未来值,但这里还没有并发。该列表中的所有事情,都是按顺序发生的,就像没有涉及到未来值一样。
|
||||
|
||||
|
||||
我们来通过发送一系列消息并在中间休眠,解决第一部分问题,如清单 17-10 所示。
|
||||
请注意这个示例中的两点。首先,消息将立即到达。其次,尽管我们在这里使用了未来值,但还并没有并发。清单中的所有操作都是按顺序发生的,就像不涉及未来值一样。
|
||||
|
||||
我们来通过发送一系列消息并在每次发送之间休眠,解决第一部分,如下清单 17-10 中所示。
|
||||
|
||||
<a name="listing_17-10"></a>
|
||||
文件名:`src/main.rs`
|
||||
|
||||
|
||||
@@ -223,26 +217,24 @@ hi 来自第一个任务的数字 9 !
|
||||
}
|
||||
|
||||
while let Some(value) = rx.recv().await {
|
||||
println!("received '{value}'");
|
||||
println!("收到 '{value}'");
|
||||
}
|
||||
```
|
||||
|
||||
*请单 17-10:通过异步通道发送及接收多条消息,并在每条消息之间对 `await` 休眠*
|
||||
**请单 17-10**:通过异步信道发送并接收多条消息,并在每条消息之间通过 `await` 休眠
|
||||
|
||||
除发送消息外,我们还需要接收他们。在本例中,由于我们知道有多少条消息进来,因此我们可通过调用 `rx.recv().await` 四次,手动完成接收。但在真实世界中,我们一般会等待 *未知* 数量的消息,因此我们需要一直等待,直到确定没有更多信息为止。
|
||||
除了发送消息外,我们还需要接收他们。在这一情形下,由于我们知道有多少条消息传入,因此可以通过手动调用 `rx.recv().await` 四次完成接收。但在现实世界中,我们通常将等待 *未知* 数量的消息,因此需要持续等待,直到确定不再有信息为止。
|
||||
|
||||
在 [清单 16-10](../concurrency/message_passing.md#listing_16-10) 中,我们使用 `for` 循环来处理从同步通道接收到的所有项目。然而,Rust 还没有对 *异步生成的* 项目序列使用 `for` 循环的方法,因此我们需要使用一种以前从未见过的循环:`while let` 条件循环。这是我们在第 6 章中 [`if let` 与 `let else` 下的简明控制流](../enums_and_pattern_matching/if-let_control_flow.md) 小节中,见过的 `if let` 结构的循环版本。只要其指定的模式继续与值匹配,该循环就会继续执行。
|
||||
|
||||
其中的 `rx.recv` 调用会生成一个未来值,我们等待该未来值。运行时将暂停该未来值,直到他准备就绪。一旦消息到达,该未来值就将解析为 `Some(message)`,且解析次数与消息到达次数相同。当信道关闭时,无论是否 *有* 消息到达,该未来值都将解析为 `None`,以表示不再有值,因此我们应该停止轮询 -- 即停止等待(未来值)。
|
||||
|
||||
`while let` 循环将所有这一切整合在一起。当调用 `rx.recv().await` 的结果是 `Some(message)` 时,我们得到对消息的访问,并可以在循环体中使用他,就像在 `if let` 下一样。当结果为 `None` 时,该循环结束。循环每次完成时,他都会再次遇到等待点,因此运行时会再次暂停他,直到另一条消息到达。
|
||||
|
||||
该代码现在成功地发送并接收所有消息。遗憾的是,仍然存在一些问题。其中之一便是,消息没有以半秒为间隔到达。他们会在我们启动程序 2 秒(2000 毫秒)后一次性全部到达。另外,这个程序还永远不会退出!相反,他会无限期等待新消息。咱们将需要用时 `Ctrl-c` 关闭他。
|
||||
|
||||
|
||||
在 [清单 16-10](../concurrency/message_passing.md#listing_16-10) 中,我们使用了个 `for` 循环,处理从同步通道接收到的所有项目。然而,Rust 还没有一种,为 *异步的* 序列项目编写 `for` 循环的方法,因此我们需要使用一种以前从未见过的循环:`while let` 条件循环。这是我们曾在 [“使用 `if let` 与 `let else` 的简明控制流”](../enums_and_pattern_matching/if-let_control_flow.md) 小节中,看到的 `if let` 结构的循环版本。只要循环所指定的模式继续匹配值,其就会继续执行。
|
||||
|
||||
|
||||
其中的 `rx.recv` 调用,会产生一个我们所等待的未来值。在其准备好前,运行时将暂停该未来值。一旦有消息到达,这个未来值将解析为 `Some(message)`,解析次数与消息到达次数相同。在通道关闭时,无论 *有多少* 消息到达,该未来值都会解析为表示不再有值的 `None`,并因此我们就应停止轮询 -- 即停止等待。
|
||||
|
||||
|
||||
那个 `while let` 循环,将所有这一切联系在了一起。在调用 `rx.recv().await` 的结果为 `Some(message)` 时,我们就可以访问到消息,并在循环体中使用他,就跟使用 `if...let` 一样。在结果为 `None` 时,则该循环结束。每次循环完毕时,其都会再次遇到等待点,因此运行时会再度将其暂停,直到另一条消息到达。
|
||||
|
||||
|
||||
代码现在可以成功发送并接收所有信息。遗憾的是,其间仍然存在一些问题。首先,消息不是以半秒为间隔到达的。他们会在我们启动程序 2 秒(2000 毫秒)后,一次性到达。另外,这个程序永远不会退出!相反,他会一直等待新信息。咱们需要以 `Ctrl-c` 关闭他。
|
||||
|
||||
### 同一个异步代码块内的代码会线性地执行
|
||||
|
||||
我们先来看看,为什么消息会在全部延迟后一次性发送,而不是每条消息之间都有延迟。在某个给定异步代码块中,`await` 关键字在代码中出现的顺序,就是他们在程序运行时,被执行的顺序。
|
||||
|
||||
|
||||
Reference in New Issue
Block a user