Rust 高并发实战(2):Future、Poll/Waker 与 Tokio 结构化并发
async fn返回的是一台惰性状态机;只有 executor 持续 poll,它才会向前执行。
本期理解 Tokio task 的调度与取消语义,并实现不泄漏子任务的 fan-out/fan-in。
一、Future 如何前进
trait Future {
type Output;
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>;
}
Poll::Ready(value):完成;Poll::Pending:尚未完成,并安排 waker 在资源就绪时通知 executor;Pin:确保可能自引用的 Future 被 poll 后不再随意移动;- executor 不应忙轮询 Pending Future,错误自唤醒会造成 CPU 空转。
1.1 Tokio runtime 的组成
一次 .await 不是自动切线程:当前 worker poll Future,Future 返回 Pending 后 worker 才能运行其他 task;I/O driver/timer 在就绪时调用 waker 把 task 放回队列。CPU 循环不返回 Pending,就一直占住 worker。
二、协作式调度的红线
// 错误:没有 await 的大循环长期占住 worker。
async fn crunch(mut n: u64) -> u64 {
while n > 0 { n -= 1; }
n
}
CPU 任务进入专用池;只能切片的算法定期 yield_now().await,但 yield 不是容量控制。同步 sleep 必须换成 tokio::time::sleep。
三、请求内并发:优先组合 Future
async fn aggregate(state: AppState) -> Result<(User, Orders, Score), Error> {
let user = state.load_user();
let orders = state.load_orders();
let score = state.load_score();
tokio::try_join!(user, orders, score)
}
try_join! 在当前 task 内并发 poll 子 Future,不产生 detached task;外层 Future 被 drop 时子 Future 一起 drop。适合固定少量子操作。
四、动态任务与 JoinSet
async fn fan_out(ids: Vec<u64>) -> Result<Vec<Item>, Error> {
let mut tasks = tokio::task::JoinSet::new();
for id in ids.into_iter().take(32) {
tasks.spawn(async move { fetch_item(id).await });
}
let mut items = Vec::new();
while let Some(joined) = tasks.join_next().await {
match joined {
Ok(Ok(item)) => items.push(item),
Ok(Err(error)) => {
tasks.shutdown().await;
return Err(error);
}
Err(join_error) => {
tasks.shutdown().await;
return Err(join_error.into());
}
}
}
Ok(items)
}
限制输入数量,观察每个 JoinError,并在失败路径停止剩余任务。tokio::spawn 返回的 handle 被直接丢弃会让任务继续执行,错误和 panic 也失去 owner。
五、select! 与取消
tokio::select! {
biased;
_ = cancel.cancelled() => Err(Error::Cancelled),
result = receive_and_process() => result,
_ = tokio::time::sleep(deadline) => Err(Error::Timeout),
}
未胜出的分支通常被 drop。检查每个操作的 cancellation safety:读取是否已消费部分数据、发送是否丢失队列位置、副作用是否已经提交。
timeout 停止等待不等于远程事务回滚,也不等于 spawn_blocking 停止。写操作要有幂等键、事务状态或补偿。
六、Channel 与背压
let (tx, mut rx) = tokio::sync::mpsc::channel::<Job>(1024);
tx.send(job).await?; // 满时异步等待,传播背压
while let Some(job) = rx.recv().await {
process(job).await;
}
有界 channel 把等待显式化。决定满时等待多久、是否 try_send 快速拒绝、receiver 关闭后 producer 如何退出。不要用 unbounded_channel 隐藏容量问题。
七、Future 状态机展开
下面的 async 函数大致保存 socket、buffer 和当前执行阶段:
async fn read_frame(socket: &mut TcpStream) -> io::Result<Frame> {
let header = socket.read_u32().await?;
let mut body = vec![0; header as usize];
socket.read_exact(&mut body).await?;
decode(body)
}
每个跨 await 仍存活的局部变量都会成为 Future 状态的一部分。大 buffer 跨多个 await 会增大每个 task 内存;可把 buffer 放入池、缩短作用域或流式处理,但先测量,不为减少 Future 大小牺牲清晰性。
八、Runtime Builder 与 worker 数
let runtime = tokio::runtime::Builder::new_multi_thread()
.worker_threads(4)
.max_blocking_threads(64)
.thread_name("api-runtime")
.enable_all()
.build()?;
worker threads 从 CPU quota 附近开始压测;async I/O 连接数不要求同等 worker 数。max_blocking_threads 不是 CPU 并行度,应另外用 semaphore。过多 worker 会增加争用和 cache miss,过少会放大单 task 阻塞。
九、Wake Storm 与忙轮询
自定义 Future/Stream 若在没有新进展时反复 wake_by_ref,executor 会不停 poll,形成高 CPU、低吞吐。正确实现要在资源状态变化时唤醒,并处理“检查状态—注册 waker”之间的竞态,防止丢唤醒。
应用层常见类似问题:零间隔重试循环、总是 ready 的 select 分支、channel closed 后未退出的 recv 循环。tokio-console 的 poll 次数和 busy duration 能提供证据。
十、任务局部变量与追踪上下文
线程本地变量不适合标识 async 请求,因为 task 会迁移线程。使用 tracing span 或 task-local;跨服务通过标准 trace headers 传播。spawn 新 task 时明确继承/进入 span:
let span = tracing::info_span!("child", item_id);
tasks.spawn(async move { process(item_id).await }.instrument(span));
十一、结构化取消测试矩阵
| 取消点 | 应释放 | 需验证副作用 |
|---|---|---|
| 等 semaphore | waiter/排队位置 | 无业务调用 |
| 网络请求中 | permit、response future | 远端是否继续 |
| 收到响应后解析 | buffer、permit | 是否已记账 |
| DB commit 后 | 本地 task | 事务已成功,重试需幂等 |
| shutdown | JoinSet 中全部 task | drain/强退边界 |
十二、后台任务的 Owner 模式
struct BackgroundTasks {
cancel: CancellationToken,
tasks: JoinSet<Result<(), Error>>,
}
impl BackgroundTasks {
async fn shutdown(mut self, wait: Duration) {
self.cancel.cancel();
let drain = async {
while let Some(result) = self.tasks.join_next().await {
if let Err(error) = result { tracing::error!(?error, "background task failed"); }
}
};
if tokio::time::timeout(wait, drain).await.is_err() {
self.tasks.abort_all();
}
}
}
后台刷新、消息消费者、健康检查等必须登记到 owner。启动失败要回收已启动任务,运行中 panic 要进入监控,停机有取消、drain 和强退边界。
12.1 Drop Future 的精确含义
Drop 会销毁 Future 状态和它持有的 RAII 资源,例如 semaphore permit、buffer、socket handle 引用;它不会撤回已经发送的网络包、数据库 commit 或已经开始的 spawn_blocking closure。
因此取消正确性要分两层:本进程资源是否释放;外部副作用是否幂等、事务化或可补偿。只看到 permit 回收,不能证明订单没有重复创建。
12.2 select! 的错误用法
循环中把非 cancel-safe 的 Future 每次重新创建,可能反复丢失部分进度。正确做法是把需要保留状态的 Future pin 在循环外,或使用文档明确 cancel-safe 的 recv/read API。对消息协议,给每条消息序号并在处理完成后 ack,避免取消造成“已取出但未处理”静默丢失。
十三、避免锁住 Runtime
不要从 async 上下文随意调用 Handle::block_on 或嵌套 runtime;这可能 panic 或死锁。同步/异步桥接应在边界选择 spawn_blocking、专用线程,或由同步入口拥有 runtime,再通过 channel 交换数据。
十四、本期验收
- 用
try_join!实现固定三路聚合; - 用 JoinSet 实现最多 32 路动态 fan-out,并覆盖错误/panic/取消;
- 在每个
.await注入取消,记录副作用状态; - 压满 mpsc,确认发送等待和 queue depth 可观测。
上一篇:Rust-1:所有权与并发成本 · 下一篇:Rust-3:Axum 生产级有界服务。