编程 从零手写一个 Rust 异步运行时:吃透 Future、Waker 与 Reactor,造一个能跑 TCP 服务的 mini-tokio(附完整实现)

2026-08-01 05:15:47 +0800 CST views 13

一、背景:为什么 Rust 的 async 让人又爱又恨

写过 Go 的人第一次碰 Rust 的 async/await,几乎都会经历同一段心路历程:

「不就是加个 async 关键字、await 一下吗?跟 goroutine 有啥区别?」

结果一编译,the trait Future is not implementedPin<Box<dyn Future>>cannot borrow across await pointfuture cannot be sent between threads safely 一大堆报错糊脸而来。

原因很简单:Go 的并发模型是「语言 + 运行时」深度绑死的,而 Rust 的异步是「语言只给了一套协议,运行时你自己选」。

Go 里的 goroutine 是有栈协程(stackful coroutine),调度器、netpoller、GC 全部内置在 runtime 里,你 go func(){} 一下就完事了,看不见底层。Rust 反过来——它把 async/await 编译成一个无栈状态机(stackless coroutine),语言标准库里只定义了 FuturePollWaker 这几个「接口」,至于「谁来驱动这些状态机往前跑」「I/O 事件怎么通知」「任务在哪个线程上跑」,标准库一概不管,全部甩给第三方运行时(tokio、async-std、smol、glommio……)。

这套设计的好处是「零成本抽象 + 运行时可插拔」:嵌入式设备可以用 embassy,高吞吐服务用 tokio,追求 thread-per-core 用 glommio。坏处就是——对使用者来说,运行时是个黑盒。你天天 #[tokio::main],但 .await 到底发生了什么?一个 TCP 读操作是怎么「挂起」又被「唤醒」的?Waker 这个从天而降的东西是干嘛的?Pin 为什么阴魂不散?

2026 年的 Rust 异步生态已经高度成熟,tokio 几乎是事实标准,各种教程教你怎么「用」它。但真正吃透 async 的唯一办法,是自己动手把运行时造一遍。本文就干这件事:从 Future 协议的第一性原理出发,用几百行代码手写一个能跑真实 TCP echo 服务的 mini-tokio,把 Executor(执行器)、Reactor(反应堆)、Waker(唤醒器)这三件套的每一根线都接明白。

读完你会得到:

  • 彻底看懂 async fn 被编译成了什么,.await 在展开后长什么样;
  • 亲手实现 RawWaker 的 vtable,理解 Waker::wake() 到底唤醒了谁;
  • 用 epoll(通过 mio)写一个 Reactor,把「I/O 就绪」翻译成「唤醒某个任务」;
  • 组合出单线程和多线程 work-stealing 两个版本的执行器;
  • 一份 async 领域最容易踩的坑清单(Pin、忘注册 Waker、伪唤醒、Send/Sync)。

代码基于 Rust 1.85(2024 edition),全部可编译运行。


二、核心概念:把 async/await 拆到底层

2.1 Future 就是一个「可以被反复问『你好了没』的状态机」

一切的起点是 std::future::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 问它:「好了没?」它要么返回 Ready(结果),要么返回 Pending(意思是「没好,别傻等,等我好了会 wake() 你」)。

关键点在于:Future 本身是惰性的(lazy)。你写 let fut = async { ... }; 这一行,里面的代码一行都不会执行。只有当某个东西(执行器)不断地调用 fut.poll(),它才会往前推进。这跟 JavaScript 的 Promise(创建即执行,eager)是根本区别。

2.2 async fn 被编译成了什么?状态机 desugar

理解 async 最大的一道坎,是意识到 async fn 只是语法糖。编译器会把它翻译成一个实现了 Future匿名状态机结构体

看这段代码:

async fn read_two(a: &TcpStream, b: &TcpStream) -> usize {
    let x = read_len(a).await;   // 暂停点 1
    let y = read_len(b).await;   // 暂停点 2
    x + y
}

编译器大致会把它变成这样一个枚举状态机(伪代码,帮助理解):

enum ReadTwoState<'s> {
    Start { a: &'s TcpStream, b: &'s TcpStream },
    WaitingA { fut_a: ReadLenFut<'s>, b: &'s TcpStream },
    WaitingB { fut_b: ReadLenFut<'s>, x: usize },
    Done,
}

impl<'s> Future for ReadTwoState<'s> {
    type Output = usize;
    fn poll(mut self: Pin<&mut Self>, cx: &mut Context) -> Poll<usize> {
        loop {
            match &mut *self {
                ReadTwoState::Start { a, b } => {
                    // 进入 await 1:创建子 Future,切到 WaitingA
                    let fut_a = read_len(a);
                    *self = ReadTwoState::WaitingA { fut_a, b };
                }
                ReadTwoState::WaitingA { fut_a, b } => {
                    match Pin::new(fut_a).poll(cx) {
                        Poll::Ready(x) => {
                            let fut_b = read_len(b);
                            *self = ReadTwoState::WaitingB { fut_b, x };
                        }
                        Poll::Pending => return Poll::Pending, // 把 Pending 冒泡上去
                    }
                }
                ReadTwoState::WaitingB { fut_b, x } => {
                    match Pin::new(fut_b).poll(cx) {
                        Poll::Ready(y) => {
                            let out = *x + y;
                            *self = ReadTwoState::Done;
                            return Poll::Ready(out);
                        }
                        Poll::Pending => return Poll::Pending,
                    }
                }
                ReadTwoState::Done => panic!("polled after completion"),
            }
        }
    }
}

三个洞察一次性打通:

  1. .await = 「poll 子 Future,如果 Pending 就把 Pending 往上冒泡(连带保存现场,下次从这个状态继续)」。这就是「暂停」的真相——不是线程阻塞,而是函数直接 return 出去,把「我执行到哪了」记录在状态机的 enum tag 里。

  2. 状态机是自引用的(self-referential)。注意 WaitingA 里存的 fut_a 可能借用了状态机自己持有的数据(比如局部变量的引用跨越了 await 点)。一旦这个结构体在内存里被移动(move),内部的自引用指针就会变成野指针——这就是 Pin 存在的唯一理由:它是一个类型级别的承诺「这块内存不会再被移动」,让自引用状态机可以安全存在。

  3. Poll::Pending 会一路冒泡到最顶层,最终被执行器接住。执行器一看返回 Pending,就知道「这个任务卡住了,先放一边」。

2.3 Waker:Pending 之后,谁来把我叫醒?

如果只有 poll,执行器面临一个尴尬问题:任务返回了 Pending,那我要不要再 poll 它一次?什么时候 poll?

朴素做法是「忙轮询」(busy-poll):不停地循环 poll 所有任务。 这会把 CPU 烧到 100%,因为绝大多数时候任务根本没就绪,你 poll 一万次可能才有一次是 Ready。这显然不能忍。

Waker 就是来解决这个问题的。看 poll 的签名,它接收一个 &mut Context,而 Context 里最重要的东西就是一个 Waker

impl Context<'_> {
    pub fn waker(&self) -> &Waker { ... }
}

impl Waker {
    pub fn wake(self) { ... }        // 消费掉自己,触发唤醒
    pub fn wake_by_ref(&self) { ... } // 不消费,触发唤醒
    pub fn clone(&self) -> Waker { ... }
}

约定是这样的:

当一个 Future 返回 Pending 时,它必须先把 cx 里的 Waker 克隆一份存起来(通常交给底层的 I/O 事件源或定时器)。等到那个「让它 Pending 的条件」满足了(比如 socket 可读了、定时器到点了),底层就调用 waker.wake(),通知执行器:「喂,这个任务可以再 poll 了!」

于是整条链路闭环了:

执行器 poll 任务
      │
      ├─ Ready(v)  → 任务完成,丢弃
      │
      └─ Pending   → 任务把 Waker 交给 Reactor 保管,执行器把任务放回「休眠区」
                        │
                     (某个时刻 I/O 就绪)
                        │
                     Reactor 调用 waker.wake()
                        │
                     任务被重新放进「就绪队列」
                        │
                     执行器再次 poll 它

这就是异步运行时的核心飞轮。整个运行时无非就是在实现这个飞轮的三个零件:

  • Executor(执行器):维护就绪队列,不停从队列取任务 poll
  • Reactor(反应堆):基于 OS 的 epoll/kqueue/IOCP,监听所有 I/O fd,就绪时调用对应的 Waker
  • Waker(唤醒器):连接前两者的桥梁——wake() 的动作,本质就是「把任务塞回执行器的就绪队列」。

想明白这三者的关系,你就已经理解了 tokio 的骨架。下面开始动手造。


三、架构总览:三件套如何咬合

先把我们要造的东西画清楚。mini-tokio 的数据流:

                    ┌─────────────────────────────────────────┐
                    │              Executor 执行器              │
                    │  ┌────────────────────────────────────┐  │
   spawn(future) ──▶│  │  ready_queue: 就绪任务的 MPSC 通道   │  │
                    │  └────────────────────────────────────┘  │
                    │        │ recv                             │
                    │        ▼                                  │
                    │   task.poll(cx)   cx 里带着 task 自己的   │
                    │        │           Waker(wake=把自己塞  │
                    │        │           回 ready_queue)        │
                    │   ┌────┴─────┐                            │
                    │  Ready      Pending                       │
                    │  (丢弃)      (等 Reactor 唤醒)             │
                    └───────────────────┬──────────────────────┘
                                        │ 注册 fd + waker
                                        ▼
                    ┌─────────────────────────────────────────┐
                    │              Reactor 反应堆               │
                    │   mio::Poll  ──  epoll_wait(events)       │
                    │   事件就绪 → 取出对应 waker → wake()      │
                    └─────────────────────────────────────────┘

我们分五步实现,由浅入深:

  1. 手写一个最小 Waker(不依赖任何库,纯 RawWaker vtable);
  2. 手写单线程执行器 block_on(能跑纯计算型 Future);
  3. 引入 spawn + 就绪队列;
  4. mio 写 Reactor,实现异步 TcpStream
  5. 组合出 echo server;再给一个多线程 work-stealing 版本。

四、代码实战

4.1 第一步:手写一个 Waker(vtable 硬核版)

很多教程直接用 futures 库的 waker_fnArcWake,那样看不到本质。我们从最底层的 RawWaker 开始。

Waker 内部其实是一个 RawWaker,而 RawWaker = 一个 *const () 数据指针 + 一个 &'static RawWakerVTable 函数表。vtable 有四个函数指针:clonewakewake_by_refdrop。这套设计等价于「手写的 dyn 虚表」,之所以不直接用 dyn Trait,是为了让 Waker 能在 no_std / FFI 场景下也能用裸指针玩转。

我们让 waker 内部的数据指针指向一个 Arc<T>T 是一个能「把任务重新排进就绪队列」的东西。先写一个通用的、基于 Arc 的 waker 构造器:

use std::sync::Arc;
use std::task::{RawWaker, RawWakerVTable, Waker};

/// 任何实现了 Wake 的类型,都能被包成一个标准库 Waker
pub trait ArcWake: Send + Sync + Sized {
    /// wake 的动作:通常是把自己(一个任务)塞回就绪队列
    fn wake_by_arc(self: Arc<Self>);
}

/// 把 Arc<W> 转成一个 std 的 Waker
pub fn waker_from_arc<W: ArcWake + 'static>(w: Arc<W>) -> Waker {
    // 把 Arc 变成裸指针塞进 RawWaker 的 data 字段
    let ptr = Arc::into_raw(w) as *const ();
    let raw = RawWaker::new(ptr, vtable::<W>());
    // SAFETY: 我们保证 vtable 的四个函数都按 Arc 语义正确实现
    unsafe { Waker::from_raw(raw) }
}

fn vtable<W: ArcWake + 'static>() -> &'static RawWakerVTable {
    &RawWakerVTable::new(
        clone::<W>,          // 克隆:Arc 引用计数 +1
        wake::<W>,           // wake:消费 Arc 并触发唤醒
        wake_by_ref::<W>,    // wake_by_ref:不消费 Arc,触发唤醒
        drop_arc::<W>,       // drop:Arc 引用计数 -1
    )
}

unsafe fn clone<W: ArcWake + 'static>(ptr: *const ()) -> RawWaker {
    // 从裸指针恢复 Arc,clone 一份(计数 +1),再把两份都 forget 回裸指针状态
    let arc = Arc::from_raw(ptr as *const W);
    let cloned = arc.clone();                 // 引用计数 +1
    let _ = Arc::into_raw(arc);               // 原来的那份别 drop,还回裸指针
    RawWaker::new(Arc::into_raw(cloned) as *const (), vtable::<W>())
}

unsafe fn wake<W: ArcWake + 'static>(ptr: *const ()) {
    // 消费语义:直接把 Arc 拿回来(计数不变,所有权转移),调用 wake_by_arc
    let arc = Arc::from_raw(ptr as *const W);
    ArcWake::wake_by_arc(arc); // arc 在这里被消费掉(drop 或转移进队列)
}

unsafe fn wake_by_ref<W: ArcWake + 'static>(ptr: *const ()) {
    // 借用语义:不能消费掉这份 Arc,所以先 clone 出一份来 wake
    let arc = Arc::from_raw(ptr as *const W);
    let cloned = arc.clone();
    let _ = Arc::into_raw(arc);   // 原份还回裸指针,别 drop
    ArcWake::wake_by_arc(cloned); // 用 clone 出来的那份去 wake
}

unsafe fn drop_arc<W: ArcWake + 'static>(ptr: *const ()) {
    // 把裸指针变回 Arc 然后正常 drop,引用计数 -1
    drop(Arc::from_raw(ptr as *const W));
}

这段代码是整个运行时里唯一需要 unsafe 的地方,也是最容易写崩的地方。注意每个函数对 Arc 引用计数的处理:

  • clone:必须 +1(多了一份 Waker 在外面),同时保证原指针还活着 → 所以要 into_raw 两份;
  • wake消费语义,直接把 Arc 的所有权拿回来(不 +1 也不 forget),用完自然 drop,计数净 -1,和当初 into_raw 时的 +1 抵消;
  • wake_by_ref借用语义,不能动原来那份的所有权,所以 clone 一份去 wake,原份 into_raw 还回去;
  • drop:单纯 -1。

引用计数只要错一个,就是「double free」或「内存泄漏」二选一。这也是为什么实际项目里大家宁愿用 futures::task::ArcWake 派生,也不想手写这段——但手写一遍,你才真正懂 Waker 是什么。

小结:Waker 本质就是「一个类型擦除的 Arc<某个能把任务重新入队的东西>」。wake() 就是「触发那个入队动作」。

4.2 第二步:最小执行器 block_on

有了 waker,我们可以写一个最简单的执行器:block_on —— 在当前线程上把一个 Future 跑到完成。它的逻辑是「poll → 如果 Pending 就用一个 Condvar 睡觉等 wake → 被 wake 后再 poll」。

use std::future::Future;
use std::pin::pin;
use std::sync::{Arc, Condvar, Mutex};
use std::task::{Context, Poll};

/// 一个最简单的 waker:wake 时把 signaled 置 true 并通知 Condvar
struct BlockWaker {
    mutex: Mutex<bool>,
    condvar: Condvar,
}

impl ArcWake for BlockWaker {
    fn wake_by_arc(self: Arc<Self>) {
        let mut signaled = self.mutex.lock().unwrap();
        *signaled = true;
        self.condvar.notify_one();
    }
}

pub fn block_on<F: Future>(future: F) -> F::Output {
    // future 必须被 Pin 住才能 poll。pin! 宏把它钉在当前栈帧上。
    let mut future = pin!(future);

    let block_waker = Arc::new(BlockWaker {
        mutex: Mutex::new(false),
        condvar: Condvar::new(),
    });
    let waker = waker_from_arc(block_waker.clone());
    let mut cx = Context::from_waker(&waker);

    loop {
        match future.as_mut().poll(&mut cx) {
            Poll::Ready(val) => return val,
            Poll::Pending => {
                // 睡觉,直到有人 wake 我们(把 signaled 设为 true)
                let mut signaled = block_waker.mutex.lock().unwrap();
                while !*signaled {
                    signaled = block_waker.condvar.wait(signaled).unwrap();
                }
                *signaled = false; // 重置,准备下一轮
            }
        }
    }
}

这个 block_on 已经能跑真实的异步代码了。测一个自定义 Future——一个「被 poll 第 N 次才 Ready」的 Future:

struct Countdown(u32);

impl Future for Countdown {
    type Output = &'static str;
    fn poll(mut self: std::pin::Pin<&mut Self>, cx: &mut Context) -> Poll<&'static str> {
        if self.0 == 0 {
            Poll::Ready("发射!")
        } else {
            self.0 -= 1;
            // 关键:返回 Pending 前必须安排 wake,否则永远不会被再次 poll
            cx.waker().wake_by_ref();
            Poll::Pending
        }
    }
}

fn main() {
    let r = block_on(Countdown(3));
    println!("{r}"); // 发射!
}

注意 Countdown 里那行 cx.waker().wake_by_ref() —— 这是新手最容易漏、也是最致命的一行。如果你返回 Pending 却不安排任何 wake,执行器会陷入永久睡眠,任务 hang 死。这条铁律后面还会反复强调。

4.3 第三步:加上 spawn 和就绪队列

block_on 一次只能跑一个 Future。真正的运行时需要 spawn:把成千上万个任务丢进去并发跑。这就需要一个「就绪队列」。

我们把每个 spawn 进来的 Future 包成一个 TaskTask 自己实现 ArcWake —— 它的 wake 动作就是「把自己(Arc<Task>)通过 channel 发回就绪队列」。执行器主循环不停从队列 recv,取出 Task 就 poll。这是 mini-tokio 最经典的设计(tokio 官方教程也是这个骨架)。

use std::future::Future;
use std::pin::Pin;
use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
use std::sync::{Arc, Mutex};
use std::task::{Context, Poll};

type BoxFuture = Pin<Box<dyn Future<Output = ()> + Send>>;

struct Task {
    // 用 Mutex 包住 Future:因为 wake 可能来自任意线程,poll 需要 &mut
    future: Mutex<Option<BoxFuture>>,
    // wake 时把 Arc<Task> 塞回这个发送端
    sender: SyncSender<Arc<Task>>,
}

impl ArcWake for Task {
    fn wake_by_arc(self: Arc<Self>) {
        // 把自己重新排进就绪队列。try_send 满了就丢弃(队列容量要给够)
        let _ = self.sender.try_send(self.clone());
    }
}

pub struct Executor {
    ready_queue: Receiver<Arc<Task>>,
}

#[derive(Clone)]
pub struct Spawner {
    sender: SyncSender<Arc<Task>>,
}

pub fn new_executor_and_spawner() -> (Executor, Spawner) {
    const MAX_QUEUED: usize = 100_000;
    let (sender, ready_queue) = sync_channel(MAX_QUEUED);
    (Executor { ready_queue }, Spawner { sender })
}

impl Spawner {
    pub fn spawn(&self, future: impl Future<Output = ()> + Send + 'static) {
        let task = Arc::new(Task {
            future: Mutex::new(Some(Box::pin(future))),
            sender: self.sender.clone(),
        });
        // 新任务立即入队,等待第一次 poll
        self.sender.try_send(task).expect("ready queue full");
    }
}

impl Executor {
    pub fn run(&self) {
        // 主循环:从就绪队列取任务,poll 它
        while let Ok(task) = self.ready_queue.recv() {
            let mut slot = task.future.lock().unwrap();
            if let Some(mut future) = slot.take() {
                // 用这个 task 自己造一个 waker:wake=把 task 塞回队列
                let waker = waker_from_arc(task.clone());
                let mut cx = Context::from_waker(&waker);

                match future.as_mut().poll(&mut cx) {
                    Poll::Ready(()) => { /* 任务完成,future 已被 take 走,不放回 */ }
                    Poll::Pending => {
                        // 没完成,把 future 放回 slot,等下次被 wake 时再 poll
                        *slot = Some(future);
                    }
                }
            }
        }
    }
}

这几十行就是一个能并发跑任意多个任务的执行器了。核心机制再强调一遍:

  • 每个 Task 既是「被执行的任务」,又是「自己的 Waker」。这是最优雅的设计——waker.wake() 就是 task.sender.send(task),把自己重新排队。
  • 主循环是单线程串行 poll,但因为每个任务遇到 I/O 就返回 Pending 让出,所以能在一个线程上并发处理海量任务。这就是「异步 = 用户态协作式调度」的精髓。
  • futureMutex<Option<...>> 包着:Option 是为了 poll 到 Ready 后能 take() 掉,Mutex 是因为 wake 可能从别的线程来,需要线程安全(尽管 poll 本身是单线程的)。

4.4 第四步:Reactor —— 把 epoll 接进来

前面的例子里,Future 靠 cx.waker().wake_by_ref() 主动重新入队,这其实是「忙轮询」的变种,只适合演示。真实的异步 I/O 必须靠 Reactor:它用一个线程守着 epoll_wait(Linux)/ kqueue(macOS),谁的 fd 就绪了,就调用谁的 waker。

跨平台 epoll 封装我们用 mio(tokio 底层也用它)。Reactor 维护一个 HashMap<Token, Waker>Token 是 fd 的唯一标识,value 是「等这个 fd 就绪的那个任务的 waker」。

use mio::{Events, Interest, Poll as MioPoll, Token};
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::task::Waker;

pub struct Reactor {
    registry: mio::Registry,
    // token -> 等待该 fd 就绪的 waker(读/写分开,这里简化只存一个)
    wakers: Arc<Mutex<HashMap<Token, Waker>>>,
}

impl Reactor {
    pub fn new() -> std::io::Result<Arc<Self>> {
        let poll = MioPoll::new()?;
        let registry = poll.registry().try_clone()?;
        let wakers: Arc<Mutex<HashMap<Token, Waker>>> = Arc::new(Mutex::new(HashMap::new()));

        // 起一个后台线程专门跑 epoll_wait 事件循环
        let wakers_bg = wakers.clone();
        std::thread::Builder::new()
            .name("mini-tokio-reactor".into())
            .spawn(move || Self::event_loop(poll, wakers_bg))
            .unwrap();

        Ok(Arc::new(Reactor { registry, wakers }))
    }

    fn event_loop(mut poll: MioPoll, wakers: Arc<Mutex<HashMap<Token, Waker>>>) {
        let mut events = Events::with_capacity(1024);
        loop {
            // 阻塞等待任意 fd 就绪;epoll_wait 的本体
            poll.poll(&mut events, None).expect("epoll poll failed");
            let mut guard = wakers.lock().unwrap();
            for event in events.iter() {
                // 某个 fd 就绪了,取出对应 waker,唤醒等它的任务
                if let Some(waker) = guard.remove(&event.token()) {
                    waker.wake(); // ← 这一下,任务被塞回执行器就绪队列
                }
            }
        }
    }

    /// 任务在返回 Pending 前调用:注册「我在等这个 fd,就绪了请 wake 我」
    pub fn register_waker(&self, token: Token, waker: Waker) {
        self.wakers.lock().unwrap().insert(token, waker);
    }

    pub fn registry(&self) -> &mio::Registry {
        &self.registry
    }
}

Reactor 的精髓就一句话:它是「fd 就绪事件」到「Waker::wake()」的翻译器。 epoll 告诉它「token=5 的 fd 可读了」,它就去 map 里找「谁在等 token=5」,把那个 waker 一 wake(),对应任务立刻被塞回执行器队列,下一轮就会被 poll,这次 poll 里读 socket 就能成功拿到数据了。

4.5 第五步:异步 TcpStream 与 echo server

现在把 Reactor 用起来,包装一个异步的 TcpStream。核心是实现一个 read 方法返回的 Future:poll 时尝试非阻塞 read,WouldBlock 就注册 waker 并返回 Pending

use mio::net::TcpStream as MioTcpStream;
use std::io::{self, Read, Write};
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::task::{Context, Poll};

pub struct TcpStream {
    inner: MioTcpStream,
    token: Token,
    reactor: Arc<Reactor>,
}

impl TcpStream {
    pub fn new(mut inner: MioTcpStream, token: Token, reactor: Arc<Reactor>) -> io::Result<Self> {
        // 向 epoll 注册这个 fd,关注可读可写
        reactor.registry().register(
            &mut inner,
            token,
            Interest::READABLE | Interest::WRITABLE,
        )?;
        Ok(TcpStream { inner, token, reactor })
    }

    /// 返回一个 Future:异步读,读到数据 Ready(n),没数据就挂起
    pub fn read<'a>(&'a mut self, buf: &'a mut [u8]) -> ReadFuture<'a> {
        ReadFuture { stream: self, buf }
    }
}

pub struct ReadFuture<'a> {
    stream: &'a mut TcpStream,
    buf: &'a mut [u8],
}

impl<'a> Future for ReadFuture<'a> {
    type Output = io::Result<usize>;

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<usize>> {
        let this = self.get_mut();
        // 尝试非阻塞读
        match this.stream.inner.read(this.buf) {
            Ok(n) => Poll::Ready(Ok(n)), // 读到了(n=0 表示对端关闭)
            Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => {
                // 没数据可读:把「读就绪时叫醒我」登记给 Reactor,然后挂起
                this.stream
                    .reactor
                    .register_waker(this.stream.token, cx.waker().clone());
                Poll::Pending
            }
            Err(e) => Poll::Ready(Err(e)),
        }
    }
}

这段代码把前面所有零件串起来了,完整的一次「异步读」生命线:

  1. 执行器 poll 这个 ReadFuture
  2. 非阻塞 read 返回 WouldBlock(socket 里现在没数据);
  3. Future 把 cx.waker()(也就是「重新排队本任务」的能力)clone 一份,登记给 Reactor:「token 就绪了 wake 我」,然后返回 Pending
  4. 执行器把任务放回休眠区,转去 poll 别的任务;
  5. 某时刻客户端发来数据,内核 socket 缓冲区可读,epoll 触发,Reactor 后台线程 epoll_wait 返回;
  6. Reactor 从 map 里找到这个 token 的 waker,调用 wake() → 任务被塞回执行器就绪队列;
  7. 执行器再次 poll 这个 ReadFuture,这次 read 成功返回数据,Poll::Ready(Ok(n))

没有一个线程在 read 上阻塞,一个执行器线程可以同时「照看」几万个连接。 这就是 C10K/C10M 问题的用户态答案,也是 io_uring 出现之前所有高性能网络库(Netty、libuv、tokio)的共同底座。

组装一个 echo server(省略部分样板,展示主逻辑):

fn main() -> io::Result<()> {
    let reactor = Reactor::new()?;
    let (executor, spawner) = new_executor_and_spawner();

    // 用 mio 的 TcpListener 接受连接(同样是非阻塞 + epoll)
    let addr = "127.0.0.1:8080".parse().unwrap();
    let mut listener = mio::net::TcpListener::bind(addr)?;
    // ... 把 listener 也注册进 reactor,用一个 accept 循环 spawn 出连接任务 ...

    let reactor2 = reactor.clone();
    spawner.spawn(async move {
        let mut next_token = 100usize;
        loop {
            let (stream, _peer) = accept(&mut listener, &reactor2).await.unwrap();
            let token = Token(next_token);
            next_token += 1;
            let mut conn = TcpStream::new(stream, token, reactor2.clone()).unwrap();

            // 每个连接一个任务:读到啥就写回啥
            spawner_clone.spawn(async move {
                let mut buf = [0u8; 4096];
                loop {
                    match conn.read(&mut buf).await {
                        Ok(0) => break,                 // 对端关闭
                        Ok(n) => { let _ = conn.write_all(&buf[..n]).await; }
                        Err(_) => break,
                    }
                }
            });
        }
    });

    executor.run(); // 主线程跑执行器
    Ok(())
}

cargo run 起来后,nc 127.0.0.1 8080 打字回车,服务器原样回显。你已经拥有一个从 Future 协议到 epoll、纯手工打造的异步运行时了。 麻雀虽小,但 tokio 的五脏俱全。

4.6 进阶:多线程 work-stealing 执行器

上面的执行器是单线程的(executor.run() 跑在一个线程)。要吃满多核,需要多线程执行器。最朴素的做法是「N 个线程共享一个 MPMC 就绪队列」,但单一全局队列会有严重的锁竞争。tokio 用的是 work-stealing(工作窃取) 调度器:每个 worker 线程有自己的本地队列,优先跑本地任务,本地空了就去别的 worker 的任务。

crossbeam-deque 可以快速搭一个原型:

use crossbeam_deque::{Injector, Stealer, Worker};
use std::sync::Arc;

struct MultiThreadExecutor {
    global: Arc<Injector<Arc<Task>>>,    // 全局注入队列(spawn 落这里)
    stealers: Vec<Stealer<Arc<Task>>>,   // 所有 worker 的偷取句柄
}

fn worker_loop(
    local: Worker<Arc<Task>>,
    global: Arc<Injector<Arc<Task>>>,
    stealers: Vec<Stealer<Arc<Task>>>,
) {
    loop {
        // 1. 优先从本地队列拿(LIFO,缓存局部性好)
        let task = local.pop()
            // 2. 本地空了,从全局队列批量搬一批到本地
            .or_else(|| std::iter::repeat_with(|| global.steal_batch_and_pop(&local))
                .find(|s| !s.is_retry())
                .and_then(|s| s.success()))
            // 3. 全局也空,随机偷别的 worker
            .or_else(|| stealers.iter()
                .map(|s| s.steal())
                .find(|s| s.is_success())
                .and_then(|s| s.success()));

        match task {
            Some(task) => poll_task(task),  // poll 逻辑同单线程版
            None => {
                // 全空,短暂 park 避免空转烧 CPU(真实实现用条件变量/事件)
                std::thread::yield_now();
            }
        }
    }
}

work-stealing 的三个关键设计点,也是面试常考:

  1. 本地队列用 LIFO(栈序):刚 spawn 的子任务大概率数据还在 CPU cache 里,立刻执行缓存命中率高(时间局部性)。
  2. 偷取用 FIFO(队首):偷「最老」的任务,减少和 owner 线程的争抢(owner 从队尾拿,thief 从队头偷,两端操作冲突最小)。
  3. 偷取要随机化起点:所有空闲线程都从 worker[0] 开始偷会造成「羊群效应」,随机起点能把偷取压力打散。

真实的 tokio 调度器还有更多精妙设计:LIFO slot(专门优化「A 唤醒 B,B 马上要跑」的消息传递场景)、全局队列的定期检查(防止本地任务把全局任务饿死)、park/unpark 的精细化管理(避免线程空转又避免唤醒延迟)。但骨架就是上面这套。


五、性能优化:从能跑到跑得快

手写运行时能跑之后,真正拉开与生产级运行时差距的是这些细节:

5.1 减少无谓的唤醒与 poll

每一次 wake() → 入队 → poll() 都有成本(加锁、context 构造、状态机 match)。常见浪费:

  • 伪唤醒(spurious wakeup):Future 被 poll 了,但其实没就绪,只能再次返回 Pending。要尽量让「就绪」和「wake」严格对应。
  • 重复注册 waker:每次 Pending 都无脑 register_waker + clone 是有开销的。优化是先用 waker.will_wake(cx.waker()) 判断是不是同一个 waker,相同就不必重复 clone/注册。tokio 内部大量用这个技巧。
// 优化前:每次 Pending 都 clone + 注册
this.reactor.register_waker(token, cx.waker().clone());

// 优化后:只在 waker 变了时才更新
if !stored_waker.as_ref().is_some_and(|w| w.will_wake(cx.waker())) {
    *stored_waker = Some(cx.waker().clone());
}

5.2 批量 syscall 与边缘触发

  • epoll 用边缘触发(ET)而非水平触发(LT):LT 模式下只要 fd 还有数据没读完,epoll 每次都报就绪,容易重复唤醒;ET 只在「状态变化」时通知一次,配合「一次 poll 里 loop read 到 WouldBlock 为止」,能大幅减少 epoll 事件数量。代价是逻辑更容易写错(漏读会永久 hang)。
  • 一次 epoll_wait 处理一批事件Events::with_capacity(1024) 让单次系统调用返回尽可能多的就绪事件,摊薄 syscall 开销。这也是 io_uring 进一步优化的方向——连 epoll_wait 这次系统调用都省掉。

5.3 内存布局与 Pin

  • 大 Future 的 move 成本async fn 生成的状态机大小等于「所有 await 点上活跃变量的最大集合」。一个巨大的 async fn(比如里面有个 [u8; 65536] 局部数组跨越 await)会生成一个巨大的状态机,spawnBox::pin 一次堆分配 + move 成本高。优化:把大缓冲区放到堆上(Vec/Box)而非栈上局部变量,让状态机瘦身。
  • Box::pin vs 内联spawn 顶层任务必须 Box::pin(类型擦除进 dyn Future),但内部的 .await 子 Future 是内联在状态机里的,零额外分配。别在热路径上手动 Box::pin 每个子 Future。

5.4 避免 false sharing 与锁竞争

  • 多线程执行器里,若多个 worker 频繁访问相邻内存(比如任务计数器数组),会因为 CPU cache line(64 字节)共享导致 false sharing。用 crossbeam_utils::CachePadded 把每个 worker 的热数据填充到独立 cache line。
  • 全局就绪队列是最大的锁竞争点。work-stealing 的本地队列本质就是「把全局锁拆成 N 个几乎无竞争的本地结构」。

5.5 一个反直觉的点:block_in_place 与阻塞任务

异步运行时最怕的不是慢,是在异步任务里干阻塞的活(比如同步文件 I/O、std::thread::sleep、CPU 密集计算)。因为执行器线程数是固定的(通常等于核数),一个任务阻塞住一个执行器线程,就等于永久废掉一个核。tokio 的解法是把阻塞活儿丢到专门的 spawn_blocking 线程池。手写运行时时也要意识到:执行器线程上跑的每一段代码,两个 await 之间都必须是「几乎立即返回」的,否则整个飞轮就卡住了。


六、踩坑清单:async 领域的「新手劝退大礼包」

这份清单是无数人用编译错误和线上 hang 死换来的,逐条对照能省你几天时间。

坑 1:返回 Pending 却忘了安排 wake → 任务永久 hang。
最经典的错误。铁律:任何返回 Poll::Pending 的代码路径,都必须保证在此之前已经把 cx.waker() 交给了某个「将来会 wake() 它」的东西(Reactor、定时器、channel)。如果你只是 return Poll::Pending 而没注册 waker,这个任务再也不会被 poll,静默死掉,且不报任何错——极难排查。

坑 2:Pin 与自引用——cannot borrow across await point
状态机是自引用的,所以 Future 一旦开始 poll 就不能移动。表现为:!Unpin 的 Future 必须先 Box::pinpin! 才能 poll。理解「为什么需要 Pin」(保护自引用指针)比死记语法重要。

坑 3:.await 期间持有 std::sync::Mutex 的锁 → 死锁或 !Send

let guard = std_mutex.lock().unwrap();
some_async_fn().await;   // ❌ guard 跨越 await 点

标准库 MutexGuard 不是 Send,跨 await 会让整个 Future 变 !Send,无法 spawn 到多线程运行时。更糟的是,持锁 await 时任务被挂起,别的任务想拿这把锁就死锁。解法:await 前先 drop(guard),或改用 tokio::sync::Mutex(但它更慢,能不用就不用)。

坑 4:伪唤醒(spurious wakeup)——poll 不保证「一定就绪」。
执行器可能因为各种原因(同一 waker 被 wake 多次、别的 fd 事件误触)重新 poll 你的 Future。所以 poll 的实现必须是幂等的:每次进来都要重新检查真实状态(重新 try_read),不能假设「被 poll 了就一定有数据」。

坑 5:select! 里的 Future 被取消 → 状态丢失。
tokio::select! 中没胜出的分支 Future 会被 drop 掉。如果那个 Future 已经消费了一部分数据(比如从 channel 收了一半),drop 会导致数据丢失。这就是「取消安全性(cancellation safety)」问题,是 async Rust 最微妙的坑之一。写 select! 前务必确认每个分支的 Future 是 cancel-safe 的。

坑 6:忘了 Send + 'static → spawn 编译不过。
多线程执行器的 spawn 要求 Future 是 Send + 'static。任何跨 await 持有的东西(局部变量、引用)都要满足。Rc、裸指针、非 'static 引用都会让 Future !Send。这也是为什么 async Rust 里 Arc<Mutex<T>> 满天飞。

坑 7:在 Drop 里做异步操作 → 做不到。
Drop::drop 是同步的,不能 .await。想在资源释放时做异步清理(比如优雅关闭连接、flush 缓冲),只能靠显式的 async close() 方法或后台任务,不能指望 Drop

坑 8:CPU 密集任务霸占执行器线程。
前面讲过,一段跑 500ms 的纯计算(无 await)会卡死一个执行器线程 500ms,期间它照看的所有连接全部延迟。解法:spawn_blocking,或在长循环里手动插 tokio::task::yield_now().await 主动让出。


七、总结与展望

我们从「Future 是一个可被反复 poll 的惰性状态机」这一句话出发,一路手写出了 Waker(RawWaker vtable)、单线程执行器、spawn + 就绪队列、基于 epoll 的 Reactor、异步 TcpStream,最后组装成一个能跑的 echo server,还给了多线程 work-stealing 的原型。回头看,整个 tokio 的骨架无非就是三个咬合的齿轮:

  • Executor:不停从就绪队列取任务 poll,是「驱动力」;
  • Reactor:把 OS 的 I/O 就绪事件翻译成 Waker::wake(),是「感知器」;
  • Wakerwake() = 把任务塞回就绪队列,是连接前两者的「传动轴」。

.await 的本质是「poll 子 Future,Pending 就带着现场往上冒泡,等 Waker 把自己叫回来」。想通这一点,之前那些 PinSend'static、伪唤醒的报错,都会从「玄学」变成「理所当然」。

站在 2026 年往前看,异步运行时这套模型正在被两股力量重塑:

  1. io_uring 正在替换 epoll。epoll 是「就绪通知」模型(告诉你能读了,你再去 read,两次 syscall),io_uring 是「完成通知」模型(你提交读请求,内核读完直接把数据放好通知你,syscall 大幅减少甚至归零)。tokio 的 tokio-uring、以及 monoio、glommio 这些 thread-per-core 运行时,正在把 Reactor 换成 io_uring 的 completion 模型。本文的 Reactor 那一层,未来会被彻底改写。

  2. async 的人体工程学还在补课。async trait 已经稳定,async closures 在 2024 edition 落地,但「async Drop」「返回位置的 impl Trait 细节」「更好的取消语义」仍在演进。Rust 的异步故事远没讲完。

但无论底层怎么变,「惰性状态机 + poll + waker」这个核心协议是稳定的。你手写过一遍运行时,就拥有了看穿任何异步框架的透视眼——无论是 tokio、monoio,还是明天出现的新东西。这,就是「造轮子」不可替代的价值:不是为了用自己的轮子,而是为了永远不再怕任何轮子。

下次再看到 Pin<Box<dyn Future<Output = ()> + Send>> 这一长串,希望你想到的不再是「什么鬼」,而是「哦,一个被钉住的、类型擦除的、能跨线程的惰性状态机而已」。共勉。

推荐文章

Elasticsearch 文档操作
2024-11-18 12:36:01 +0800 CST
2024年公司官方网站建设费用解析
2024-11-18 20:21:19 +0800 CST
微信小程序热更新
2024-11-18 15:08:49 +0800 CST
程序员茄子在线接单