Skip to content

Tokio 实战

用 Tokio 写真实并发任务:spawn、超时、并发组合与取消。

Tokio 实战

tokio::fstokio::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'—— 先说「没有这个方法」,再提示「trait AsyncReadExt is 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:8080new connection from ..., 客户端打印 server says: pingserver says: rust

一个服务器也可以接受多个客户端同时连接,每个连接一个任务。用 tokio::spawnstreammove 进任务,因此不需要锁——这正是 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::syncstd 对照

tokio::sync 里的类型都有异步版本,因为它们的「等待」也要能挂起任务而不是挂起线程。 同一个功能有多少个「版本」非常容易混,先记这张表:

需求std 对应物Tokio 对应物什么时候用 Tokio 版本
互斥访问共享数据std::sync::Mutextokio::sync::Mutex需要在.await 的临界区里持锁
读写锁std::sync::RwLocktokio::sync::RwLock同上,且读多写少
一次性值/完成通知std::sync::OnceLockCondvartokio::sync::OnceCellNotify等待发生在 async 上下文
多生产者单消费者队列std::sync::mpsctokio::sync::mpsc需要异步 recv().await,或需要背压
单次应答(请求-响应)(无直接对应)tokio::sync::oneshot一个值、一次、异步等待
广播(无)tokio::sync::broadcast一个发送者、多个订阅者
多播/可多消费者Arc<Mutex<VecDeque>>tokio::sync::watchasync_channel状态广播、配置热更新
并发限流信号量需自己写tokio::sync::Semaphore限制同时进行的请求数
屏障std::sync::Barriertokio::sync::Barrier等待 N 个任务到齐
只读共享Arc<T>Arc<T>(相同)——

tokio::sync::Mutexstd::sync::Mutex 的真正差别

std::sync::Mutextokio::sync::Mutex
等待时行为阻塞 OS 线程挂起任务,让出 worker 线程
守卫是否 Send否(MutexGuard: !Send是(MutexGuard: Send
性能(无竞争)更快(futex,几十 ns)更慢(要登记 waker)
.await 持锁编译错误(future 变成 !Send允许,但要警惕死锁

⚠️ 陷阱(重要规则)不要跨 .await 持有 std::sync::Mutex 的守卫MutexGuard 不是 Send,一旦跨挂起点存活,整个 future 就不再 Sendtokio::spawn 会直接拒绝编译(错误细节见「跨 .awaitMutexGuardfuture 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(&notify);
        tokio::spawn(async move {
            notify.notified().await;
            "woken"
        })
    };
    notify.notify_one();
    println!("{}", waiter.await.expect("join failed"));
}

输出(顺序因调度略有不同):3 行 got 0/1/2reply: 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用途一行说明安装
axumWeb 框架Tokio 官方团队维护,基于 tower + hyper,路由与提取器极简cargo add axum
reqwestHTTP 客户端异步 + 同步双形态,默认 rustls/tls,json() 直接反序列化cargo add reqwest --features json
sqlx数据库纯 Rust 异步驱动,编译期校验 SQL(需 DATABASE_URLcargo add sqlx --features runtime-tokio,sqlite
tower中间件/服务抽象Service trait:超时、重试、限流、负载均衡都可叠加cargo add tower --features timeout,retry
tracing结构化日志span + event 模型,异步任务追踪靠 tracing::instrumentcargo add tracing tracing-subscriber
serde序列化#[derive(Serialize, Deserialize)],几乎所有 Web/DB crate 的共同依赖cargo add serde --features derive
hyper底层 HTTPaxum 的底层实现,需要极致控制时直接用cargo add hyper --features server,http1
futures组合子工具库join_allStreamExtblock_on,与运行时无关cargo add futures
tokio-streamStream 适配mpsc::Receiver 等变成 Stream,与 futures::Stream 打通cargo add tokio-stream
async-traittrait 里的 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 + JoinSetJoinSet 会自动 在 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 channelbiased: tick wins

取消安全(cancel safety) 指的是:一个 future 在被 drop(在任意挂起点之后)时, 是否还能保证「没完成也不留下半成品状态」。

  • 安全(cancel safe):mpsc::Receiver::recv(没收到就没有消费)、TcpStream::connecttokio::time::sleepMutex::lockread_exact 之外的「要么全做、要么不做」的操作。 Tokio 文档对每个方法都标了 cancel safe 与否,用 select! 前必须查一遍。
  • 不安全:AsyncReadExt::read / read_exactAsyncBufReadExt::read_linelines().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::fsstd::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 对应「一串值」。它不在 stdstd::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 ≈ 异步的 Iteratornext() 返回 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


内容以 rustc 1.98.1 · Rust 2024 edition 为基准