Skip to main content

Rust 高并发实战(2):Future、Poll/Waker 与 Tokio 结构化并发

Rainy
雨落无声,代码成诗 —— 致力于技术与艺术的极致平衡
Rainy
7 MIN READ... VIEWS

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 函数大致保存 socketbuffer 和当前执行阶段:

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));

十一、结构化取消测试矩阵

取消点应释放需验证副作用
等 semaphorewaiter/排队位置无业务调用
网络请求中permit、response future远端是否继续
收到响应后解析buffer、permit是否已记账
DB commit 后本地 task事务已成功,重试需幂等
shutdownJoinSet 中全部 taskdrain/强退边界

十二、后台任务的 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 生产级有界服务

参考资料

Logo
RainLib

Exploring the frontiers of technology, design, and distributed systems. Building tools for the future developers.

Suggestions & Feedback

© 2026 RainLib. Built for the Future.
All rights reserved.
System Normal