阶段五 · 进阶与陷阱

Stream:异步数据流

一句话总结

Stream 是迭代器(Iterator)的异步版本:一次产出一个值,但每个值可能需要 .await 等待.
它之于"一连串异步到达的值",正如 Iterator 之于"一连串同步的值".配合 StreamExtmap/filter/next 等方法,能用熟悉的链式风格处理异步数据流.

从 Iterator 到 Stream

回顾基础篇的迭代器:Iterator 的核心是 next(),每次同步地返回 Option<T>,None 表示结束.Stream 几乎一样,只是 next() 变成了异步的--取下一个值可能要等(等网络包,等定时器):

Iterator(同步) Stream(异步)
取下一个值 next() -> Option<T> next().await -> Option<T>
结束信号 返回 None 返回 None
遍历 for x in iter while let Some(x) = s.next().await
适配器 map/filter/... 同名,来自 StreamExt

Iterator 像一摞已经印好的卡片,翻一张拿一张,立刻就有.
Stream 像陆续到货的快递:排队取件,下一件可能要等一会儿(.await),但取件,处理的流程和翻卡片别无二致.

基本用法

Stream 不在标准库的 prelude 里,常用 tokio-stream crate.遍历一个 Stream:

use tokio_stream::StreamExt;

#[tokio::main]
async fn main() {
let mut stream = tokio_stream::iter(vec![1, 2, 3]);

while let Some(v) = stream.next().await {
println!("收到 {v}");
}
}

记得 use StreamExt

和第 11 篇的 AsyncReadExt 同理:next,map,filter 这些方法定义在 StreamExt 扩展 trait 上,必须 use tokio_stream::StreamExt 才能调用,否则会提示方法不存在.

链式适配器

Stream 复用了熟悉的迭代器适配器,只是它们处理的是异步流:

use tokio_stream::StreamExt;

let mut stream = tokio_stream::iter(1..=10)
.filter(|n| n % 2 == 0) // 留偶数
.map(|n| n * n); // 平方

while let Some(v) = stream.next().await {
println!("{v}"); // 4, 16, 36, 64, 100
}

惰性也一脉相承

和迭代器一样,Stream 的适配器是惰性的:map/filter 只是描述变换,直到用 next().awaitwhile let 去消费,数据才真正流动.

第 04 篇关于"惰性 + 消费器"的理解可以直接平移过来.

把 channel 当 Stream 消费

实践中 Stream 常来自第 09 篇的 channel.例如把 mpsc::Receiver 包装成 Stream,就能用适配器处理源源不断的消息:

use tokio_stream::wrappers::ReceiverStream;
use tokio_stream::StreamExt;

let (tx, rx) = tokio::sync::mpsc::channel(16);
let mut stream = ReceiverStream::new(rx);

tokio::spawn(async move {
for i in 0..5 { tx.send(i).await.unwrap(); }
});

while let Some(v) = stream.next().await {
println!("流式处理:{v}");
}

这让"持续到达的事件"和"集合数据"享有统一的处理范式--无论数据来自网络,定时器还是 channel,都能用同一套链式 API.

与 Go 的对照

Go 里"持续到达的值"通常直接用 for v := range ch 遍历 channel.
Rust 的 Stream 在这个用法之上,叠加了迭代器那一整套 map/filter/take 链式变换能力--相当于"可遍历的 channel" + "函数式流水线"的合体.

系列收尾

至此,已经走完异步 Rust 的完整路径:从"为什么需要异步"的动机,到 Future / async / 运行时的机制,再到 Tokio 的任务,原语,编排,IO 与实战,最后是取消,避坑与数据流.

现在具备读懂主流异步 Rust 代码,并独立构建异步服务的能力.后续可深入的方向:具体框架(如 axum web 框架),Pin 与手写 Future,以及高性能调优.

快速回顾

  • Stream = 异步 Iterator:核心是 next().await,每个值可能需要等待,None 表示结束.
  • 遍历:用 while let Some(x) = s.next().await,对应同步的 for x in iter.
  • 适配器:map/filter 等来自 StreamExt,需 use;同样是惰性的,靠消费驱动.
  • 常见来源:用 ReceiverStream 等把 channel 包成 Stream,统一处理持续到达的事件.
  • 对照 Go:相当于"可遍历的 channel"叠加迭代器式函数流水线.

动手练习

  1. 基本遍历:用 tokio_stream::iter 造一个 Stream,用 while let + next().await 遍历打印.
  2. 链式变换:给上题加上 filter + map 链,只保留奇数并乘以 10.
  3. 复现报错:故意不 use StreamExt,观察 next 方法找不到的报错,再修复.
  4. ReceiverStream:用 ReceiverStream 把一个 mpsc::Receiver 变成 Stream,并用适配器处理收到的消息.
  5. 本质区别:用一句话说明 Stream 与 Iterator 的唯一本质区别,以及它和 Go for range ch 的关系.