从 Future 到工作窃取调度器:Rust 异步运行时源码级深挖与 Tokio 实战
很多人写了两年
async/await,却说不清一个.await到底发生了什么。本文不讲概念背书,直接从状态机、Waker、Runtime、Reactor、工作窃取调度器一路拆到底,配可运行代码,最后给出一份真正能落地的性能优化清单。读完你对 Rust 异步的心智模型会重建一次。
一、背景:Rust 的异步为什么长这样
写 Go 的人第一次看 Rust 异步会懵:Go 里 go func(){} 一行搞定,运行时全给你兜好了;Rust 却要你选运行时(Tokio、async-std、smol)、纠结 Send + 'static、被 Pin 折磨、还得理解 Waker。
这不是 Rust 故意为难人,而是它做了一个和别的语言完全不同的取舍:
- 零成本抽象:
async fn不分配堆、不引入 GC,编译成一个状态机结构体。你不用异步就不付任何代价。 - 运行时可插拔:语言只定义
Future这个"接口协议",具体怎么调度、怎么 IO 多路复用,全丢给第三方库。这是 Rust 没有官方运行时的根本原因。 - 协作式调度:Rust 的任务不会被强制抢占(对比 Go 1.14 后的异步抢占),一个任务不
.await就会一直霸占线程。这是双刃剑:极致性能,但也容易写出"饿死"别人的代码。
理解这三点,后面所有"反直觉"的设计都能说通。
二、核心概念:Future 到底是什么
2.1 Future trait 只有一个方法
标准库里的定义精简到令人发指:
pub trait Future {
type Output;
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>;
}
pub enum Poll<T> {
Ready(T),
Pending,
}
一句话概括:Future 是一个可以被反复"询问"的状态机。运行时调用 poll:
- 返回
Poll::Ready(v):我算完了,结果是v。 - 返回
Poll::Pending:我还没好,但我会在准备好时"通知你",别傻等,去干别的。
关键在"通知你"这个约定。谁来通知?靠 Context 里携带的 Waker。
2.2 async/await 只是语法糖
这段代码:
async fn read_two(a: &File, b: &File) -> io::Result<usize> {
let x = read_len(a).await?;
let y = read_len(b).await?;
Ok(x + y)
}
编译器会把它翻译成一个手写状态机,概念上等价于:
enum ReadTwo<'a> {
Start { a: &'a File, b: &'a File },
WaitingA { fut: ReadLen<'a>, b: &'a File },
WaitingB { x: usize, fut: ReadLen<'a> },
Done,
}
impl<'a> Future for ReadTwo<'a> {
type Output = io::Result<usize>;
fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
loop {
match &mut *self {
ReadTwo::Start { a, b } => {
let fut = read_len(a);
*self = ReadTwo::WaitingA { fut, b };
}
ReadTwo::WaitingA { fut, b } => {
// 把 poll 往下传递(真实代码这里要 Pin 投影)
match Pin::new(fut).poll(cx) {
Poll::Ready(Ok(x)) => {
let fut = read_len(b);
*self = ReadTwo::WaitingB { x, fut };
}
Poll::Ready(Err(e)) => return Poll::Ready(Err(e)),
Poll::Pending => return Poll::Pending, // 挂起,交还控制权
}
}
ReadTwo::WaitingB { x, fut } => {
match Pin::new(fut).poll(cx) {
Poll::Ready(Ok(y)) => {
let sum = *x + y;
*self = ReadTwo::Done;
return Poll::Ready(Ok(sum));
}
Poll::Ready(Err(e)) => return Poll::Ready(Err(e)),
Poll::Pending => return Poll::Pending,
}
}
ReadTwo::Done => panic!("polled after completion"),
}
}
}
}
看懂这个你就懂了一半:每个 .await 就是状态机的一个"暂停点"。函数里所有跨越 .await 的局部变量,都会被存进状态机结构体的对应变体里——这就是为什么异步函数的局部变量必须 Send 才能跨线程调度。
2.3 为什么需要 Pin
状态机里可能存在自引用:比如你 let buf = [0u8; 1024]; read(&mut buf).await,那个 ReadLen 里持有指向 buf 的指针,而 buf 也在同一个状态机结构体里。一旦这个结构体被 move(内存地址变了),内部指针就成了悬垂指针。
Pin<&mut T> 的作用就是:向类型系统承诺"这个值的内存地址在被 poll 期间不会移动"。这不是运行时检查,是编译期契约。绝大多数业务代码不用手写 Pin,交给 tokio::pin! 或 Box::pin 即可,但你得知道它为什么存在。
三、Waker 机制:整个异步的心脏
Poll::Pending 之后,运行时凭什么知道"什么时候再来 poll 一次"?答案是 Waker。
3.1 Waker 的本质
Waker 本质是一个类型擦除的回调句柄,它的底层是手写的虚表(vtable),因为标准库不能依赖任何具体运行时:
pub struct RawWakerVTable {
clone: unsafe fn(*const ()) -> RawWaker,
wake: unsafe fn(*const ()),
wake_by_ref: unsafe fn(*const ()),
drop: unsafe fn(*const ()),
}
data指针:通常指向"任务"本身(比如 Tokio 里的Task引用计数指针)。wake:被调用时,把这个任务重新丢回运行时的就绪队列。
流程闭环是这样的:
- 运行时 poll 一个任务,传入包含
Waker的Context。 - 任务 poll 到最底层——比如一个 TCP 读——发现数据没到,返回
Pending。在返回前,它把Waker克隆一份,**注册到 IO 事件源(Reactor)**上。 - 运行时看到
Pending,把这个任务搁置,去跑别的任务。 - 网卡来数据了,操作系统的
epoll/kqueue/io_uring唤醒 Reactor。 - Reactor 查表找到对应的
Waker,调用wake()。 wake()把任务塞回就绪队列,调度器下一轮再次 poll 它。这次读得到数据,返回Ready。
没有轮询、没有忙等,全靠事件驱动。这就是异步高并发的物理基础。
3.2 手写一个能跑的最小运行时
光说不练假把式。下面这个 60 行的运行时能真正驱动 Future 跑完,帮你把 Waker 彻底理解透:
use std::future::Future;
use std::pin::Pin;
use std::sync::{Arc, Mutex};
use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
use std::task::{Context, Poll, Wake, Waker};
/// 一个任务:持有 Future + 一个能把自己送回队列的发送端
struct Task {
future: Mutex<Pin<Box<dyn Future<Output = ()> + Send>>>,
sender: SyncSender<Arc<Task>>,
}
/// 实现 Wake:被唤醒时把自己 clone 一份塞回 channel
impl Wake for Task {
fn wake(self: Arc<Self>) {
let _ = self.sender.send(self.clone());
}
}
struct MiniRuntime {
ready_queue: Receiver<Arc<Task>>,
spawner: SyncSender<Arc<Task>>,
}
impl MiniRuntime {
fn new() -> Self {
let (spawner, ready_queue) = sync_channel(1024);
MiniRuntime { ready_queue, spawner }
}
fn spawn(&self, fut: impl Future<Output = ()> + Send + 'static) {
let task = Arc::new(Task {
future: Mutex::new(Box::pin(fut)),
sender: self.spawner.clone(),
});
let _ = self.spawner.send(task);
}
fn run(self) {
// spawner 必须 drop,否则 recv 永远阻塞
drop(self.spawner);
while let Ok(task) = self.ready_queue.recv() {
let waker = Waker::from(task.clone());
let mut cx = Context::from_waker(&waker);
let mut fut = task.future.lock().unwrap();
// poll 一次,Pending 就等下次 wake,Ready 就自然丢弃
let _ = fut.as_mut().poll(&mut cx);
}
}
}
fn main() {
let rt = MiniRuntime::new();
rt.spawn(async {
println!("hello from mini runtime");
});
rt.run();
}
这段代码点破了运行时的核心真相:运行时 = 一个就绪队列 + 一个 poll 循环 + Waker 把任务在两者之间搬运。Tokio 的复杂之处不在这个骨架,而在于把"单队列单线程"升级成"多线程工作窃取 + 高性能 Reactor"。
四、架构分析:Tokio 的多线程调度器
Tokio 的多线程运行时(multi_thread flavor)是它性能的灵魂。核心是工作窃取(work-stealing)调度器,思路借鉴自 Go runtime 和 Java ForkJoinPool。
4.1 三级队列结构
Tokio 每个 worker 线程持有三种队列:
- 本地队列(local queue):固定容量 256 的环形缓冲区(LIFO 附近有优化),无锁,只有本线程 push/pop。绝大多数 spawn 出来的任务先进这里。
- 全局队列(global/injection queue):加锁的共享队列。本地队列满了会溢出到这里;从运行时外部 spawn 的任务也进这里。
- LIFO slot:每个 worker 有一个单槽"最近任务"缓存。刚被 wake 的任务优先放这里,利用缓存局部性(刚唤醒的任务数据大概率还在 CPU cache 里)。
4.2 worker 的调度循环(简化逻辑)
loop {
// 1. 优先跑 LIFO slot(缓存热)
if let Some(task) = self.lifo_slot.take() { run(task); continue; }
// 2. 从本地队列取
if let Some(task) = self.local_queue.pop() { run(task); continue; }
// 3. 周期性检查全局队列(防止全局队列饿死)
if tick % GLOBAL_INTERVAL == 0 {
if let Some(task) = global_queue.pop() { run(task); continue; }
}
// 4. 本地空了,去别的 worker 那"偷"一半任务
if let Some(task) = self.steal_from_others() { run(task); continue; }
// 5. 实在没活干,注册到 Reactor,park 线程休眠
self.park();
}
几个魔鬼细节:
- 为什么偷"一半":一次偷一半能减少偷取频率(偷取要 CAS 竞争),也让负载更均衡。这是经典 work-stealing 的最优策略。
- 为什么周期性查全局队列:如果 worker 只顾着跑本地队列,全局队列的任务可能永远轮不上,导致公平性问题。Tokio 用
tick计数强制"雨露均沾"。 - LIFO slot 的抢占保护:为防止两个任务互相 wake 对方形成"乒乓"独占一个 worker,Tokio 给 LIFO slot 加了预算限制,超了就踢回本地队列尾部。
4.3 Reactor:IO 事件的中枢
调度器负责"跑任务",Reactor(基于 mio 库)负责"等 IO"。Reactor 内部就是对 epoll(Linux)/kqueue(macOS/BSD)/IOCP(Windows) 的封装:
- 每个异步 IO 资源(TcpStream 等)注册时,会拿到一个 token,并把自己的
Waker存进一张表。 - 一个专门的线程(或复用 worker)调用
epoll_wait阻塞等待。 - 事件到达,根据 token 找到
Waker,wake()之,对应任务重回调度。
新版 Tokio 也在实验 io_uring 后端(tokio-uring),能把"提交 IO 请求"和"收割完成事件"批处理化,进一步降低 syscall 开销,这是未来高性能 IO 的方向。
五、代码实战:从入门到踩坑
5.1 并发不等于并行:join 与 spawn 的区别
新手最大的误区:以为写了 async 就是并发跑。看这段:
// ❌ 串行!total ≈ 2 秒
async fn serial() {
slow_task(1).await;
slow_task(2).await;
}
// ✅ 并发!在同一个任务内并发,total ≈ 1 秒
async fn concurrent() {
tokio::join!(slow_task(1), slow_task(2));
}
// ✅ 并行!分发到线程池,可跨多核
async fn parallel() {
let h1 = tokio::spawn(slow_task(1));
let h2 = tokio::spawn(slow_task(2));
let _ = tokio::join!(h1, h2);
}
.await依次写 = 串行,第一个不完成第二个不开始。join!= 并发,在单个任务内交替 poll 多个 Future,适合 IO 密集但不需要多核。spawn= 并行,任务被丢进调度器可被不同 worker 抢,能吃满多核。
判断依据:如果子任务之间没有共享可变状态、想吃满 CPU,用 spawn;只是想同时等多个 IO、逻辑上仍是"一个请求",用 join!。
5.2 致命陷阱:在异步里跑 CPU 密集 / 阻塞调用
Tokio 是协作式调度,一个任务不 .await 就不会交还线程。下面是生产事故重灾区:
// ❌ 灾难:阻塞了整个 worker 线程,其他任务全部饿死
async fn bad_handler() {
let data = std::fs::read_to_string("huge.log").unwrap(); // 同步阻塞 IO
let hash = expensive_hash(&data); // CPU 密集
save(hash).await;
}
在默认的多线程运行时里,如果 worker 数量有限(默认 = CPU 核数),几个这样的 handler 就能让整个服务卡死。正确姿势:
async fn good_handler() {
// 阻塞 IO / CPU 密集统统丢到专用阻塞线程池
let data = tokio::task::spawn_blocking(|| {
std::fs::read_to_string("huge.log")
}).await.unwrap().unwrap();
let hash = tokio::task::spawn_blocking(move || expensive_hash(&data))
.await
.unwrap();
save(hash).await;
}
spawn_blocking 会把闭包丢到一个独立的、可动态扩容的阻塞线程池(默认上限 512),不占用核心异步 worker。记住一条铁律:核心异步线程上,两次 .await 之间的代码不应该超过 ~10-100μs。
5.3 用 select! 实现超时与取消
use tokio::time::{timeout, Duration};
async fn with_deadline() -> Result<String, &'static str> {
match timeout(Duration::from_secs(3), fetch_remote()).await {
Ok(Ok(body)) => Ok(body),
Ok(Err(_)) => Err("fetch failed"),
Err(_) => Err("timed out"), // 3 秒没完成,fetch_remote 被 drop 取消
}
}
Rust 异步取消是**"drop 即取消"**:timeout 超时后直接丢弃里面的 Future,状态机停在某个 .await 点不再被 poll。这带来一个隐藏陷阱——取消安全性(cancellation safety):
// ⚠️ 取消不安全的例子
loop {
tokio::select! {
// 如果 read_exact 读了一半被取消,已读进 buf 的字节就丢了
res = stream.read_exact(&mut buf) => handle(res),
_ = shutdown.recv() => break,
}
}
read_exact 不是取消安全的:被 select 的另一分支取消时,它可能已经消费了一部分数据但没返回。要用取消安全的 API(如 tokio 的 read_buf),或把状态存在 select 外部。这是写健壮异步服务的必修课。
5.4 一个完整的并发限流爬虫
把上面的知识串起来,写个真正能跑的东西——用信号量限流的并发抓取:
use std::sync::Arc;
use tokio::sync::Semaphore;
async fn crawl_all(urls: Vec<String>, max_concurrency: usize) -> Vec<usize> {
let sem = Arc::new(Semaphore::new(max_concurrency));
let mut handles = Vec::with_capacity(urls.len());
for url in urls {
let sem = sem.clone();
handles.push(tokio::spawn(async move {
// 拿到许可才能继续,超过并发上限的任务在此挂起(不是忙等)
let _permit = sem.acquire().await.unwrap();
match fetch_len(&url).await {
Ok(len) => len,
Err(_) => 0,
}
// _permit 在此 drop,自动归还许可
}));
}
let mut results = Vec::with_capacity(handles.len());
for h in handles {
results.push(h.await.unwrap_or(0));
}
results
}
Semaphore::acquire().await 在没许可时返回 Pending 并挂起,有许可释放时通过 Waker 唤醒——又是 Waker 机制在底层默默工作。这就是异步限流优雅之处:一万个任务全 spawn 出去,实际并发被信号量精确控制,没有一个线程在忙等。
六、性能优化:把运行时榨干
6.1 减少 spawn 数量与任务体积
每个 spawn 都有分配 + 调度开销,任务的状态机越大,move 成本越高。优化手段:
- 批处理:与其为每条消息 spawn 一个任务,不如用一个任务循环 + channel 批量消费。
- Box 大 Future:如果一个
async块状态机特别大(比如包含大数组),spawn时会整体拷贝。用Box::pin把它放堆上,减少栈上 move 成本。编译器还会警告large_futures。
6.2 选对锁:别在 .await 期间持有 std::Mutex
// ❌ 持有 std Mutex 跨越 .await:可能死锁 + 阻塞 worker
let mut guard = std_mutex.lock().unwrap();
do_async(&mut *guard).await; // guard 跨 await 存活
// ✅ 方案 A:缩小临界区,await 前释放锁
let value = { std_mutex.lock().unwrap().clone() };
do_async(&value).await;
// ✅ 方案 B:确实要跨 await 持锁,用 tokio::sync::Mutex
let mut guard = tokio_mutex.lock().await;
do_async(&mut *guard).await;
原则:能不跨 .await 持锁就不持;短临界区、无 await 就用 std::sync::Mutex(更快);必须跨 await 才用 tokio::sync::Mutex。别无脑全用异步锁,它比 std 锁慢。
6.3 channel 选型
tokio::sync::mpsc:有界 channel,带背压(发送方满了会 await 挂起),生产环境首选。tokio::sync::broadcast:一发多收,注意慢消费者会导致消息丢失(Lagged)。flume/crossbeam:同步场景或需要 select 时更灵活。
永远优先有界 channel。无界 channel(unbounded_channel)在生产者快于消费者时会无限堆积内存,是线上 OOM 的经典原因。
6.4 运行时调参
let rt = tokio::runtime::Builder::new_multi_thread()
.worker_threads(num_cpus::get()) // 通常等于物理核数
.max_blocking_threads(512) // 阻塞池上限
.thread_stack_size(2 * 1024 * 1024) // 深递归 Future 可调大
.enable_all()
.build()
.unwrap();
- IO 密集型服务:worker 数不必超过核数,多了只增加上下文切换。
- 有大量
spawn_blocking:适当调大max_blocking_threads。 - 单核容器(如 K8s 里 limit=1 core):直接用
new_current_thread()单线程运行时,省掉工作窃取的同步开销,反而更快。
6.5 用 tokio-console 定位问题
抽象的性能问题最怕拍脑袋。tokio-console 能实时看到每个任务的 poll 次数、poll 耗时、被 wake 的频率、是否长时间占用 worker。发现某个任务 busy 时间异常高,基本就是它在同步阻塞——这是异步性能排查的第一工具,务必上。
七、总结与展望
回顾整条链路,Rust 异步的心智模型其实很清晰:
async fn编译成状态机,每个.await是一个可暂停点,跨 await 的变量存进状态机——所以要Send、要Pin。- Future 靠
poll被驱动,没好就返回Pending并注册Waker。 Waker是连接"任务"和"运行时就绪队列"的回调,IO 就绪时由 Reactor 触发,实现零忙等的事件驱动。- Tokio 用工作窃取调度器 + mio Reactor,把单线程 poll 循环扩展成多核高并发引擎。
- 协作式调度是双刃剑:性能极致,但阻塞/CPU 密集代码会毒害整个 worker,必须
spawn_blocking隔离。
展望未来,几个值得盯的方向:
io_uring后端普及:tokio-uring成熟后,Linux 上的异步 IO 会再上一个台阶,syscall 批处理能显著降低高 QPS 场景的 CPU 占用。- async trait 稳定化:
async fn in trait已在稳定 Rust 落地,异步生态的抽象能力大幅增强,动态分发的dyn异步 trait 也在推进。 - 结构化并发:
TaskTracker、JoinSet、CancellationToken等模式让任务生命周期管理更安全,"spawn 出去就失控"的时代正在过去。
Rust 异步的学习曲线确实陡,但一旦你把状态机和 Waker 这两块拼图装进脑子,之前所有"玄学报错"都会变得可解释。工具会变、API 会升级,但 poll + Waker 这套底层协议是稳定的地基。把地基打牢,剩下的都是熟练度问题。
动手建议:别只看不写。把本文第三节那个 60 行 mini runtime 敲一遍跑起来,再往里加一个基于
mio的 TCP echo,你对 Rust 异步的理解会彻底不一样。