Tokio 实战
用 Tokio 写真实并发任务:spawn、超时、并发组合与取消。
Tokio 实战
tokio::fs 与 tokio::io
tokio::fs 的每个函数其实都内部派发到 spawn_blocking——因为绝大多数操作系统 没有真正异步的文件 IO(Linux 的 io_uring 除外,Tokio 默认不用)。 这一点决定了「文件 IO 用 async 收益有限」。
rust
// 依赖:tokio = { version = "1", features = ["full"] }
use tokio::io::{AsyncReadExt, AsyncWriteExt};
#[tokio::main]
async fn main() -> std::io::Result<()> {
// tokio::fs 是 async 的 std::fs
tokio::fs::write("demo.txt", b"hello async").await?;
// 一次性读全部内容
let content = tokio::fs::read_to_string("demo.txt").await?;
println!("{content}");
// 需要读 trait 方法(read / write / read_exact ...)时必须 use AsyncReadExt
let mut file = tokio::fs::File::open("demo.txt").await?;
let mut buf = [0u8; 5];
let n = file.read(&mut buf).await?;
println!("read {n} bytes: {:?}", &buf[..n]);
// 流式追加写
let mut out = tokio::fs::File::create("demo2.txt").await?;
out.write_all(b"line 1\n").await?;
out.write_all(b"line 2\n").await?;
out.flush().await?;
tokio::fs::remove_file("demo.txt").await?;
tokio::fs::remove_file("demo2.txt").await?;
Ok(())
}💡 对照:
read/write这些方法来自AsyncReadExt/AsyncWriteExt(扩展 trait), 忘记use时的报错形如error[E0599]: no method named 'read' found for struct 'tokio::fs::File'—— 先说「没有这个方法」,再提示「traitAsyncReadExtis implemented but not in scope」。
TCP 回声服务器(完整可运行)
toml
# Cargo.toml
[package]
name = "echo-server"
version = "0.1.0"
edition = "2024"
[[bin]]
name = "server"
path = "src/server.rs"
[[bin]]
name = "client"
path = "src/client.rs"
[dependencies]
tokio = { version = "1", features = ["full"] }rust
// src/server.rs —— 完整可运行的 Tokio TCP 回声服务器
// 依赖:tokio = { version = "1", features = ["full"] }
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::TcpListener;
#[tokio::main]
async fn main() -> std::io::Result<()> {
// 监听本机回环地址的 8080 端口(换成 0.0.0.0:8080 才是监听所有网卡)
let listener = TcpListener::bind("127.0.0.1:8080").await?;
println!("listening on {}", listener.local_addr()?);
loop {
// accept 是异步的:没有新连接时这里不会占住线程
let (stream, peer) = listener.accept().await?;
println!("new connection from {peer}");
// 每个连接起一个任务:任务数量不受线程数限制
tokio::spawn(async move {
if let Err(e) = handle(stream).await {
eprintln!("connection {peer} error: {e}");
}
});
}
}
async fn handle(stream: tokio::net::TcpStream) -> std::io::Result<()> {
// 把读写两半拆开:读用带缓冲的行读取,写直接用原 stream
let (read_half, mut write_half) = stream.into_split();
let mut lines = BufReader::new(read_half).lines();
while let Some(line) = lines.next_line().await? {
// 回声:原样写回
write_half.write_all(line.as_bytes()).await?;
write_half.write_all(b"\n").await?;
write_half.flush().await?;
}
println!("client disconnected");
Ok(())
}rust
// src/client.rs —— 一次连接、发两行、读两行
// 依赖:tokio = { version = "1", features = ["full"] }
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::TcpStream;
#[tokio::main]
async fn main() -> std::io::Result<()> {
let mut stream = TcpStream::connect("127.0.0.1:8080").await?;
stream.write_all(b"ping\n").await?;
stream.write_all(b"rust\n").await?;
stream.flush().await?;
let mut lines = BufReader::new(stream).lines();
while let Some(line) = lines.next_line().await? {
println!("server says: {line}");
}
Ok(())
}运行:先
cargo run --bin server,另开一个终端cargo run --bin client。 输出:服务端打印listening on 127.0.0.1:8080与new connection from ..., 客户端打印server says: ping、server says: rust。
一个服务器也可以接受多个客户端同时连接,每个连接一个任务。用 tokio::spawn 时 stream 被 move 进任务,因此不需要锁——这正是 async 比「共享状态 + 锁」更好写的地方。
定时与超时
rust
// 依赖:tokio = { version = "1", features = ["full"] }
use std::time::Duration;
use tokio::time::{interval, timeout, Instant};
async fn slow() -> &'static str {
tokio::time::sleep(Duration::from_millis(300)).await;
"done"
}
#[tokio::main]
async fn main() {
// sleep:让出线程,不占 CPU(对比 std::thread::sleep 会占住线程)
let start = Instant::now();
tokio::time::sleep(Duration::from_millis(100)).await;
println!("slept {:?}", start.elapsed());
// timeout:给任意 future 加时限,返回 Result
match timeout(Duration::from_millis(100), slow()).await {
Ok(v) => println!("in time: {v}"),
Err(_elapsed) => println!("timed out"), // slow 在这里被 drop(取消)
}
// interval:周期性 tick;MissedTickBehavior 决定落后时怎么补
let mut ticker = interval(Duration::from_millis(50));
for _ in 0..3 {
ticker.tick().await;
println!("tick at {:?}", start.elapsed());
}
}输出:
slept 100ms 左右、timed out、三行tick at ...(约 50/100/150ms)。
⚠️ 陷阱:
timeout超时后会 drop 内层 future。如果内层 future 已经做了一半的 副作用(见「取消的本质是 drop future」),这个副作用不会被回滚。
tokio::sync 与 std 对照
tokio::sync 里的类型都有异步版本,因为它们的「等待」也要能挂起任务而不是挂起线程。 同一个功能有多少个「版本」非常容易混,先记这张表:
| 需求 | std 对应物 | Tokio 对应物 | 什么时候用 Tokio 版本 |
|---|---|---|---|
| 互斥访问共享数据 | std::sync::Mutex | tokio::sync::Mutex | 需要在跨 .await 的临界区里持锁 |
| 读写锁 | std::sync::RwLock | tokio::sync::RwLock | 同上,且读多写少 |
| 一次性值/完成通知 | std::sync::OnceLock、Condvar | tokio::sync::OnceCell、Notify | 等待发生在 async 上下文 |
| 多生产者单消费者队列 | std::sync::mpsc | tokio::sync::mpsc | 需要异步 recv().await,或需要背压 |
| 单次应答(请求-响应) | (无直接对应) | tokio::sync::oneshot | 一个值、一次、异步等待 |
| 广播 | (无) | tokio::sync::broadcast | 一个发送者、多个订阅者 |
| 多播/可多消费者 | Arc<Mutex<VecDeque>> | tokio::sync::watch、async_channel | 状态广播、配置热更新 |
| 并发限流 | 信号量需自己写 | tokio::sync::Semaphore | 限制同时进行的请求数 |
| 屏障 | std::sync::Barrier | tokio::sync::Barrier | 等待 N 个任务到齐 |
| 只读共享 | Arc<T> | Arc<T>(相同) | —— |
tokio::sync::Mutex 与 std::sync::Mutex 的真正差别:
std::sync::Mutex | tokio::sync::Mutex | |
|---|---|---|
| 等待时行为 | 阻塞 OS 线程 | 挂起任务,让出 worker 线程 |
守卫是否 Send | 否(MutexGuard: !Send) | 是(MutexGuard: Send) |
| 性能(无竞争) | 更快(futex,几十 ns) | 更慢(要登记 waker) |
跨 .await 持锁 | 编译错误(future 变成 !Send) | 允许,但要警惕死锁 |
⚠️ 陷阱(重要规则):不要跨
.await持有std::sync::Mutex的守卫。MutexGuard不是Send,一旦跨挂起点存活,整个 future 就不再Send,tokio::spawn会直接拒绝编译(错误细节见「跨.await持MutexGuard:future cannot be sent between threads safely」)。即使你在current_thread运行时里绕过了这个限制,还会制造第二个问题: 锁被一个任务长期持有,其他任务在lock()里阻塞 OS 线程,runtime 的 worker 被白白占住——高并发下会退化成「线程池被锁死」。
修法有三种,按推荐顺序:
rust
// 依赖:tokio = { version = "1", features = ["full"] }
use std::sync::{Arc, Mutex};
use tokio::sync::Mutex as AsyncMutex;
// 1) 缩小临界区:让守卫在 await 之前就 drop —— 首选
async fn good_short_critical_section(shared: &Mutex<u32>) {
{
let mut guard = shared.lock().unwrap();
*guard += 1;
} // 守卫在这里 drop
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
}
// 2) 换成 tokio::sync::Mutex:允许跨 await,但临界区越大越容易死锁
async fn with_async_mutex(shared: &AsyncMutex<u32>) {
let mut guard = shared.lock().await;
*guard += 1;
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
drop(guard);
}
// 3) 只在 await 后需要结果:先用作用域取出值,再 await
async fn take_then_await(shared: &Mutex<u32>) -> u32 {
let snapshot = { *shared.lock().unwrap() };
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
snapshot
}
#[tokio::main]
async fn main() {
let a = Arc::new(Mutex::new(0));
let b = Arc::new(AsyncMutex::new(0));
good_short_critical_section(&a).await;
with_async_mutex(&b).await;
println!("{}", take_then_await(&a).await);
}输出:
1。
其他两个高频类型的用法:
rust
// 依赖:tokio = { version = "1", features = ["full"] }
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::{mpsc, oneshot, Notify, Semaphore};
#[tokio::main]
async fn main() {
// mpsc:多生产者单消费者,buffer=16 提供背压;send 满了会 await
let (tx, mut rx) = mpsc::channel::<u32>(16);
for i in 0..3 {
let tx = tx.clone();
tokio::spawn(async move {
tx.send(i).await.expect("receiver dropped");
});
}
drop(tx); // 所有发送者都 drop 后,recv() 才会返回 None
while let Some(v) = rx.recv().await {
println!("got {v}");
}
// oneshot:一次性的请求-响应,非常适合「任务完成后回报结果」
let (reply_tx, reply_rx) = oneshot::channel::<String>();
tokio::spawn(async move {
reply_tx.send(String::from("finished")).expect("caller gone");
});
println!("reply: {}", reply_rx.await.expect("sender dropped"));
// Semaphore:限制同时进行的操作数(拿到 permit 才能继续)
let sem = Arc::new(Semaphore::new(2));
let mut handles = Vec::new();
for id in 0..5 {
let sem = Arc::clone(&sem);
handles.push(tokio::spawn(async move {
let _permit = sem.acquire_owned().await.expect("semaphore closed");
tokio::time::sleep(Duration::from_millis(20)).await;
id
}));
}
for h in handles {
println!("task {} done", h.await.expect("join failed"));
}
// Notify:最轻量的「等一下、我叫你」;notify_one 在无人等待时会存一次许可
let notify = Arc::new(Notify::new());
let waiter = {
let notify = Arc::clone(¬ify);
tokio::spawn(async move {
notify.notified().await;
"woken"
})
};
notify.notify_one();
println!("{}", waiter.await.expect("join failed"));
}输出(顺序因调度略有不同):3 行
got 0/1/2、reply: finished、 5 行task 0..4 done(受Semaphore::new(2)限制,每批最多 2 个)、woken。
🚀 进阶:
tokio::sync::watch是「配置热更新」的标准答案——发送方tx.send(new_cfg),每个接收方rx.changed().await后*rx.borrow()读最新值。 它只保留最新一份值,不排队,因此不会像broadcast那样因消费者慢而丢历史。
spawn 的 'static + Send 要求
rust
// 依赖:tokio = { version = "1", features = ["full"] }
#[tokio::main]
async fn main() {
let local = String::from("borrowed");
// 错:async 块借用 local,future 带生命周期,不满足 'static
// tokio::spawn(async { println!("{local}"); });
// 对:move 把所有权交给任务
tokio::spawn(async move { println!("{local}"); }).await.ok();
// 需要共享时用 Arc;需要修改时 Arc<Mutex<_>> 或 Arc<tokio::sync::Mutex<_>>
let shared = std::sync::Arc::new(std::sync::atomic::AtomicU32::new(0));
let mut handles = Vec::new();
for _ in 0..4 {
let shared = std::sync::Arc::clone(&shared);
handles.push(tokio::spawn(async move {
shared.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}));
}
for h in handles {
h.await.expect("join failed");
}
println!("{}", shared.load(std::sync::atomic::Ordering::Relaxed));
}输出:
borrowed,然后4。
如果确实需要「借用局部变量做并发」,又不想 spawn,就用 join!:它在同一个任务内并发, 不要求 'static(也不要求 Send,在 current_thread 下):
rust
// 依赖:tokio = { version = "1", features = ["full"] }
#[tokio::main]
async fn main() {
let data = vec![1u32, 2, 3];
// 借用 data 的不可变引用,安全且不需要 Arc
let (a, b) = tokio::join!(sum(&data), sum(&data));
println!("{a} {b}");
}
async fn sum(xs: &[u32]) -> u32 {
tokio::task::yield_now().await;
xs.iter().sum()
}🚀 进阶:Rust 1.85 稳定的 async 闭包(
async || { ... })让「把 async 逻辑当参数传」 自然得多,例如stream.map(async |x| fetch(x).await)。它解决了 2024 之前「闭包不能是 async」 只能退回|x| async move { ... }的别扭写法。
Tokio 生态常用 crate
| crate | 用途 | 一行说明 | 安装 |
|---|---|---|---|
axum | Web 框架 | Tokio 官方团队维护,基于 tower + hyper,路由与提取器极简 | cargo add axum |
reqwest | HTTP 客户端 | 异步 + 同步双形态,默认 rustls/tls,json() 直接反序列化 | cargo add reqwest --features json |
sqlx | 数据库 | 纯 Rust 异步驱动,编译期校验 SQL(需 DATABASE_URL) | cargo add sqlx --features runtime-tokio,sqlite |
tower | 中间件/服务抽象 | Service trait:超时、重试、限流、负载均衡都可叠加 | cargo add tower --features timeout,retry |
tracing | 结构化日志 | span + event 模型,异步任务追踪靠 tracing::instrument | cargo add tracing tracing-subscriber |
serde | 序列化 | #[derive(Serialize, Deserialize)],几乎所有 Web/DB crate 的共同依赖 | cargo add serde --features derive |
hyper | 底层 HTTP | axum 的底层实现,需要极致控制时直接用 | cargo add hyper --features server,http1 |
futures | 组合子工具库 | join_all、StreamExt、block_on,与运行时无关 | cargo add futures |
tokio-stream | Stream 适配 | 把 mpsc::Receiver 等变成 Stream,与 futures::Stream 打通 | cargo add tokio-stream |
async-trait | trait 里的 async | 需要 dyn Trait 时的兼容方案(见「async fn 在 trait 中:1.75+ 能写,但 dyn 仍需变通」) | cargo add async-trait |
powershell
# 一次性加上本章演练常用的几个(PowerShell)
cargo add tokio --features full
cargo add futures
cargo add tokio-stream
cargo add serde --features derive
cargo add tracing tracing-subscriber并发组合与取消
join! vs try_join! vs spawn
细看这三者的差别,关键在于 future 在哪个任务里跑:
rust
// 依赖:tokio = { version = "1", features = ["full"] }
use std::time::Duration;
async fn fetch(id: u32) -> Result<String, String> {
tokio::time::sleep(Duration::from_millis(30)).await;
if id == 3 { Err(format!("id {id} failed")) } else { Ok(format!("body {id}")) }
}
#[tokio::main]
async fn main() {
// join!:一个任务内并发;3 个请求同时发起,总耗时 ≈ 30ms
let (a, b, c) = tokio::join!(fetch(1), fetch(2), fetch(4));
println!("{a:?} {b:?} {c:?}");
// try_join!:任一失败立刻返回;已经完成的其它 future 结果被丢弃
let r = tokio::try_join!(fetch(1), fetch(3), fetch(4));
println!("{r:?}");
// spawn:每个请求独立任务,可跨线程并行;代价是 'static + Send
let handles: Vec<_> = (1..=3u32)
.map(|id| tokio::spawn(async move { fetch(id).await }))
.collect();
for h in handles {
// JoinHandle<Result<String, String>> 的 await 得到 Result<Result<..>, JoinError>
match h.await {
Ok(Ok(body)) => println!("ok: {body}"),
Ok(Err(e)) => println!("业务错误: {e}"),
Err(join_err) => println!("任务失败: {join_err}"),
}
}
}输出:
Ok("body 1") Ok("body 2") Ok("body 4")、Err("id 3 failed"), 然后三行请求结果。
选型建议:
- 同一批可枚举的异步操作、需要借用局部数据 →
join!/try_join!。 - 数量不定、需要独立生命周期、可能 panic 隔离 →
spawn+JoinSet(JoinSet会自动 在 drop 时 abort 所有任务,比手动收集JoinHandle更不容易泄漏)。 - 只想快、不要慢的 →
select!或timeout。
select! 的语义与取消安全
rust
// 依赖:tokio = { version = "1", features = ["full"] }
use std::time::Duration;
use tokio::sync::mpsc;
#[tokio::main]
async fn main() {
let (tx, mut rx) = mpsc::channel::<&'static str>(4);
tokio::spawn(async move {
tx.send("from channel").await.expect("receiver dropped");
});
// select! 的语义:并发轮询所有分支,谁先 Ready 就执行谁的分支体;
// 输掉的分支(以及它们已经轮询过的 future)会被 drop
tokio::select! {
Some(msg) = rx.recv() => println!("branch 1: {msg}"),
_ = tokio::time::sleep(Duration::from_millis(500)) => println!("branch 2: timeout"),
}
// biased;:按书写顺序优先,而不是随机;适合「先检查取消信号」
let mut tick = tokio::time::interval(Duration::from_millis(10));
tokio::select! {
biased;
_ = tick.tick() => println!("biased: tick wins"),
_ = tokio::time::sleep(Duration::from_millis(1)) => println!("biased: sleep wins"),
}
}输出:
branch 1: from channel、biased: tick wins。
取消安全(cancel safety) 指的是:一个 future 在被 drop(在任意挂起点之后)时, 是否还能保证「没完成也不留下半成品状态」。
- 安全(cancel safe):
mpsc::Receiver::recv(没收到就没有消费)、TcpStream::connect、tokio::time::sleep、Mutex::lock、read_exact之外的「要么全做、要么不做」的操作。 Tokio 文档对每个方法都标了 cancel safe 与否,用select!前必须查一遍。 - 不安全:
AsyncReadExt::read/read_exact、AsyncBufReadExt::read_line、lines().next_line()、io::copy这类「已经消费了一部分输入」的操作。 被取消时,已读走的那部分字节就丢了,且无法放回流里。
取消的本质是 drop future
Rust 没有「取消一个 future」的专门 API:取消就是把它 drop 掉。所有析构逻辑照常运行 (这一点比 Go 的 goroutine 泄漏、Java 的 interrupt 都更干净),但副作用不会回滚。
rust
// 依赖:tokio = { version = "1", features = ["full"] }
use std::time::Duration;
use tokio::io::AsyncReadExt;
use tokio::time::timeout;
// 危险:先读 header,再读 body。若在第二次 read 时被取消,
// 前 4 个字节已经从流里消费掉了,无法放回。
async fn read_message(stream: &mut tokio::net::TcpStream) -> std::io::Result<Vec<u8>> {
let mut header = [0u8; 4];
stream.read_exact(&mut header).await?; // 若被取消,header 已消耗
let len = u32::from_be_bytes(header) as usize;
let mut body = vec![0u8; len];
stream.read_exact(&mut body).await?; // 若被取消,body 读了一半
Ok(body)
}
// 危险用法演示:超时取消会丢掉已经读到的半条消息,下一次调用从流中间开始解析
async fn broken(stream: &mut tokio::net::TcpStream) -> std::io::Result<Option<Vec<u8>>> {
match timeout(Duration::from_millis(100), read_message(stream)).await {
Ok(r) => Ok(Some(r?)),
Err(_) => Ok(None), // 消息被吞掉了:字节已经不在流里
}
}
// 安全:把「被取消也不怕」的部分交给独立任务,取消只丢弃结果
async fn read_message_cancel_safe(
mut stream: tokio::net::TcpStream,
) -> std::io::Result<Vec<u8>> {
// spawn 出去的任务不会因为调用者被取消而停止;它自己跑完
let handle = tokio::spawn(async move {
let mut header = [0u8; 4];
stream.read_exact(&mut header).await?;
let len = u32::from_be_bytes(header) as usize;
let mut body = vec![0u8; len];
stream.read_exact(&mut body).await?;
Ok::<_, std::io::Error>(body)
});
// handle.await 得到的是 Result<Result<Vec<u8>, io::Error>, JoinError>,
// 所以要把 panic/取消这一层 io 化,才能和函数返回类型对齐
handle.await.map_err(|e| std::io::Error::other(e.to_string()))?
}
#[tokio::main]
async fn main() {
// 演示形状:真实用法是 listener.accept() 拿到的 stream
let Ok(mut stream) = tokio::net::TcpStream::connect("127.0.0.1:8080").await else {
println!("没有服务器在监听,跳过运行演示");
return;
};
match broken(&mut stream).await {
Ok(Some(body)) => println!("got {} bytes", body.len()),
Ok(None) => println!("超时:半条消息已丢失"),
Err(e) => println!("io error: {e}"),
}
}🧠 原理:
tokio::spawn出来的任务与调用者的生命周期无关。 因此「取消」和「丢弃结果」可以分开:timeout(d, handle)超时只会取消handle.await, 后台任务继续把数据读完,之后再handle.await一次就能取到完整结果。
spawn_blocking 的位置:Tokio 的 worker 线程不能被阻塞,因为一个 worker 上可能 排着成百上千个任务。所以:
| 操作类型 | 放哪里 |
|---|---|
| 短小、纯 CPU(< 10~100 µs) | 直接在 async 里做(不值得切换开销) |
| 长 CPU 计算(解析大文件、图像处理、加密) | tokio::task::spawn_blocking |
同步 IO(std::fs、std::net、FFI 调用) | spawn_blocking |
| 调用会阻塞的外部库 | spawn_blocking |
rust
// 依赖:tokio = { version = "1", features = ["full"] }
#[tokio::main]
async fn main() {
// spawn_blocking 在专用的阻塞线程池里跑,不占 worker 线程
let handle = tokio::task::spawn_blocking(|| {
// 这里可以放心用同步阻塞 API
std::fs::read_to_string("Cargo.toml").unwrap_or_default()
});
let content = handle.await.expect("blocking task panicked");
println!("{} bytes", content.len());
}⚠️ 陷阱:
spawn_blocking的默认线程池上限是 512 个线程,且不会因为队列变长而 无限扩容。把大量长任务丢进去会排队;而spawn_blocking里再await是编译不过的 (闭包不是 async),这反而是好东西——它强迫你把阻塞逻辑写干净。
Stream:异步版的迭代器
Future 对应「一个值」,Stream 对应「一串值」。它不在 std 里 (std::async_iter 尚未稳定),实际用的是 futures::Stream:
最小的消费循环只有一句:while let Some(x) = stream.next().await。下面三种写法演示 「串行」「并发但保序」「并发不保序」的区别:
rust
// 依赖:tokio = { version = "1", features = ["full"] }, futures = "0.3"
use std::time::{Duration, Instant};
use futures::stream::{self, StreamExt};
async fn fetch(id: u32) -> u32 {
tokio::time::sleep(Duration::from_millis(50)).await;
id * 10
}
#[tokio::main]
async fn main() {
// 1) 串行:先建出 5 个 future(还没跑),buffered(1) 保证同时只有一个在飞
let start = Instant::now();
let serial: Vec<u32> = stream::iter(1..=5).map(fetch).buffered(1).collect().await;
println!("serial: {serial:?} in {:?}", start.elapsed());
// 2) 并发但保序: buffered(3) 最多同时跑 3 个,输出顺序与输入一致
let start = Instant::now();
let kept: Vec<u32> = stream::iter(1..=5).map(fetch).buffered(3).collect().await;
println!("buffered: {kept:?} in {:?}", start.elapsed());
// 3) 并发且不保序:谁先完成谁先出,吞吐最高
let start = Instant::now();
let unordered: Vec<u32> = stream::iter(1..=5).map(fetch).buffer_unordered(3).collect().await;
println!("unordered: {unordered:?} in {:?}", start.elapsed());
}输出(耗时为例):
serial: [10, 20, 30, 40, 50] in 250ms 左右、buffered: [10, 20, 30, 40, 50] in 100ms 左右、unordered: [10, 20, 30, 40, 50] in 100ms 左右——本例每个 future 等长, 所以buffer_unordered看起来也是升序;一旦耗时不同(例如让fetch的 sleep 与id反相关),它会按「完成顺序」输出。
💡 对照:
futures::Stream≈ 异步的Iterator:next()返回Future<Output = Option<T>>。 因此while let Some(x) = stream.next().await是标准消费写法(需要use futures::StreamExt)。async_stream::stream!宏则让你用yield写自定义流。
tokio_stream 的角色是打通 Tokio 与 futures:Tokio 的 mpsc::Receiver 实现了 tokio_stream::Stream,用 tokio_stream::wrappers::ReceiverStream 或直接 ReceiverStream::new(rx) 就能交给 StreamExt 的组合子:
rust
// 依赖:tokio = { version = "1", features = ["full"] }, tokio-stream = "0.1"
use tokio::sync::mpsc;
use tokio_stream::wrappers::ReceiverStream;
use tokio_stream::StreamExt;
#[tokio::main]
async fn main() {
let (tx, rx) = mpsc::channel::<u32>(8);
tokio::spawn(async move {
for i in 1..=3 {
tx.send(i).await.expect("receiver dropped");
}
});
// 把 Receiver 变成 Stream,然后像迭代器一样消费
let mut stream = ReceiverStream::new(rx);
while let Some(v) = stream.next().await {
println!("stream got {v}");
}
}输出:3 行
stream got 1/2/3。