Agent X-Ray
RuntimeNotesAbout
Notes/代码工程/Rust 深度教材/第10章

第10章:异步编程 —— Future、Pin、Stream 与任务调度

13 分钟 · 更新于 2026-09-01

第10章:异步编程 —— Future、Pin、Stream 与任务调度

异步编程处理的核心矛盾是:任务经常需要等待网络、定时器或消息,但等待期间不应独占一条线程。Rust 将可暂停的计算表示为 Future,由运行时反复推进;事件准备好后通过 Waker 把任务送回就绪队列。本章以遥测采集服务为场景,建立从状态机到取消、背压和优雅停机的完整模型。

本章目标

  • 区分异步、并发与并行。
  • 理解 async fn 的惰性、.await 的暂停语义与 Future::poll 协议。
  • 解释 Reactor、Executor、Scheduler 和 Waker 的协作关系。
  • 理解 Pin 为何出现在 Future 接口中,以及它不保证什么。
  • 根据任务关系选择 join!spawnselect!
  • 识别取消不安全、阻塞运行时、持锁跨 .await 和无界生成任务等风险。
  • 使用 Stream 与有界并发构建可控的数据流水线。

一、异步优化的是等待,不是计算

1.1 三个容易混淆的概念

概念回答的问题可能的实现
并发多项工作能否在同一时间段内交替推进事件循环、任务调度、多线程
并行多项工作能否在同一时刻真实执行多核上的多个线程
异步等待外部事件时是否释放承载线程Future、事件通知、协作让出

单线程运行时可以并发驱动成千上万个等待网络的任务,却无法在同一时刻并行执行两段 CPU 指令。多线程程序也可能采用同步阻塞 I/O,并不具备异步语义。

1.2 等待时发生了什么

一个异步任务通常经历如下循环:

text
运行 -> 遇到尚未完成的 await -> 返回 Pending
     -> 运行时执行其他任务
     -> I/O 或计时器准备好 -> wake
     -> 重新 poll -> 继续运行

异步并没有让网络变快,也没有让等待消失。它只让线程在等待期间可以承载其他任务。

1.3 使用 async 仍可能串行

rust
use tokio::time::{sleep, Duration};

async fn collect(sensor: &str) -> String {
    sleep(Duration::from_millis(200)).await;
    format!("sample:{sensor}")
}

async fn serial_collection() {
    let cpu = collect("cpu").await;
    let memory = collect("memory").await;
    println!("{cpu}, {memory}");
}

第二次调用要等第一次结束才开始。代码具有异步等待能力,但业务流程仍是顺序的。

并发是额外的结构选择,而不是 async 关键字的副作用。

1.4 CPU 密集工作需要另一类执行资源

rust
let report = tokio::task::spawn_blocking(move || {
    compress_samples(samples)
})
.await
.expect("blocking task panicked");

大规模压缩、图像处理、密码计算或长时间同步阻塞会占住运行时 worker。应使用 spawn_blocking、Rayon、专用线程池或独立计算服务,并对并发数量设置上限。


二、async fn 创建惰性的 Future

2.1 调用与执行是两件事

rust
async fn load_threshold() -> u32 {
    println!("loading threshold");
    80
}

#[tokio::main]
async fn main() {
    let pending = load_threshold();
    println!("future constructed");

    let value = pending.await;
    println!("threshold = {value}");
}

先打印 future constructed,之后才进入函数体。调用 async fn 只是构造一个 Future;Future 被 .awaitspawn 或其他组合器驱动时,计算才开始推进。

这类惰性有两个直接影响:

  • 忘记 .await 可能意味着操作根本没有执行;
  • 可以先构造多个 Future,再决定顺序、并发、超时或竞速关系。

2.2 概念上的函数签名

rust
async fn read_metric(name: String) -> Result<u64, MetricError> {
    query_backend(name).await
}

可以近似理解为:

rust
fn read_metric(
    name: String,
) -> impl std::future::Future<Output = Result<u64, MetricError>> {
    async move {
        query_backend(name).await
    }
}

编译器会为匿名 Future 生成具体类型。参数、跨暂停点仍存活的局部变量以及子 Future 都可能成为状态机字段。

2.3 async move 把捕获值放进状态机

rust
let tenant = String::from("tenant-a");

let handle = tokio::spawn(async move {
    send_snapshot(tenant).await
});

handle.await??;

独立任务可能比当前函数活得更久,因此通常不能借用当前栈上的短期变量。async move 把捕获值移入 Future。

tokio::spawn 常要求 Future 满足 Send + 'static

  • Send:任务暂停后可能转移到其他 worker 线程;
  • 'static:Future 不依赖寿命不足的外部借用。

'static 不表示任务必须运行到进程结束。


三、Future 是可暂停的一次性计算

3.1 核心接口

rust
use std::pin::Pin;
use std::task::{Context, Poll};

pub trait Future {
    type Output;

    fn poll(
        self: Pin<&mut Self>,
        cx: &mut Context<'_>,
    ) -> Poll<Self::Output>;
}

poll 只有两类结果:

  • Poll::Ready(output):计算完成,产生最终输出;
  • Poll::Pending:目前无法继续,未来可能再次推进。

Future 通常是一次性的。进入 Ready 后,调用者不应再次轮询;具体实现可能 panic,也可能有其他不保证行为。

3.2 两个 await 点如何变成状态

rust
async fn build_snapshot() -> Result<Snapshot, Error> {
    let config = load_config().await?;
    let metrics = load_metrics(&config).await?;
    Ok(Snapshot { config, metrics })
}

概念上可以拆为:

text
Start
  -> WaitingConfig { config_future }
  -> WaitingMetrics { config, metrics_future }
  -> Complete

每次 poll 会推进当前状态:

  1. 第一次进入时创建配置 Future;
  2. 配置未完成则保存状态并返回 Pending
  3. 配置完成后保存结果,创建指标 Future;
  4. 指标未完成时再次返回 Pending
  5. 最终组装快照并返回 Ready

.await 仍需使用的 config 必须保存在状态机中;只在暂停点前使用的局部值则无需保存。

3.3 .await 不一定真正让出

.await 会轮询子 Future。如果子 Future 立即返回 Ready,当前任务可能继续执行,并不会切换到其他任务。

因此,以下循环即使写在 async fn 中,也可能长期霸占 worker:

rust
async fn compute_forever() {
    loop {
        run_large_cpu_batch();
    }
}

协作式调度要求任务在合理时间内返回运行时。CPU 工作应分块、显式 yield_now(),或移交到专用执行资源。


四、Waker:Pending 之后如何回来

4.1 返回 Pending 必须安排后续通知

如果 Future 只返回 Pending,却没有保存并触发 Waker,它可能永远不会再被轮询。反过来,如果执行器不停主动轮询所有 Pending Future,就会形成耗费 CPU 的忙等。

正确协议是:

  1. poll 检查当前状态;
  2. 未完成时登记最新 Waker;
  3. 返回 Pending
  4. 外部事件完成后调用 wake
  5. 执行器把任务重新放入就绪队列。

4.2 一个教学用信号 Future

rust
use std::future::Future;
use std::pin::Pin;
use std::sync::{Arc, Mutex};
use std::task::{Context, Poll, Waker};

struct SignalState {
    ready: bool,
    waker: Option<Waker>,
}

struct Signal {
    state: Arc<Mutex<SignalState>>,
}

impl Future for Signal {
    type Output = ();

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> {
        let mut state = self.state.lock().unwrap();

        if state.ready {
            Poll::Ready(())
        } else {
            let should_replace = state
                .waker
                .as_ref()
                .map_or(true, |old| !old.will_wake(cx.waker()));

            if should_replace {
                state.waker = Some(cx.waker().clone());
            }

            Poll::Pending
        }
    }
}

外部生产者将 ready 设为 true 后取出并唤醒 Waker:

rust
fn trigger(state: &Arc<Mutex<SignalState>>) {
    let waker = {
        let mut state = state.lock().unwrap();
        state.ready = true;
        state.waker.take()
    };

    if let Some(waker) = waker {
        waker.wake();
    }
}

先释放锁再调用 wake,可以避免唤醒链条中的代码与当前锁形成意外耦合。

示例边界 真实 I/O 运行时不会为每个等待创建一条线程,而会借助 epoll、kqueue、IOCP、timer wheel 等机制批量观察事件。这里展示的是 poll 与唤醒协议,而非高性能实现。

4.3 运行时组件的分工

组件主要职责
Reactor观察网络、计时器、信号等外部就绪事件
Waker将就绪事件映射回具体任务
Executor取得就绪任务并调用 poll
Scheduler决定任务放在哪个线程、以何种顺序运行
Task包装 Future 及其调度状态

标准库定义 Future 协议,不规定运行时实现。Tokio 等运行时将这些能力组合为可使用的系统。


五、Pin:保护依赖地址稳定的状态机

5.1 为什么状态机可能不允许移动

rust
async fn parse_packet() {
    let buffer = String::from("alpha,beta");
    let first = &buffer[..5];
    send_field(first).await;
}

在暂停期间,生成的状态机可能同时保存:

  • 拥有数据的 buffer
  • 指向 buffer 内部的 first
  • 被等待的子 Future。

若这个状态机在建立内部引用后被移动,内部地址关系可能被破坏。

5.2 Pin<P> 提供的承诺

Pin<P> 限制通过指针 P 移动目标值。常见形态是:

rust
Pin<&mut T>
Pin<Box<T>>

它不是把内存固定到物理页,也不是阻止操作系统移动虚拟内存。它限制的是 Rust 语义中的值移动。

Future::poll 接收 Pin<&mut Self>,使实现可以在假设自身不再移动的前提下访问状态。

5.3 Unpin 表示移动不会破坏不变量

大部分普通类型都自动实现 Unpin。对于 T: Unpin,固定不会增加实质限制,可以安全取得普通可变访问。

编译器生成的某些 Future 可能是 !Unpin,需要先固定:

rust
let operation = async {
    collect_once().await
};

tokio::pin!(operation);
let output = operation.await;

需要在堆上保存或做动态分发时:

rust
use std::future::Future;
use std::pin::Pin;

type BoxedJob = Pin<Box<dyn Future<Output = u64> + Send>>;

fn boxed_job() -> BoxedJob {
    Box::pin(async { 42 })
}

5.4 Pin 的边界

Pin 不会自动证明:

  • 内部裸指针初始化正确;
  • 固定投影没有移动受保护字段;
  • 引用别名合法;
  • 对象释放顺序正确;
  • 类型具备 SendSync

手写 unsafe 固定投影或自引用结构时,仍要维护完整不变量。应用层通常应依赖成熟抽象,而不是自行实现。


六、组合 Future:join、spawn 与 select

6.1 join!:同一任务内等待全部分支

rust
async fn collect_dashboard() -> Result<Dashboard, Error> {
    let cpu = fetch_cpu();
    let memory = fetch_memory();

    let (cpu, memory) = tokio::try_join!(cpu, memory)?;
    Ok(Dashboard { cpu, memory })
}

两个 Future 在同一个 Task 中交替被轮询。它们具有并发性,但不必成为独立任务,也不保证多核并行。

适合:

  • 分支寿命与当前函数一致;
  • 需要等待全部结果;
  • 不需要独立任务句柄。

6.2 spawn:创建独立调度单元

rust
let writer = tokio::spawn(async move {
    persist_batch(batch).await
});

continue_collecting().await?;
writer.await??;

独立 Task 有自己的调度状态和 JoinHandle。它适合可以独立推进、需要并发执行或需要单独取消与观察的工作。

不要把 spawn 当作“让函数异步”的装饰。创建 Task 会改变错误传播、取消关系、所有权边界和关停责任。

6.3 select!:等待第一个可处理分支

rust
use tokio::time::{sleep, Duration};

async fn collect_with_deadline() -> Result<Sample, CollectError> {
    tokio::select! {
        result = read_sensor() => result,
        _ = sleep(Duration::from_secs(1)) => Err(CollectError::Timeout),
    }
}

select! 常用于:

  • 请求与超时竞速;
  • 数据输入与关闭信号竞速;
  • 多个接收端择一处理;
  • 主任务与健康检查竞速。

某个分支完成并匹配后,其他分支通常被丢弃。因此 select! 同时也是取消边界。


七、取消:drop Future 不等于回滚世界

7.1 Future 被丢弃后

Future 内部尚未执行的步骤不会继续;局部值和子 Future 会被析构。但已经发生的外部副作用不会自动撤回:

  • 已发送的字节仍可能抵达服务器;
  • 已写入的文件可能只完成一部分;
  • 已提交的数据库语句可能已经生效;
  • spawn 的独立 Task 不会自动随父 Future 消失。

因此,取消设计需要明确“提交点”。提交点之前可以安全重试或丢弃,提交点之后则需要幂等键、状态查询或补偿操作。

7.2 循环中的取消安全

rust
loop {
    tokio::select! {
        message = receiver.recv() => {
            // 处理消息
        }
        _ = shutdown.changed() => break,
    }
}

如果 recv() 在被取消时会丢失已经读取但尚未返回的数据,这个循环就不安全。应阅读 API 文档中的 cancellation safety 说明,并避免把不可中断的多步骤协议直接放入反复竞速的分支。

7.3 结构化并发

稳健的任务树应让父作用域知道所有子任务:

  • 保存 JoinHandle
  • 使用 JoinSet 管理动态任务;
  • 关闭时停止接收新工作;
  • 广播取消信号;
  • 给等待子任务设置总超时;
  • 记录未能正常结束的任务。

“发射后不管”的 detached task 会让错误、取消和资源寿命失去归属。


八、Stream:异步地产生多项结果

8.1 与 Iterator 的差异

同步迭代器的下一项立即可计算:

rust
fn next(&mut self) -> Option<Self::Item>;

Stream 的下一项可能尚未到达:

rust
use std::pin::Pin;
use std::task::{Context, Poll};

pub trait Stream {
    type Item;

    fn poll_next(
        self: Pin<&mut Self>,
        cx: &mut Context<'_>,
    ) -> Poll<Option<Self::Item>>;
}

返回值含义:

  • Ready(Some(item)):产生一项;
  • Ready(None):流结束;
  • Pending:下一项未准备好,稍后唤醒。

8.2 遥测流水线

rust
use futures::{stream, StreamExt};

let results = stream::iter(sensor_ids)
    .map(|sensor_id| async move {
        query_sensor(sensor_id).await
    })
    .buffer_unordered(12)
    .collect::<Vec<_>>()
    .await;

buffer_unordered(12) 最多同时驱动 12 个查询,并允许结果按完成顺序输出。这个数字是系统容量参数:过小会浪费可用资源,过大可能压垮下游或耗尽连接池。

8.3 适配器顺序就是业务含义

rust
let stream = tokio_stream::iter(1..=100)
    .filter(|value| value % 5 == 0)
    .take(3);

这里先过滤再取前三项,结果是 5、10、15。若交换顺序,则先截取 1、2、3,随后一项也没有。

Stream 与 Iterator 一样是惰性管道。每个适配器放置的位置都会影响资源使用、取消边界和输出语义。

8.4 背压不会自动覆盖全链路

消费者通过 next().await 拉取下一项,能对当前 Stream 形成一定背压。但如果上游来自无界通道、消息代理或不断 spawn 的生产者,数据仍可能在别处积压。

要形成端到端背压,需要同时限制:

  • 输入队列容量;
  • 同时进行的 I/O 数量;
  • 单项处理产生的子任务数量;
  • 输出缓存和下游写入速度。

九、Send、锁与异步作用域

9.1 非 Send 值跨 await

rust
use std::rc::Rc;

async fn local_work() {
    let label = Rc::new(String::from("local"));
    tokio::task::yield_now().await;
    println!("{label}");
}

如果把该 Future 交给多线程运行时的 tokio::spawn,会因 Rc<String> 跨越 .await 而不满足 Send

若该值只在暂停前使用,可以缩小作用域:

rust
async fn sendable_work() {
    {
        let label = Rc::new(String::from("temporary"));
        println!("{label}");
    }

    tokio::task::yield_now().await;
}

如果业务确实需要跨任务共享,可改用 Arc;如果任务本来就必须留在单线程,可使用 LocalSetspawn_local。两种方案表达的是不同架构,不能只按“哪种能编译”选择。

9.2 不要让锁 guard 跨 await

rust
let snapshot = {
    let guard = state.lock().await;
    guard.clone_for_upload()
};

upload(snapshot).await?;

先在短临界区提取快照,再释放锁进行远程调用。若把 guard 保留到 upload().await 之后,等待时间会被放大为锁占用时间,甚至与其他任务形成循环等待。

标准库 Mutex 还可能让整个运行时 worker 阻塞。异步代码应根据临界区特征选择:

  • 极短、绝不跨 .await:标准同步锁可能合适;
  • 必须异步等待锁:使用运行时提供的异步锁;
  • 所有权可以移交:优先 channel,减少共享状态。

十、一个可关闭的遥测采集循环

rust
use tokio::sync::{mpsc, watch};
use tokio::task::JoinSet;

#[derive(Debug)]
struct SensorRequest {
    id: u64,
}

async fn run_collector(
    mut requests: mpsc::Receiver<SensorRequest>,
    mut shutdown: watch::Receiver<bool>,
) {
    let mut tasks = JoinSet::new();
    let max_in_flight = 16;

    loop {
        tokio::select! {
            changed = shutdown.changed() => {
                if changed.is_err() || *shutdown.borrow() {
                    break;
                }
            }
            Some(result) = tasks.join_next(), if !tasks.is_empty() => {
                match result {
                    Ok(Ok(sample)) => store_sample(sample).await,
                    Ok(Err(error)) => eprintln!("collection failed: {error}"),
                    Err(join_error) => eprintln!("task failed: {join_error}"),
                }
            }
            request = requests.recv(), if tasks.len() < max_in_flight => {
                match request {
                    Some(request) => {
                        tasks.spawn(async move {
                            collect_sensor(request).await
                        });
                    }
                    None => break,
                }
            }
        }
    }

    requests.close();

    while let Some(result) = tasks.join_next().await {
        if let Err(error) = result {
            eprintln!("task failed during shutdown: {error}");
        }
    }
}

这个循环体现了几项工程约束:

  • mpsc::Receiver 提供有界输入和背压基础;
  • JoinSet 让动态子任务仍受父循环管理;
  • tasks.len() 限制同时进行的采集数量;
  • 关闭后停止接收新请求,再等待存量任务完成;
  • 任务错误和 JoinError 分开处理。

实际系统还应增加总停机超时。超过期限后,是中止剩余任务、写入恢复日志还是延迟退出,应由业务一致性要求决定。


十一、常见陷阱

  1. 创建 Future 后忘记 .await 操作可能从未开始。
  2. 在 async 任务中调用长时间阻塞 API。 一个任务会冻结承载它的 worker。
  3. 认为 .await 每次都会切换任务。 子 Future 立即完成时可能连续执行。
  4. 无限 spawn 输入峰值会变成任务、内存和连接数量峰值。
  5. 丢弃 JoinHandle 任务失败和完成状态失去管理者。
  6. 持锁跨 .await 等待时间被转化为锁占用时间。
  7. 把超时理解为事务回滚。 未胜出的 Future 被取消,但外部副作用仍存在。
  8. 忽略 cancellation safety。 循环中的 select! 可能让半完成操作丢失进度。
  9. 使用 Arc 就假定 Future 为 Send。.await 的所有状态都必须满足线程迁移要求。
  10. 只限制 Stream 消费速度,不限制上游。 积压可能转移到通道、代理或 detached task 中。

十二、实践练习

  1. 编写两个各等待 300ms 的 Future,比较顺序 .awaitjoin!spawn 的耗时与错误传播方式。
  2. 为遥测采集增加单任务超时和总停机超时,明确两者触发时的处理差异。
  3. 让一个 Rc<String> 跨越 .await,阅读完整编译错误,再分别用缩小作用域、ArcLocalSet 修复。
  4. 写一个自定义 Future,在返回 Pending 时保存 Waker;加入测试验证触发后能完成。
  5. 使用 Stream 处理 100 个模拟传感器,比较 bufferedbuffer_unordered 的结果顺序。
  6. 构造一个持异步锁进行网络等待的例子,测量其他任务的排队时间,再改为快照模式。
  7. 设计一个具有明确提交点的文件写入操作,说明在提交点前后被取消分别如何恢复。
  8. JoinSet 版本改为无界 spawn,压测并观察内存与任务数量,再恢复并发限制。
  9. 画出一个含三个 .await 和一个局部借用的异步函数状态图,标记各状态保存的字段。

十三、总结

  1. 异步让任务在等待时释放线程;它不等于多线程,也不自动产生业务并发。
  2. async fn 返回惰性的 Future,真正执行依赖 .await 或执行器轮询。
  3. Future 是一次性状态机,Ready 表示终态,Pending 表示暂时无法推进。
  4. 返回 Pending 的实现必须安排唤醒;Waker 将外部事件重新连接到任务队列。
  5. Reactor 观察事件,Executor 驱动 Future,Scheduler 分配任务执行资源。
  6. Pin 保护依赖地址稳定的值不被安全代码移动,Unpin 表示移动不会破坏不变量。
  7. join! 适合共同寿命的全部完成,spawn 创建独立任务,select! 处理竞速和取消。
  8. 取消通常通过 drop Future 实现,但已产生的外部副作用不会自动回滚。
  9. Stream 异步地产生多项结果;并发上限和端到端背压必须显式设计。
  10. 多线程运行时中的任务是否可 spawn,取决于跨 .await 保存的状态是否满足 Send + 'static
  11. 异步系统最常见的工程问题是阻塞、任务失管、持锁、取消不安全和无界资源增长。

十四、章节导航

上一章使用线程、消息和锁组织同步并发;本章将等待过程转换成可调度状态机。下一章进入更底层的实现边界,讨论 unsafe、宏、内存布局和高级类型系统工具。

  • 上一章:第9章:并发编程
  • 下一章:第11章:底层能力 —— Unsafe、宏、类型系统与编译陷阱

  • Rust 深度教程
  • 第9章:并发编程
  • 第11章:底层能力
  • 第12章:Cargo 工程化

本章目录
一、异步优化的是等待,不是计算二、async fn 创建惰性的 Future三、Future 是可暂停的一次性计算四、Waker:Pending 之后如何回来五、Pin:保护依赖地址稳定的状态机六、组合 Future:join、spawn 与 select七、取消:drop Future 不等于回滚世界八、Stream:异步地产生多项结果九、Send、锁与异步作用域十、一个可关闭的遥测采集循环十一、常见陷阱十二、实践练习十三、总结十四、章节导航Related Documents
苏ICP备2025204887号-2