Rust 无栈协程如何暂停和恢复?从 Future 与 poll 走完一次 HTTP 请求
调用一个 async 函数,为什么 HTTP 请求还没发出去?函数等网络时,线程去了哪里?调用栈都退出了,变量又怎样留到下一次执行?
Rust 的无栈协程要解决的就是这组问题:开发者按顺序写代码,编译器把跨暂停的计算变成 Future 状态机;执行器或父 Future 调用 poll,让它获得执行机会。暂时不能继续时返回 Pending,之后按保存的状态推进。
“无栈”不代表不使用线程栈。每次执行 poll 仍有普通函数调用栈,只是不为每份等待中的计算保留一份专属协程栈。下面沿一次 HTTP 请求,把调用者、状态保存、通知与恢复分别讲清,不要求你先了解 Go 或 Tokio 源码。
示例使用 rustc 1.94.0、reqwest 0.12.28 与 Tokio 1.53.2。reqwest 使用异步客户端并开启适当的 TLS 功能;涉及 spawn、定时器的片段需要已有 Tokio 运行时及对应功能。正文中的 enum 和 poll 控制流是示意,真实 MIR 单独折叠展示。本文解释执行机制,不做 HTTP 性能比较。
一、先看一次请求要完成什么
我们要完成一件简单的事:请求一个地址,把响应正文读成字符串。下面使用 reqwest 0.12.28,两个问号将请求或读取错误交给调用者:
async fn fetch_data(url: &str) -> Result<String, reqwest::Error> {
let resp = reqwest::get(url).await?; // 等待取得响应
let text = resp.text().await?; // 等待读完正文
Ok(text)
}第一个 await 完成时,已经取得 Response,可以访问状态码和响应头,但正文未必读完。第二个 await 才收集完整正文并返回字符串。两个等待点对应两个容易辨认的阶段:等待响应,读取正文。 为了突出转换过程,这里使用 get 简写;连续请求通常应复用 Client,而不是每次创建一个。
Rust 的 async/await 可以理解为编译器提供的语法糖:底层用实现 Future 的状态对象表示计算,再通过 poll 推进。转换在编译期完成,不是 Tokio 在运行时解释这段源码。
二、调用 async 函数,为什么请求还没有发出?
let future = fetch_data("https://example.com");这次调用构造并返回一个实现 Future 的匿名状态机,暂不执行 fetch_data 的函数体,也不会立刻发出这次 HTTP 请求。参数表达式仍正常求值;这里捕获的是传入的 url 引用。真正执行函数体,要等调用者来 poll 这份 Future。
三、谁调用 poll,执行 Future 里的代码?
先看标准库 Future trait 的核心定义:
pub trait Future {
type Output;
fn poll(
self: Pin<&mut Self>,
cx: &mut Context<'_>,
) -> Poll<Self::Output>;
}poll 是 Future 提供给调用者的推进方法,不会因为“属于 Future”就自动执行。对 async 函数生成的 Future,编译器把函数体的执行逻辑放进这个入口:第一次从开头执行,以后根据保存的阶段继续,不是重新调用 fetch_data(url)。
顶层 Future:由执行器调用
把 fetch_data 的 Future 交给 tokio::spawn 后,Tokio 将它作为 task 管理。工作线程取得这个可运行 task,才调用其顶层 Future 的 poll。调度过程可以简化为下面几步,注意这不是 Tokio 源码:
取得一个可运行 task
用该 task 关联的 Waker 构造 Context
调用已固定的顶层 Future.poll(context)
Ready(result):记录任务完成,交付结果
Pending:结束本轮推进,等待或处理再次推进的通知执行器不是每隔一会儿询问所有 Future。任务因提交或通知获得执行机会;Pending 后的再次推进由通知和调度状态协调。Future 本身不是一条线程,也不要求为每份计算创建专属线程;执行器或父 Future 决定何时调用它。
子 Future:由父 Future 的 poll 调用
fetch_data 执行到 reqwest::get(url).await 时,父状态机调用请求子 Future 的 poll;执行到 resp.text().await 时,则调用正文子 Future 的 poll。运行时不需要把这些子 Future 逐个当成独立任务调度。
在这个例子作为顶层任务提交给 Tokio 后,调用方向是:
flowchart TB
A["执行器<br/>安排顶层 task"]
B["父 Future.poll<br/>执行到 await"]
C["子 Future.poll<br/>推进当前异步操作"]
A <-->|调用向下,结果向上| B
B <-->|Context 向下,Poll 结果向上| C图 1:调用关系,不是独立任务之间的消息流。子 Future 的 Ready 先交给父层处理,父层不一定立即返回 Ready;Pending 沿当前直接 await 的路径向外返回。
await 的作用就在这层父子关系中:父状态机推进子 Future,检查它返回的结果。
- 子 Future 返回 Ready:这份子计算已经完成,父状态机取出结果,执行 await 后面的代码。子层 Ready 不等于父层也已完成。
- 子 Future 返回 Pending:父状态机保留当前阶段与子 Future,也返回 Pending,退出本轮调用,等待调用者以后再次推进。
因此,调用方向自上而下,返回方向自下而上。子 Future 可以通过登记的 Waker 请求所属 task 再获执行机会,但不会借此直接执行父 Future 的后续代码。
一次 poll 不一定只跨过一个 await。如果响应和正文都已经可以取得,父状态机可以一路执行到 Ok(text),向调用者返回最终的 Ready。Ready 表示对应 Future 已完成,而不只是“现在能推进一点”。
四、响应没到时,Pending 和阻塞有什么不同?
这里的异步 HTTP 操作不会为了等网络一直占住这条工作线程。暂时拿不到响应时,等待路径安排好通知,请求子 Future 返回 Pending,外层 fetch_data 的 poll 也返回 Pending,把控制权交还给执行器。线程因此可以处理其他工作。
这与同步阻塞调用不同:同步调用没完成,线程仍停在调用里面;返回 Pending 则意味着本轮调用已经退出。async 也不能把任意阻塞函数变成这种行为——如果在函数体里直接调用同步阻塞接口,poll 仍可能被堵住,下面用两种等待方式对照。
同样等十秒,为什么一种卡线程,另一种可以让出?
use std::time::Duration;
async fn blocking_wait() {
std::thread::sleep(Duration::from_secs(10));
}
async fn async_wait() {
tokio::time::sleep(Duration::from_secs(10)).await;
}第一段被 poll 后,执行 sleep 的工作线程会停在同步调用里十秒。第二段使用 Tokio 定时器:期限未到时安排通知并返回 Pending,期限到达后再获得推进机会。区别不在函数有没有 async,而在实际调用的等待接口能否交还控制权。
在 current_thread 运行时里,一个同步阻塞就能拖住该运行时的异步任务;在多线程运行时里,其他 worker 可能继续工作,但被占住的 worker 无法及时返回执行器,阻塞工作过多也会耗尽执行资源。
必须使用同步接口时,把阻塞工作移出去
例如下面这段局部代码应放在能够处理 JoinError 和 I/O 错误的 async 函数中,path 是移入闭包的文件路径:
let bytes = tokio::task::spawn_blocking(move || std::fs::read(path))
.await??;同步读取在阻塞线程池执行,原异步任务等待返回的 JoinHandle,不必占着异步 worker。两层问号分别处理任务执行错误和文件读取错误。这是移动阻塞工作,不是消灭阻塞线程。
长计算也不能靠 async 自动拆开。大量 CPU 密集工作应限制并发,或交给专用计算线程池;Tokio 阻塞池的默认线程上限很大,不宜无限提交计算。短任务合理分段也可以,但主动让出不替代负载控制。
同步锁同样可能阻塞。优先缩短临界区,避免跨 await 持有同步锁;确实需要跨 await 持锁时再选择合适的异步锁,也仍要避免死锁和长时间占锁。
五、poll 已经返回,变量和读取进度保存在哪里?
调用帧退出了,但 HTTP 请求还没完成,后面还得继续。请求对象、已读数据以及“现在等到哪一步”,不能依赖这轮 poll 的栈帧保留。它们由仍然存在的 Future 及其子 Future 持有。
编译器因此分析:应该记住哪个阶段?哪些参数、局部变量和子 Future 需要跨暂停保留?
可以用下面这个 enum 理解生成结果。RequestFuture、BodyFuture 是说明性名称;这不是可直接编译的实现,也不承诺编译器采用同样的物理布局:
enum FetchData<'a> {
Init { url: &'a str },
WaitingForResponse { request: RequestFuture },
WaitingForText { body: BodyFuture },
Done,
}四个状态分别表示:
| 阶段 | 保存什么? | 再次推进时做什么? |
|---|---|---|
| Init | 传入的 url | 创建请求子 Future,开始推进 |
| WaitingForResponse | 尚未完成的请求子 Future | 继续等待 Response |
| WaitingForText | 尚未完成的正文子 Future | 继续收集正文 |
| Done | 已完成标记 | 不应再 poll 这个 async Future |
这里不需要保存所有局部变量。取得响应后,resp 被移入 resp.text() 返回的子 Future;等待正文时,由这个子 Future 持有响应及读取进度,父层不必再独立保留一份 resp。最终字符串通过 Ready(Ok(text)) 交出,不保留在 Done 中供反复领取。
让一个普通局部变量跨过 await
原版 fetch_data 主要把状态放进子 Future。为了直接看见普通局部变量的保存,可以让同一次请求额外返回状态码:
async fn fetch_with_status(
url: &str,
) -> Result<(reqwest::StatusCode, String), reqwest::Error> {
let resp = reqwest::get(url).await?;
let status = resp.status();
let text = resp.text().await?;
Ok((status, text))
}执行到第二个 await 时,status 已经取得,后面的返回值还要用它;text 则尚未取得。因此,这一阶段的逻辑状态是:
WaitingForText {
status: 已取得的状态码,
body: 正在读取正文的子 Future,
}下次 poll 时,继续推进 body。它返回 Ready 后,才有 text 可以与保存的 status 组成最终结果。不能把尚未取得的 text 画进等待状态,也不需要重新计算 status。
保存依据不只是“后面有没有读取”:还要考虑所有权、析构以及借用关系。一个之后不再读取、但仍需在退出作用域时析构的值,也可能继续存活。编译器再结合优化决定具体字段和布局,不能简单按源码变量数估算 Future 大小。
也不要理解成每次 Pending 都把整份线程栈复制到堆上。状态对象是跨多次 poll 持续存在的存储,生成的代码可以直接在其中建立和更新字段;只有本轮临时调用帧退出。
下面只看正常完成路径。Pending 时,保留当前阶段并退出本轮推进;错误则可以提前结束,不要求两个 await 都完成:
flowchart TB
A["Init:尚未开始<br/>保存传入的 url"]
B["WaitingForResponse<br/>保存请求子 Future"]
C["WaitingForText<br/>保存正文子 Future"]
D["Done:已完成<br/>结果已经交出"]
A -->|首次 poll 创建请求| B
B -->|取得 resp,移入正文 Future| C
C -->|读完正文,返回 Ready| D图 2:两处 await 对应两个可能暂停的阶段。请求、响应和读取进度按阶段交接,不靠保留整条未返回的函数调用链来恢复。
六、再次 poll 时,怎样从原来的阶段继续?
状态和数据已经留下了。运行时再次调用顶层 fetch_data Future 的 poll,生成的状态机先读取自己的阶段:如果在等待响应,就再次调用原请求子 Future 的 poll;如果在读取正文,就再次调用原正文子 Future 的 poll。子 Future 返回 Ready 后,父状态机才执行对应 await 后面的代码,不会因一次唤醒就跳过尚未完成的操作。若 fetch_data 本身嵌套在另一个 Future 中,这次推进则由它的父 Future 发起。
前面 Future trait 中的 Output,在这里是 Result<String, reqwest::Error>。self 指向保存的状态;Context 携带当前任务的 Waker,也就是稍后请求再次推进的通知入口。Pin 的地址稳定约束留到后面解释,先看状态如何改变。
生成的 poll 可以按下面的逻辑理解。它是控制流示意,不是可直接编译的 Rust:
poll(cx):
Init:
创建请求子 Future
进入 WaitingForResponse
WaitingForResponse:
调用保存的请求子 Future.poll(cx)
Pending:保存当前阶段,返回 Pending
Ready(Err):清理,标记 Done,返回 Ready(Err)
Ready(Ok(resp)):把 resp 移入正文子 Future,继续下一阶段
WaitingForText:
调用保存的正文子 Future.poll(cx)
Pending:保存当前阶段,返回 Pending
Ready(result):清理,标记 Done,返回 Ready(result)调用自父层向子层展开,Pending 则沿直接 await 的调用关系向外返回。之后仍由调用者再次发起 poll,父状态机按保存的阶段推进原来的子 Future,不会重新发起请求或从头收集正文。
无栈省掉的是哪份栈?
返回 Pending 是一次真正的函数返回:本轮 poll 的普通调用帧退出,工作线程不必停在网络调用里。还没完成的计算,由 Future 字段和其中的子 Future 继续持有。
把观察点停在“已经让出执行权,正在等待网络”的瞬间,两种方案的差别就清楚了:
| 此刻观察什么? | Go goroutine | Rust Future |
|---|---|---|
| 未完成的调用现场 | G 的专属调用栈仍保留 | 本轮 poll 的调用帧已退出 |
| 进度保存在哪里? | 调用栈与恢复上下文 | 阶段、字段和当前子 Future |
| 再次执行时做什么? | 恢复上下文,继续原调用 | 新调用 poll,按阶段进入对应路径 |
两者都能让线程去执行其他工作。无栈的区别在于:等待中的计算不必各自留下一份专属调用栈,执行时才临时使用工作线程的栈。 它仍需要保存状态、连接和正文缓冲,不能由此推断所有任务都更省内存或执行更快。 这类转换发生在 rustc 的 MIR coroutine lowering 阶段,之后再生成机器码。下面摘自 rustc 1.94.0、Windows x86-64、未优化的 fetch_data 编译结果,依赖 reqwest 0.12.28。编号与布局只解释本次输出,不是稳定 ABI。 入口读取状态并分派;0 对应尚未开始,3 对应请求等待,4 对应正文等待,1 和 2 分别用于完成与 panic 后不可继续的状态: 两处子 Future 暂时不能完成时,外层记录对应阶段,直接返回:编译结果补充:状态如何变成分派与返回
bb0: {
_34 = copy (_1.0: &mut {async fn body of fetch_data()});
_33 = discriminant((*_34));
switchInt(move _33) -> [0: bb1, 1: bb35, 2: bb34, 3: bb32, 4: bb33, otherwise: bb7];
}bb8: {
_0 = Poll::<Result<String, reqwest::Error>>::Pending;
discriminant((*_34)) = 3;
return;
}
bb19: {
_0 = Poll::<Result<String, reqwest::Error>>::Pending;
discriminant((*_34)) = 4;
return;
}
七、谁触发 Waker,Tokio 怎样再次安排执行?
现在已经知道,Future 能保存进度,poll 能按进度继续。还差一环:它不能自己知道什么时候再执行,也不应靠线程不停地询问网络。Tokio 与 Waker 负责把等待事件接回后续推进。
把 Future 交给运行时
在已有 Tokio 运行时中,可以这样提交这份计算:
let handle = tokio::spawn(fetch_data("https://example.com"));Future 保存计算状态;task 是执行器安排的工作单位;thread 真正执行 poll。Tokio 为 spawn 提交的 Future 建立 task,提供地址稳定的存储与调度状态。它不是每次 await 都创建一个新任务。
工作线程获得这个 task 后,用关联到它的 Waker 构造 Context,再调用顶层 fetch_data Future 的 poll。spawn 是提交工作,不是在调用点同步执行完函数体。
Waker 从哪里来,又保存在哪里?
假设 Pending 后执行器立刻循环重试,网络还没变化也不断调用 poll,这就变成忙等待。另一种做法是先停止推进,等条件可能改变时再收到通知。Waker 是后者需要的通知句柄。
对于 Tokio 的 task,运行时提供关联到该 task 的 Waker,并用它构造 Context。父 Future 在 await 处把 Context 传给子 Future,底层等待路径通过 cx.waker() 取得同一任务的通知入口。为了在本次 poll 返回后继续使用,等待路径通常克隆并保存它:
let waker = cx.waker().clone();这不会克隆 Future 或复制计算状态,得到的句柄仍关联同一个任务。Waker 可以跨线程传递和调用;但这不等于它所关联的 Future 也自动满足 Send,二者的约束不同。
Waker 封装了执行器定义的唤醒行为,底层 RawWaker 用数据指针和函数表描述相应操作,不能简单当成一个普通的 dyn Wake 对象。应用通常不必实现这些底层接口,只需保存、更新并调用运行时提供的句柄。
谁真正调用 wake?
事件源或负责推进底层工作的组件来调用,而不是父状态机不停检查。不同等待对应不同来源:
| 在等什么? | 谁安排后续通知? |
|---|---|
| socket 可读、可写 | 运行时 I/O 驱动与库的等待路径 |
| 定时器到期 | 定时器驱动 |
| channel 出现消息 | 发送端更新共享状态后通知接收者 |
| 阻塞任务结束 | 任务完成路径通知等待 JoinHandle 的异步任务 |
第一次推进 fetch_data 时,请求路径可能经历连接、TLS、发送与读取。暂时不能继续时,底层先安排通知,再逐层返回 Pending。reqwest 及 HTTP 库还可能经内部连接任务和通知渠道转接,不保证每个 socket 事件都直接唤醒 fetch_data 所属 task。
以 Tokio 1.53.2 的 socket 就绪等待为例,ScheduledIo::poll_readiness 保存当前 Waker,并在登记后重新检查状态,避免“检查时没数据,登记之前数据却到了”造成丢通知。不是每层 Future 都要各自保存一份 Waker,保存位置也不一定在 Future 内部。
为什么要先登记通知,再返回 Pending?
考虑下面的竞争:poll 检查到尚未完成,另一个线程立刻完成操作,但 poll 还没保存 Waker。完成线程找不到等待者;随后 poll 返回 Pending,执行器却再也收不到通知。
所以检查条件、登记等待者与完成通知必须有协调机制:可以在同一把短锁下检查并登记,也可以用原子状态和登记后的重新检查。Waker 本身不负责替你消除这场竞争。
下一次 poll 还可能传入不同的 Waker,等待实现应使用最新的通知入口。wake 请求的是再次推进,不是“操作肯定已经成功”:一次通知后仍可能 Pending,多次通知也可能被合并成一次 poll。
手写一个 Future,看清保存与唤醒的位置
下面的 Delay 只用于观察机制:构造时启动一个线程产生“时间到了”的事件,poll 只检查状态、更新 Waker,不会每次被调用都再开线程。实际定时等待应使用 tokio::time::sleep 的共享驱动,不应为每个等待创建 OS 线程。
use std::{
future::Future,
pin::Pin,
sync::{Arc, Mutex},
task::{Context, Poll, Waker},
thread,
time::Duration,
};
struct Delay {
shared: Arc<Mutex<DelayState>>,
}
struct DelayState {
done: bool,
waker: Option<Waker>,
}
impl Delay {
fn new(duration: Duration) -> Self {
let shared = Arc::new(Mutex::new(DelayState {
done: false,
waker: None,
}));
let worker = Arc::clone(&shared);
thread::spawn(move || {
thread::sleep(duration);
let waker = {
let mut state = worker.lock().unwrap();
state.done = true;
state.waker.take()
};
if let Some(waker) = waker {
waker.wake();
}
});
Self { shared }
}
}
impl Future for Delay {
type Output = ();
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> {
let mut state = self.shared.lock().unwrap();
if state.done {
Poll::Ready(())
} else {
state.waker = Some(cx.waker().clone());
Poll::Pending
}
}
}这个实现中,三个位置值得对应起来:
- poll 在同一把锁下检查 done 并登记 Waker。若事件先完成,直接 Ready;若 poll 先登记,完成线程一定能看到通知入口。
- 后台线程先写 done,再取出 Waker;释放锁之后才调用 wake,避免把外部通知逻辑放在共享状态锁里。
- 之后执行器再 poll,读取到 done 才返回 Ready。wake 自己没有代替这次检查。
这里用了只保护少量字段的同步短锁,没有跨 await 持锁;Mutex 也不是 Tokio 调度器的实现。这个教学版本不支持撤销后台线程,丢弃 Delay 仍要等该线程自行结束,不能把它当生产定时器使用。
现在再看这次 HTTP 请求的完整等待循环:Waker 提供执行机会,Future 状态决定从哪里继续,两者不是同一份信息。
flowchart TB
A["调度器取得可运行 task<br/>用 Waker 构造 Context"]
B["调用 fetch_data.poll<br/>推进请求或正文 Future"]
C["顶层返回 Pending<br/>本次调用栈退出"]
D["I/O 事件触发通知<br/>Waker 请求再次推进"]
A -->|提供执行机会| B
B -->|暂时不能完成| C
C -->|等待路径已登记通知| D
D -->|task 再次获得调度| A图 3:运行期的正常等待循环。编译器已经生成状态机,Tokio 只负责驱动。图中省略 HTTP 库内部连接任务与通知转接;事件也可能在 poll 进行期间到来,由运行时协调调度状态。
这就是“反复调用 poll”的含义:任务可运行时推进,暂时不能继续时等待通知,而不是不停轮询所有等待中的 Future。 wake 只是重试机会,不保证下一次一定 Ready,也不要求一次通知对应一次 poll。
沿一次响应走到 Ready
假设服务器先返回响应头,再分批发送正文:
| 发生什么? | fetch_data 怎样推进? |
|---|---|
| 请求尚未取得响应 | 请求子 Future 返回 Pending,外层保留 WaitingForResponse |
| 取得 Response | 请求子 Future 返回 Ready,把 resp 移入正文子 Future |
| 正文只到了一部分 | 正文子 Future 返回 Pending,外层保留 WaitingForText |
| 正文收集完成 | 子 Future 返回字符串,外层返回 Ready(Ok(text)) |
这不是固定四次 poll:一次 poll 可能跨过两个 await,一个 await 也可能多次 Pending;任何一步出错都可能提前返回 Ready(Err)。顶层 Ready 后,Tokio 记录任务完成,结果可以通过 handle.await 取得,不应再推进这个已完成的 async Future。
八、两份 Future,为什么不必创建两个独立任务?
理解一份 Future 的生命周期后,再看组合。在一个 async 函数或块中,可以同时推进两个请求:
let (result_a, result_b) = tokio::join!(
fetch_data(url_a),
fetch_data(url_b),
);join 把两份 Future 内联保存并在同一个 task 中推进,不需要分别 spawn,也不需要为两条分支各准备一份协程栈。A 暂时 Pending,join 仍可推进 B;两者完成后交出两个结果。这是并发,不是让同一条线程并行执行两段代码。
这里被唤醒的是包含 join 的外层 task。再次执行从外层 Future 进入组合逻辑,再推进分支,而不是调度器把每层子 Future 都独立入队。状态可以嵌套,调度任务不必跟着层层增加。
九、task 怎样在 worker 之间流转,为什么要求 Send?
前面解释了 Future 怎么等待,现在看运行时怎样分配执行机会。Tokio 多线程运行时不仅有一条任务队列,还包含异步 worker、I/O 与定时器驱动,以及执行同步工作的阻塞线程池。#[tokio::main] 负责创建并驱动运行时,不负责把 async 函数编译成状态机。
current_thread 运行时借用调用 Runtime::block_on 的线程执行异步任务;多线程调度器则有自己的 worker 线程池。I/O 驱动是一项职责,不应因此想象成所有配置下都存在一条独占的“I/O 线程”。
本地队列、全局队列与工作窃取
按 Tokio 1.53.2 文档描述,多线程调度器可以用下面几项理解:
| 结构 | 作用 |
|---|---|
| 每个 worker 的本地队列 | 保存该 worker 可取的任务,减少共享队列竞争 |
| 全局队列 | 承接外部注入的工作及部分队列转移 |
| 独立 LIFO slot | 对某些刚被当前任务唤醒的任务提供更快的后续推进 |
| 工作窃取 | 当前无可取工作时,从其他 worker 的本地队列转移一批任务 |
worker 优先处理本地工作,但也会按策略检查全局队列与 I/O、定时器事件,不能理解成“本地永远清空后才看全局”。本地与全局都没有工作时,它可以尝试窃取;仍无工作则等待通知,而不是无限空转。
LIFO slot 与普通本地队列是不同结构,不能一概说成“Tokio 的本地队列都是 LIFO,其他线程从另一端偷一个”。当前实现的窃取会转移一批任务,LIFO slot 本身不被窃取。具体调度细节可随版本变化,也不提供固定的任务执行顺序保证。
wake 可能发生在 task 的 poll 仍进行时。Tokio 通过任务状态协调通知与执行,避免两个 worker 同时 poll 同一个 task,也避免把执行中的通知直接丢掉。再调度不等于每次 wake 都立即追加一个独立队列元素。
Send 约束的是跨线程执行,不是地址是否固定
任务这次可能由 worker A 推进,返回 Pending;之后由 worker B 再次推进。跨暂停保留的状态也因此必须适合在线程间交接。tokio::spawn 的约束包括:
Future 本身:Future + Send + 'static
最终输出:Send + 'staticFuture 要满足 Send,是因为后续执行可能换线程;输出也要满足 Send,是因为结果可能被另一个线程上的 JoinHandle 等待者取得。并不是因为 Waker 可以跨线程,所有 Future 就天然可以跨线程。
即使当前使用 current_thread 运行时,tokio::spawn 的 Send 约束也不会消失。确实需要保留非 Send 状态,可以选择 LocalSet 和 spawn_local,或直接在当前任务中组合执行。
Rc 是否跨过 await,要看它是否真正结束生命周期
下面的片段故意不能通过 tokio::spawn 的 Send 检查:rc 在暂停之后还要使用,因而成为等待状态的一部分。
tokio::spawn(async {
let rc = std::rc::Rc::new(5);
tokio::task::yield_now().await;
println!("{rc}");
});如果只需要在等待之前使用它,用明确的词法作用域结束其生命周期:
tokio::spawn(async {
{
let rc = std::rc::Rc::new(5);
println!("{rc}");
}
tokio::task::yield_now().await;
});“最后一次读取在 await 前”不等于“一定已经析构”。这里只调整 println 的顺序而不结束作用域,不能作为普遍可靠的修复。若确实需要跨线程共享数据,可根据数据与同步需求选择 Arc 等类型;Arc 也不会自动把不适合跨线程的内部数据变得安全。
'static 不是让任务永远活着
它要求任务不能依赖可能先失效的非静态借用,不要求任务永久运行。任务可以持有自己的 String,也可以使用字符串字面值这样的静态借用。
如果请求地址来自一个局部 String,可以把所有权移进父 async 块:
let url = String::from("https://example.com");
let handle = tokio::spawn(async move {
fetch_data(&url).await
});这里 url 由提交的任务拥有,内部 fetch_data 借用它,不需要借用调用者即将退出的栈帧。async move 捕获已有变量的时机也在创建这个 async 块时,不能泛化成“创建 Future 什么都不做”。
必须跨 await 保留 Rc,就让它留在同一线程
下面的片段在已有 Tokio 运行时的异步上下文中使用 LocalSet。它主动等待本地任务结束,rc 可以跨 await 保留:
let local = tokio::task::LocalSet::new();
local.run_until(async {
tokio::task::spawn_local(async {
let rc = std::rc::Rc::new(5);
tokio::task::yield_now().await;
println!("{rc}");
}).await.unwrap();
}).await;spawn_local 不要求 Future 满足 Send,但仍有 'static 边界,也需要 LocalSet 或相应本地运行时环境。它不会因为名称里有 local 就允许任务任意借用一个很快退出的栈帧。
十、保存了内部借用,为什么还需要 Pin?
前面的 HTTP 例子把 Response 按所有权移交给子 Future。更一般的 async 函数也可能创建一个缓冲区,再把它的可变引用交给异步读取;暂停时,外层 Future 同时保存缓冲区和借用它的子 Future。
这种内部借用指向对象中的字段。若外层对象随意搬动,引用可能仍指向旧地址。poll 接收 Pin<&mut Self>,就是为需要固定的类型建立地址稳定约束,直到按约束结束生命周期;并非每份 async Future 都一定有这样的内部借用。
一个内部借用,为什么会怕整个对象搬家?
把这种关系缩小到几行代码:
async fn borrow_across_wait() -> usize {
let text = String::from("hello");
let text_ref = &text;
tokio::task::yield_now().await;
text_ref.len()
}暂停后 text_ref 还要使用,所以必须保留被借用的 text 及这个借用。逻辑关系是:
外层 Future
├─ text:String 值
├─ text_ref:指向上面这个 String
└─ 当前等待的子 FutureString 的字符数据在堆上,不代表这里的借用就与外层对象地址无关:text_ref 的类型是 &String,指向的是 String 值本身。若这个值随外层 Future 搬家,原引用指向的位置就可能失效。这解释的是允许出现的内部关系,不承诺优化后一定保留同样的物理字段。
Pin 的约束在建立固定关系时就生效,而不是第一次 poll 后才自动产生。对需要固定的 !Unpin 值,必须维持地址稳定直到按约束结束生命周期;Unpin 类型则允许相应移动。执行器通常先把 Future 放进稳定存储,再首次 poll。
堆分配也不自动等于固定:普通 Box 中的值仍可能被移出。Box::pin 或执行器遵守的固定约束才限制这些操作;移动一个指针、任务句柄,和移动它指向的 Future 值,是两件事。
Pin 不把任务固定到某条线程,也不要求每层单独 Box。父 Future 可以直接包含并固定子 Future,局部执行也可以栈上固定;Future 本身大小确定,也不代表它持有的正文、连接等堆数据大小固定。
因此,Send 回答“能否交给另一条线程继续执行”,Pin 回答“被固定的值能否搬动”。一个任务可以满足 Send,同时包含需要固定的 Future;这两个约束并不矛盾。
十一、完成、取消与回收,是不是同一件事?
正常结束或返回错误时,局部状态按所有权规则清理;任务容器、结果和 Future 值的存储由各自拥有者与执行器管理,不能都等同于“函数返回”。
同样是 drop,当前阶段不同,清理对象也不同
| 当前阶段 | 清理时要区分什么? |
|---|---|
| 尚未 poll | 原例捕获的是 url 引用,丢弃 Future 不代表释放调用者拥有的字符串 |
| 等待响应 | 清理当前请求子 Future 持有的状态;共享资源或库内部工作可能仍存活 |
| 正在读取正文 | 清理正文子 Future、响应及部分正文的持有关系 |
| 已经 Ready | 最终 String 已交给任务结果或调用者,不再属于等待阶段的局部状态 |
直接嵌套的子 Future 随父对象的清理路径一起析构。独立 spawn 的任务则有自己的生命周期:父任务只持有它的 JoinHandle,丢弃这个句柄不会自动取消已提交任务。这也是“组合成一个状态对象”和“拆成独立调度任务”的区别。
等待期间丢弃 Future,会按当前阶段清理它持有的子 Future 和资源,但已经发出的 HTTP 请求不会被撤回,服务端可能继续处理。底层连接是否关闭或复用、库内部工作如何收尾,要看库的实现,不能把 drop 一概等同于关闭所有关联活动。
对于 Tokio 的异步 task,可以用 handle.abort() 请求取消,再等待 handle.await 确认任务结束;取消与完成可能竞争,取消也不是同步回滚当前代码。若 poll 正被阻塞,清理不能绕过这段正在执行的代码。已经开始运行的 spawn_blocking 工作则不能靠 abort 强制停止,需要它自己的退出机制。
回看这次请求:编译器保存阶段与跨暂停数据,执行器或父 Future 发起 poll,等待路径安排 Waker 通知,后续调用按阶段继续。暂停时留下的是可恢复的计算对象,不是一份专属调用栈;完成与取消再按所有权清理对象和结果。
关于这种方案为什么契合 Rust,以及它与线程池、事件驱动和 Go 有栈协程的取舍,见 《从并发模型的演进看:为什么无栈协程是 Rust 的最优解?》。
参考资料
- await 的推进、暂停与恢复语义:https://doc.rust-lang.org/reference/expressions/await-expr.html
- Future 的 poll 契约:https://doc.rust-lang.org/std/future/trait.Future.html
- Waker 的来源、克隆与通知保证:https://doc.rust-lang.org/std/task/struct.Waker.html
- Pin、Unpin 与地址稳定约束:https://doc.rust-lang.org/std/pin/index.html
- rustc 1.94.0 的协程状态机转换:https://github.com/rust-lang/rust/blob/1.94.0/compiler/rustc_mir_transform/src/coroutine.rs
- reqwest 0.12.28 的 get 与 Client 复用:https://docs.rs/reqwest/0.12.28/reqwest/fn.get.html
- Response 所有权与正文读取:https://docs.rs/reqwest/0.12.28/reqwest/struct.Response.html#method.text
- Tokio 1.53.2 的运行时与多线程调度细节:https://docs.rs/tokio/1.53.2/tokio/runtime/index.html
- Tokio 任务、spawn、spawn_local 与取消:https://docs.rs/tokio/1.53.2/tokio/task/index.html
- 阻塞池、CPU 并发限制与取消边界:https://docs.rs/tokio/1.53.2/tokio/task/fn.spawn_blocking.html
- socket 等待状态与 Waker 登记:https://docs.rs/tokio/1.53.2/src/tokio/runtime/io/scheduled_io.rs.html
- 异步锁与同步锁的适用范围:https://docs.rs/tokio/1.53.2/tokio/sync/struct.Mutex.html