10 流
6 分钟阅读
流(stream)是一系列异步产生的值。它是 Rust std::iter::Iterator 的异步等价物,由 Stream trait 表示。流可以在 async 函数中迭代,也可以使用适配器进行转换。Tokio 在 StreamExt trait 上提供了许多常用适配器。
Tokio 在单独的 crate tokio-stream 中提供流支持。
| |
info 目前,Tokio 的 Stream 工具位于
tokio-streamcrate 中。 一旦Streamtrait 在 Rust 标准库中稳定,tokiocrate 将承载 Tokio 的流工具。
目前,Rust 编程语言不支持异步 for 循环。
相反,迭代流使用 while let 循环配合 StreamExt::next()。
| |
与迭代器类似,next() 方法返回 Option<T>,其中 T 是流的值类型。收到 None 表示流迭代已结束。
Mini-Redis 广播
让我们通过一个使用 Mini-Redis 客户端的稍复杂示例来说明。
完整代码见此处。
| |
一个任务被生成,向 Mini-Redis 服务器的 “numbers” 通道发布消息。然后,在主任务上,我们订阅 “numbers” 通道并显示收到的消息。
订阅后,对返回的 subscriber 调用 into_stream()。这会消费 Subscriber,返回一个在消息到达时产生消息的流。在开始迭代消息之前,注意该流使用 tokio::pin! 被固定到栈上。对流调用 next() 要求流被固定。into_stream() 函数返回的流未被固定,我们必须显式固定它才能迭代。
info 当 Rust 值在内存中不能再被移动时,它就被「固定」了。固定值的一个关键性质是,可以取得指向固定数据的指针,且调用方可以确信该指针保持有效。
async/await利用这一特性来支持在.await点之间借用数据。
如果我们忘记固定流,会得到类似这样的错误:
| |
如果你遇到类似这样的错误消息,试试固定该值!
在运行之前,先启动 Mini-Redis 服务器:
| |
然后尝试运行代码。我们会看到消息输出到 STDOUT。
| |
由于订阅和发布之间存在竞争,一些早期消息可能会丢失。程序永远不会退出。对 Mini-Redis 通道的订阅会一直保持活跃,直到服务器仍在运行。
让我们看看如何使用流来扩展这个程序。
适配器
接受 Stream 并返回另一个 Stream 的函数通常称为「流适配器」,它们是「适配器模式」的一种形式。常见的流适配器包括 map、take 和 filter。
让我们更新 Mini-Redis 示例以便程序能够退出。收到三条消息后,停止迭代消息。这通过 take 实现。该适配器将流限制为最多产生 n 条消息。
| |
再次运行程序,我们会得到:
| |
这次程序会结束。
现在,让我们将流限制为个位数。我们通过检查消息长度来判断。我们使用 filter 适配器丢弃任何不匹配谓词的消息。
| |
再次运行程序,我们会得到:
| |
注意适配器应用的顺序很重要。先 filter 再 take 与先 take 再 filter 是不同的。
最后,我们通过剥离输出中的 Ok(Message { ... }) 部分来整理输出。这通过 map 完成。由于这是在 filter 之后应用的,我们知道消息是 Ok,因此可以使用 unwrap()。
| |
现在,输出为:
| |
另一种选择是使用 filter_map 将 filter 和 map 步骤合并为一次调用。
还有更多可用适配器。完整列表见此处。
实现 Stream
Stream trait 与 Future trait 非常相似。
| |
Stream::poll_next() 函数很像 Future::poll,只不过它可以被反复调用以从流中接收多个值。正如我们在 深入 async 中所见,当流尚未准备好返回值时,会返回 Poll::Pending。任务的 waker 会被注册。一旦流应该再次被轮询,waker 会收到通知。
size_hint() 方法的用法与迭代器中的相同。
通常,手动实现 Stream 是通过组合 future 和其他流来完成的。作为示例,让我们在 深入 async 中实现的 Delay future 基础上构建。我们将把它转换为一个流,以 10 ms 的间隔产生三次 ()。
| |
async-stream
使用 Stream trait 手动实现流可能很繁琐。
遗憾的是,Rust 编程语言尚不支持用于定义流的 async/await 语法。相关工作正在进行,但尚未就绪。
async-stream crate 可作为临时解决方案。该 crate 提供 stream! 宏,将输入转换为流。使用这个 crate,上面的 interval 可以这样实现:
| |