tokio:异步运行时是怎么回事
一个线程怎么同时伺候一万个连接?把 async/await、Future、运行时调度讲透,再把 forge-store 的接口改造成异步版本。
学完这节你能做到
- 解释 async fn 编译成状态机的过程
- 说清 tokio 的多线程调度与 work stealing
- 知道什么时候该用 spawn_blocking:文件 IO 的尴尬地位
从 epoll 说起:你在监控里见过的那些 fd
做运维时你八成排查过这样的问题:一台 nginx 顶着几万并发连接,ss -s 显示几万条 established,
但 top 里它只有几个 worker 进程,CPU 还很闲。几万个连接,几个进程,怎么伺候得过来?
答案你其实知道,叫 epoll:进程把几万个 fd 注册给内核,然后睡觉; 哪个 fd 上来了数据,内核把它叫醒,它只处理「有事的那几个」。 这就是「一个线程伺候一万个连接」的全部秘密 —— 不是更快地轮询,而是不轮询,等通知。
问题在于,直接用 epoll 写程序是灾难:每个连接的处理逻辑被切成碎片,
状态要自己存(读到一半的 buffer、处理到哪一步),代码变成一张回调织成的网。
async/await 就是为了解决这个:让你用顺序代码的写法,得到 epoll 的并发能力。
编译器负责把顺序代码切碎,运行时负责在事件到来时接着跑。
Future 与状态机:await 点就是让出点
一个 async fn 编译后不是普通函数,而是一个状态机(实现了 Future trait 的匿名类型):
async fn handle_get(conn: &mut Conn, key: &[u8]) -> Result<(), Error> {
let req = conn.read_request().await?; // 状态 1:等请求可读
let value = store_lookup(&req)?; // 同步计算,不让出
conn.write_response(&value).await?; // 状态 2:等写缓冲可写
Ok(())
}
编译器把它变成大致这样的东西:一个 enum,每个 .await 是一个状态;
每次被 poll,就从上次停下的状态继续跑,跑到下一个 .await 如果没就绪,
保存现场、返回 Pending、把线程让出来。
三个要点,记住它们后面全用得上:
.await是唯一的让出点。两个 await 之间的同步代码会一口气跑完, 谁也打断不了 —— 这既是性能(没有抢占开销),也是坑(下面讲)。- Future 是惰性的。不 poll 就什么都不发生,这和你熟悉的「起个线程就开跑」完全不同。
- 状态机就是那点被切碎的栈。一万个任务就是一万个小 struct,而不是一万个 8MB 的线程栈 —— 这就是 async 能扛超高并发的内存账。
tokio 运行时:worker 线程与任务队列
状态机自己不会跑,需要有人反复 poll 它 —— 这就是运行时(runtime)的活。 tokio 的多线程运行时结构,用运维的眼光看非常眼熟:
- worker 线程:默认每个 CPU 核一个,类比 nginx 的 worker_processes auto
- 每个 worker 一个本地任务队列 + 一个全局注入队列
- work stealing:自己队列空了,去偷别的 worker 的任务 —— 负载自动均衡, 不会出现「一个核 100% 其他核围观」
- 一个 IO driver:封装 epoll,fd 就绪时把对应任务标记为可运行
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let listener = tokio::net::TcpListener::bind("0.0.0.0:9090").await?;
loop {
let (socket, _addr) = listener.accept().await?;
// spawn:把任务丢进队列,不等它完成 —— 一万个连接就是一万个轻量任务
tokio::spawn(async move {
if let Err(e) = serve_conn(socket).await {
eprintln!("conn error: {e}");
}
});
}
}
worker 线程是全体任务共享的。你在一个任务里跑了 100ms 的同步计算或同步 IO,
这个 worker 上排队的所有任务都陪着卡 100ms —— 表现为 P99 毛刺,和你运维时
见过的「某核被单进程打满、软中断堆积」是同一类病。规矩:两个 await 之间
不许有慢操作;慢的丢给 spawn_blocking。
文件 IO 的真相:tokio::fs 底下是线程池
现在说一个对存储工程师至关重要的事实:epoll 对普通文件基本无效。 epoll 的模型是「fd 没就绪时等通知」,但磁盘文件永远「就绪」—— read 会直接发起, 然后阻塞到数据从盘上回来。内核没有给 buffered 文件 IO 提供像 socket 那样的就绪通知。
所以 tokio::fs::File 的实现是个「善意的谎言」:表面是 async,
底下把每个操作丢进一个阻塞线程池(spawn_blocking,默认上限 512 个线程)执行。
你得到了不卡 worker 的效果,但代价是:每次 IO 一次线程切换 + 任务在池里排队。
用你熟的量级感受一下:一块 NVMe 标称 100 万 4k 随机读 IOPS, 每个 IO 走一趟线程池的调度开销在几微秒量级,叠加线程数上限, 线程池方案单机通常在十几万到几十万 IOPS 就见顶 —— 离硬件差着倍数。 这就是下一课 io_uring 出场的理由。
为什么 tokio 处理 TCP 用 epoll 就够了,处理文件却要走线程池?
改造 forge-store:同步核心 + 异步外壳
L1 写完的 forge-store 是纯同步的,这是对的,不要重写。我们采用一个工程上非常主流的结构:
同步核心 + 异步外壳 —— 引擎内部保持同步(落盘逻辑简单可审),
对外提供 async 接口,内部用 spawn_blocking 桥接:
use std::sync::Arc;
use forge_store::{BlobStore, StoreError};
/// 异步外壳:把同步引擎包进 Arc,IO 操作转投阻塞线程池。
pub struct AsyncStore<S: BlobStore + Send + Sync + 'static> {
inner: Arc<S>,
}
impl<S: BlobStore + Send + Sync + 'static> AsyncStore<S> {
pub async fn get(&self, key: Vec<u8>) -> Result<Option<Vec<u8>>, StoreError> {
let store = Arc::clone(&self.inner);
// 同步 get 在阻塞池里跑,worker 线程不被占用
tokio::task::spawn_blocking(move || store.get(&key))
.await
.map_err(|e| StoreError::Internal(format!("blocking task failed: {e}")))?
}
}
注意两个细节,review AI 代码时专门盯:
spawn_blocking的闭包要move进所有权(所以参数改成Vec<u8>而非借用)—— 借用检查器会拦住你把引用送进另一个线程,这正是 L0 讲的「拦一下 = 别处一次事故」。.await后面那个?处理的是任务本身失败(panic 或取消), 和get的业务错误是两层,不能混在一起吞掉。
这层外壳性能不惊艳,但它让 L3 的 gRPC 服务端(天然 async)能直接调用引擎。 性能问题我们用接下来两课解决,而且解决之后只换外壳,不动核心 —— 这就是接口边界的红利。
review 三个重点:一,JoinError 的处理路径 —— AI 很爱在这里
unwrap,任务 panic 会连坐整个进程;二,Arc 是否真的必要,
让 AI 解释如果去掉会碰到哪个编译错误;三,最后那个反问题它答没答到点上
(worker 线程被同步 IO 卡住,整个运行时的其他任务陪跪)。
在 forge workspace 里给 forge-store 增加异步外壳,要求: 1. forge-store 增加 feature "async",启用时引入 tokio(features: rt-multi-thread, macros),不启用时保持纯同步零依赖; 2. 新增 async_store 模块,实现 AsyncStore,内部持有 Arc 包裹的同步引擎(约束为实现了 BlobStore + Send + Sync 的类型),对外提供 async 的 put/get/delete,内部用 tokio::task::spawn_blocking 桥接; 3. spawn_blocking 的 JoinError(任务 panic/取消)必须映射为独立的错误变体,不许与业务错误混淆,也不许 panic; 4. 单元测试用 tokio::test:并发 spawn 64 个任务同时 put 不同 key,全部完成后逐一 get 校验内容;再测一个 get 不存在的 key 返回 Ok(None); 5. 跑 cargo test -p forge-store --features async 和 cargo clippy,贴出结果。 做完后回答我:如果把 spawn_blocking 换成直接在 async fn 里调用同步 get,并发压测时会发生什么?
小结
- epoll 的本质是「等通知而不是轮询」,async/await 让你用顺序代码写出这种并发
- async fn 编译成状态机,await 点就是让出点;两个 await 之间的慢操作会卡住整个 worker
- tokio 多线程运行时 = 每核一个 worker + work stealing,负载自动均衡
- 普通文件 IO 没有就绪通知,tokio::fs 底下是阻塞线程池 —— 有并发,但离硬件极限远
- forge-store 采用同步核心 + 异步外壳:接口先 async 化,性能下一课用 io_uring 解决