编程 从 Future 到工作窃取调度器:Rust 异步运行时源码级深挖与 Tokio 实战

2026-07-26 18:43:56 +0800 CST views 9

从 Future 到工作窃取调度器:Rust 异步运行时源码级深挖与 Tokio 实战

很多人写了两年 async/await,却说不清一个 .await 到底发生了什么。本文不讲概念背书,直接从状态机、WakerRuntimeReactor、工作窃取调度器一路拆到底,配可运行代码,最后给出一份真正能落地的性能优化清单。读完你对 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:被调用时,把这个任务重新丢回运行时的就绪队列

流程闭环是这样的:

  1. 运行时 poll 一个任务,传入包含 WakerContext
  2. 任务 poll 到最底层——比如一个 TCP 读——发现数据没到,返回 Pending。在返回前,它把 Waker 克隆一份,**注册到 IO 事件源(Reactor)**上。
  3. 运行时看到 Pending,把这个任务搁置,去跑别的任务。
  4. 网卡来数据了,操作系统的 epoll/kqueue/io_uring 唤醒 Reactor。
  5. Reactor 查表找到对应的 Waker,调用 wake()
  6. 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 线程持有三种队列:

  1. 本地队列(local queue):固定容量 256 的环形缓冲区(LIFO 附近有优化),无锁,只有本线程 push/pop。绝大多数 spawn 出来的任务先进这里。
  2. 全局队列(global/injection queue):加锁的共享队列。本地队列满了会溢出到这里;从运行时外部 spawn 的任务也进这里。
  3. 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 找到 Wakerwake() 之,对应任务重回调度。

新版 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 异步的心智模型其实很清晰:

  1. async fn 编译成状态机,每个 .await 是一个可暂停点,跨 await 的变量存进状态机——所以要 Send、要 Pin
  2. Future 靠 poll 被驱动,没好就返回 Pending 并注册 Waker
  3. Waker 是连接"任务"和"运行时就绪队列"的回调,IO 就绪时由 Reactor 触发,实现零忙等的事件驱动。
  4. Tokio 用工作窃取调度器 + mio Reactor,把单线程 poll 循环扩展成多核高并发引擎。
  5. 协作式调度是双刃剑:性能极致,但阻塞/CPU 密集代码会毒害整个 worker,必须 spawn_blocking 隔离。

展望未来,几个值得盯的方向:

  • io_uring 后端普及tokio-uring 成熟后,Linux 上的异步 IO 会再上一个台阶,syscall 批处理能显著降低高 QPS 场景的 CPU 占用。
  • async trait 稳定化async fn in trait 已在稳定 Rust 落地,异步生态的抽象能力大幅增强,动态分发的 dyn 异步 trait 也在推进。
  • 结构化并发TaskTrackerJoinSetCancellationToken 等模式让任务生命周期管理更安全,"spawn 出去就失控"的时代正在过去。

Rust 异步的学习曲线确实陡,但一旦你把状态机和 Waker 这两块拼图装进脑子,之前所有"玄学报错"都会变得可解释。工具会变、API 会升级,但 poll + Waker 这套底层协议是稳定的地基。把地基打牢,剩下的都是熟练度问题。

动手建议:别只看不写。把本文第三节那个 60 行 mini runtime 敲一遍跑起来,再往里加一个基于 mio 的 TCP echo,你对 Rust 异步的理解会彻底不一样。

推荐文章

虚拟DOM渲染器的内部机制
2024-11-19 06:49:23 +0800 CST
Graphene:一个无敌的 Python 库!
2024-11-19 04:32:49 +0800 CST
Go的父子类的简单使用
2024-11-18 14:56:32 +0800 CST
程序员茄子在线接单