异步任务可以理解成一段能够暂停和继续执行的代码。在 Rust 中任务的暂停和恢复,主要和 Future、poll、Waker 这些概念有关。
前言
很多语言都提供了异步编程的能力,本质上就是在“等结果”的这段时间里,不让线程白白挂住。像 C# 里有 Task 和 async/await,JavaScript 里有 Promise 和 async/await,目的都差不多,就是想让异步逻辑写起来不要那么蛋疼。
Rust 也有 async/await,但它背后的实现思路和其他语言不完全一样。比如下面这个异步函数:
1 | async fn hallo() { |
调用 hallo() 时,并不会立刻执行完整个函数,而是先返回一个 Future。
Future 是什么
Future 可以理解为一个尚未完成,但会在未来产出结果的异步计算。比如调用 hallo() 时,拿到的只是一个 Future:
1 | fn main() { |
虽然已经调用了 hallo(),但函数体里的 println!("hallo") 还没有执行。调用 async fn 时,Rust 先创建一个表示这段计算的 Future,而不是立刻运行函数体。
Future 可以看作一个状态机。它会保存当前执行到了哪里、后续还需要哪些局部变量,以及正在等待什么。
Future 本身不会在后台自动运行。runtime 负责驱动它执行,当 Future 遇到网络请求这类暂时拿不到结果的操作时,会先保存状态并暂停,等条件满足后再继续执行。
因此,Future 描述的是一段异步计算,异步任务则是 runtime 调度这段计算的单位。
Future 是怎么被执行的
Future 本身不会主动执行,runtime 会在合适的时候调用它的 poll 方法,推动异步计算继续往下走。
先来看一下 Future 特型和 Poll 枚举的定义:
1 | pub trait Future { |
Output 表示 Future 最终会产出的结果类型。poll 每次被调用时,会返回一个 Poll,告诉 runtime 这个 Future 目前的状态。
Poll::Ready(value):Future已经执行完成,value就是最终结果。Poll::Pending:现在还不能完成,需要之后再继续执行。
举个例子,给 hallo() 加一个 sleep:
1 | use tokio::time::{sleep, Duration}; |
runtime 第一次调用 hallo() 对应的 poll 时,代码会运行到 sleep(...).await。由于定时器还没有到期,hallo() 对应的 Future 会返回 Poll::Pending,runtime 转而执行其他任务。
3 秒后,runtime 再次调用它的 poll。代码会从 .await 之后继续执行,输出 hallo,最后返回 Poll::Ready(())。
不过还有一个问题:Future 返回 Poll::Pending 后,runtime 怎么知道它什么时候可以继续执行?这就是 Waker 的作用了。
Waker 是什么
从字面上看,Waker 就是一个“唤醒器”。当 Future 等待的条件满足后,定时器这类事件源会通过 Waker 通知 runtime:当前任务可以重新调度了。
前面例子里的 sleep 看起来没有直接使用 Waker。但实际上 Tokio 的 Sleep 在 poll 时会把 Context 传给内部定时器,关键逻辑如下:
1 | // tokio/time/sleep.rs |
Sleep::poll 把 cx 传给内部定时器。定时器会从中取出当前任务的 Waker,先把它注册起来,再读取当前状态:
1 | // tokio/runtime/time/entry.rs |
如果定时器还没有到期,read_state() 会返回 Poll::Pending。如果时间驱动器已经触发了这个定时器,此时返回 Poll::Ready(Ok(()))。
1 | // tokio/runtime/time/entry.rs |
等定时器到期后,时间驱动器会取出已注册的 Waker,并在释放内部锁后批量唤醒任务:
1 | // tokio/runtime/time/mod.rs |
wake_all() 不会直接执行 hallo() 后面的代码,它只是通知 runtime 将对应任务重新放回调度队列。等任务再次被 poll 时,sleep(...).await 才会完成,代码继续从 .await 之后往下执行。
相关源码:
实现一个简单的执行器
这个执行器一次只运行一个 Future,不支持 spawn,也不处理多个任务。入口函数通常命名为 block_on,因为它会阻塞当前线程,直到传入的 Future 返回最终结果。
先定义一个 Waker。它保存当前线程,当 Future 调用 wake() 时,就解除该线程的等待状态:
1 | use std::{ |
thread::park() 会让当前线程等待,thread::unpark() 则会唤醒它。这样,Future 暂时无法继续执行时,线程不需要反复调用 poll 轮询检查。
有了 Waker,block_on 的实现就比较直接了:
1 | fn block_on<F>(future: F) -> F::Output |
block_on 会循环调用 poll。
- 如果返回
Poll::Ready(value),直接返回结果。 - 如果返回
Poll::Pending,当前线程通过thread::park()进入等待。 - 等事件源准备好后,会调用之前注册的
Waker::wake(),最终通过thread::unpark()唤醒当前线程,再次调用poll。
这里注意 Box::pin(future) 是为了得到 Pin<&mut F>,它是 Future::poll 所需要的参数。暂时只需要知道 Pin 用于保证 Future 在内存中不被随意移动就够了。
最后
Future 的执行过程有点像函数调用栈。执行器从最外层 Future 开始调用 poll,遇到 .await 时,会继续调用内部 Future 的 poll。
如果内部 Future 返回 Poll::Pending,这个结果会向上传回执行器。等事件源通过 Waker 唤醒任务后,执行器再次从最外层调用 poll,任务再从上次暂停的位置继续执行。