第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!、spawn 或 select!。
- 识别取消不安全、阻塞运行时、持锁跨 .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 被 .await、spawn 或其他组合器驱动时,计算才开始推进。
这类惰性有两个直接影响:
- 忘记 .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 会推进当前状态:
- 第一次进入时创建配置 Future;
- 配置未完成则保存状态并返回 Pending;
- 配置完成后保存结果,创建指标 Future;
- 指标未完成时再次返回 Pending;
- 最终组装快照并返回 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 的忙等。
正确协议是:
- poll 检查当前状态;
- 未完成时登记最新 Waker;
- 返回 Pending;
- 外部事件完成后调用 wake;
- 执行器把任务重新放入就绪队列。
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 不会自动证明:
- 内部裸指针初始化正确;
- 固定投影没有移动受保护字段;
- 引用别名合法;
- 对象释放顺序正确;
- 类型具备 Send 或 Sync。
手写 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;如果任务本来就必须留在单线程,可使用 LocalSet 与 spawn_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 分开处理。
实际系统还应增加总停机超时。超过期限后,是中止剩余任务、写入恢复日志还是延迟退出,应由业务一致性要求决定。
十一、常见陷阱
- 创建 Future 后忘记 .await。 操作可能从未开始。
- 在 async 任务中调用长时间阻塞 API。 一个任务会冻结承载它的 worker。
- 认为 .await 每次都会切换任务。 子 Future 立即完成时可能连续执行。
- 无限 spawn。 输入峰值会变成任务、内存和连接数量峰值。
- 丢弃 JoinHandle。 任务失败和完成状态失去管理者。
- 持锁跨 .await。 等待时间被转化为锁占用时间。
- 把超时理解为事务回滚。 未胜出的 Future 被取消,但外部副作用仍存在。
- 忽略 cancellation safety。 循环中的 select! 可能让半完成操作丢失进度。
- 使用 Arc 就假定 Future 为 Send。 跨 .await 的所有状态都必须满足线程迁移要求。
- 只限制 Stream 消费速度,不限制上游。 积压可能转移到通道、代理或 detached task 中。
十二、实践练习
- 编写两个各等待 300ms 的 Future,比较顺序 .await、join! 和 spawn 的耗时与错误传播方式。
- 为遥测采集增加单任务超时和总停机超时,明确两者触发时的处理差异。
- 让一个 Rc<String> 跨越 .await,阅读完整编译错误,再分别用缩小作用域、Arc 和 LocalSet 修复。
- 写一个自定义 Future,在返回 Pending 时保存 Waker;加入测试验证触发后能完成。
- 使用 Stream 处理 100 个模拟传感器,比较 buffered 与 buffer_unordered 的结果顺序。
- 构造一个持异步锁进行网络等待的例子,测量其他任务的排队时间,再改为快照模式。
- 设计一个具有明确提交点的文件写入操作,说明在提交点前后被取消分别如何恢复。
- 将 JoinSet 版本改为无界 spawn,压测并观察内存与任务数量,再恢复并发限制。
- 画出一个含三个 .await 和一个局部借用的异步函数状态图,标记各状态保存的字段。
十三、总结
- 异步让任务在等待时释放线程;它不等于多线程,也不自动产生业务并发。
- async fn 返回惰性的 Future,真正执行依赖 .await 或执行器轮询。
- Future 是一次性状态机,Ready 表示终态,Pending 表示暂时无法推进。
- 返回 Pending 的实现必须安排唤醒;Waker 将外部事件重新连接到任务队列。
- Reactor 观察事件,Executor 驱动 Future,Scheduler 分配任务执行资源。
- Pin 保护依赖地址稳定的值不被安全代码移动,Unpin 表示移动不会破坏不变量。
- join! 适合共同寿命的全部完成,spawn 创建独立任务,select! 处理竞速和取消。
- 取消通常通过 drop Future 实现,但已产生的外部副作用不会自动回滚。
- Stream 异步地产生多项结果;并发上限和端到端背压必须显式设计。
- 多线程运行时中的任务是否可 spawn,取决于跨 .await 保存的状态是否满足 Send + 'static。
- 异步系统最常见的工程问题是阻塞、任务失管、持锁、取消不安全和无界资源增长。
十四、章节导航
上一章使用线程、消息和锁组织同步并发;本章将等待过程转换成可调度状态机。下一章进入更底层的实现边界,讨论 unsafe、宏、内存布局和高级类型系统工具。
- 上一章:第9章:并发编程
- 下一章:第11章:底层能力 —— Unsafe、宏、类型系统与编译陷阱
- Rust 深度教程
- 第9章:并发编程
- 第11章:底层能力
- 第12章:Cargo 工程化