专栏 编程工程

1.5.2 Async Rust · pin / Future / Stream 实战

Rust 异步是 zero-cost abstraction。

1. 为什么需要 async:同步、多线程与 Green Thread 的取舍

I/O 密集型场景下,同步阻塞会浪费线程等待,多线程又会因锁与上下文切换撑爆内存,Rust 选择把异步作为 zero-cost abstraction(零开销抽象)挂在标准库 std::future::Future 上,运行时由 Tokio、async-std 等库按需实现。async/await 在编译期被展开成状态机,运行时只在 .await 点让出执行权,既能像 Green Thread 那样写线性代码,又不必为每个任务预分配 8MB 栈。

模式 单任务内存 并发 10k 连接 编程模型 适用场景
OS Thread ~8MB 栈 ~80GB(不可行) 同步 CPU 密集
Green Thread ~2KB ~20MB 同步 + 协作 脚本/嵌入式
Rust async ~百字节 state ~MB 级 状态机 + .await I/O 密集 + 生态成熟

依据:tokio 官方文档《Shared state》《Shared) 与 2024 Tokio 1.40 release notes 关于 task 内存预算的描述。

2. Future trait 详解:Poll / Pin / Waker 三件套

Future trait 只有一个 poll 方法,它返回 Poll<Self::Output>,其中 Ready(value) 表示完成、Pending 表示尚未就绪;Pin<&mut Self> 保证 future 在被轮询时地址不变,从而让内部的自引用结构(struct S { a: i32, b: *const i32 })安全地持有指向自身的指针;Waker 则是任务通知 executor「我已经准备好了,再来 poll 一次」的回调句柄,典型实现 wake() 会把任务重新压入 runtime 的就绪队列。运行时负责反复调用 poll,future 负责要么直接出结果,要么注册 waker 后立刻返回 Pending。

组件 类型 作用
poll Pin<&mut Self> -> Poll<Output> 推进状态机
Pin 不可移动的 &mut 指针 保证自引用有效
Waker Arc<dyn Wake>,线程安全 通知 executor 重新 poll
Context<'_> 含 &Waker 的环境 传给 poll 的轮询上下文

依据:Rust RFC 2592、std::future::Future 文档与 Pin RFC。

3. async/await 语法糖:从 async fn 到手工 Future impl

编译器把 async fn fetch(url: &str) -> String 改写成 fn fetch<'a>(url: &'a str) -> impl Future<Output = String> + 'a,内部生成一个匿名的状态机 enum:每个 .await 点对应一个 variant,局部变量变成字段,跨 await 持有引用的变量被 pin 在栈外堆上(Box::pin)。这也是为什么 await 必须在 async 上下文里、且返回的 future 必须是 'static 才能 tokio::spawn。

// 手写 Future 实现 poll,与 async fn 等价
struct DelayFut { when: Instant, waker: Option<Waker> }
impl Future for DelayFut {
    type Output = ();
    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> {
        if Instant::now() >= self.when {
            Poll::Ready(())
        } else {
            self.waker = Some(cx.waker().clone());
            Poll::Pending
        }
    }
}

依据:Rust Reference《Async/Await》《Implementation of generators and async lowering》与 rustc 编译器源代码 compiler/rustc_lexer/rustc_hir_lowering。

4. Tokio runtime 实战:spawn / select / timeout / JoinHandle

Tokio 默认 multi_thread runtime 提供 work-stealing 调度,tokio::spawn(future) 把任务丢进全局队列返回 JoinHandle<T>,调用 handle.await 即可拿到结果或 JoinError;tokio::select! 宏像 match 一样等待多个 future,谁先 Ready 就执行谁的分支并取消其余分支,适合「抢答」场景;tokio::time::timeout(Duration::from_secs(3), fut) 给 future 套一个时限,超时返回 Err(Elapsed)。

#[tokio::main]
async fn main() {
    let a = tokio::spawn(async { reqwest::get("https://a").await?.text().await });
    let b = tokio::spawn(async { reqwest::get("https://b").await?.text().await });
    let res = tokio::time::timeout(Duration::from_secs(2), async {
        tokio::select! {
            Ok(t) = a => t.unwrap_or_default(),
            Ok(t) = b => t.unwrap_or_default(),
        }
    }).await;
    println!("{res:?}");
}

依据:Tokio 教程《Spawning》《Shared state》与 Tokio 1.40 API docs。

5. Stream 异步迭代器 + 跨 Future 共享状态

Stream 是 Future 的迭代版本:trait Stream { type Item; fn poll_next(...) -> Poll<Option<Self::Item>>; },Option::None 代表流结束;tokio_stream 与 futures crate 提供 wrappers::ReceiverStream、StreamExt::map/filter/chunks_timeout 等适配器。跨 future 共享可变状态时,优先用 tokio::sync::Mutex 或消息通道(mpsc::channel),而不是 Rc<RefCell<T>>——后者不是 Send,无法跨任务;Arc<tokio::sync::RwLock<HashMap<K,V>>> 是只读多写一场景的常见选择。

共享方式 Send 跨 await 性能 典型场景
Arc<tokio::sync::Mutex<T>> ✅ ✅ 一般 写少读多、临界区短
Arc<tokio::sync::RwLock<T>> ✅ ✅ 读多写少更优 配置/缓存
mpsc::channel ✅ ✅ 高 任务间单向消息
Rc<RefCell<T>> ❌ ✅ 高 仅单线程 LocalSet

依据:Tokio 教程《Shared state》《Channels》与 futures::Stream trait docs。

说明 · 本站内容均为学习笔记与经验总结,所有菜谱与技法请结合实际食材、季节与个人口味灵活调整。涉及生食、营养与健康的内容仅供参考,特殊体质或疾病请咨询专业营养师/医生。