一句话总结
Stream 是迭代器(Iterator)的异步版本:一次产出一个值,但每个值可能需要 .await 等待.
它之于"一连串异步到达的值",正如 Iterator 之于"一连串同步的值".配合 StreamExt 的 map/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 StreamExt
和第 11 篇的 AsyncReadExt 同理:next,map,filter 这些方法定义在 StreamExt 扩展 trait 上,必须 use tokio_stream::StreamExt 才能调用,否则会提示方法不存在.
链式适配器
Stream 复用了熟悉的迭代器适配器,只是它们处理的是异步流:
|
惰性也一脉相承
和迭代器一样,Stream 的适配器是惰性的:map/filter 只是描述变换,直到用 next().await 或 while let 去消费,数据才真正流动.
第 04 篇关于"惰性 + 消费器"的理解可以直接平移过来.
把 channel 当 Stream 消费
实践中 Stream 常来自第 09 篇的 channel.例如把 mpsc::Receiver 包装成 Stream,就能用适配器处理源源不断的消息:
|
这让"持续到达的事件"和"集合数据"享有统一的处理范式--无论数据来自网络,定时器还是 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"叠加迭代器式函数流水线.
动手练习
- 基本遍历:用
tokio_stream::iter造一个 Stream,用while let+next().await遍历打印. - 链式变换:给上题加上
filter+map链,只保留奇数并乘以 10. - 复现报错:故意不
use StreamExt,观察next方法找不到的报错,再修复. - ReceiverStream:用
ReceiverStream把一个mpsc::Receiver变成 Stream,并用适配器处理收到的消息. - 本质区别:用一句话说明 Stream 与 Iterator 的唯一本质区别,以及它和 Go
for range ch的关系.