Tokio 追求进度而非顺序:调度 100 万个任务
摘要
解释了为何 Tokio 的调度器在大规模下不保证任务顺序,使用一个真实世界的 100 万个任务示例,并强调了进度与顺序之间的区别。
暂无内容
查看缓存全文
缓存时间: 2026/07/27 22:47
# Tokio 保证进度,不保证顺序:调度 100 万个任务
来源:https://pranitha.dev/posts/tokio-gives-progress-not-ordering/
这是我上一篇文章的前传,[你的 Rust 服务没有泄漏——可能是分配器的问题](https://pranitha.dev/posts/rust-and-memory-allocators)。在那篇文章中,我介绍了几个内存分配器在我们的工作负载下表现不同。在发现内存行为与分配器相关之前,我们首先尝试从应用程序层面减少内存使用。
我们的服务是事件驱动的:
1. 从消息队列(Kafka/Redis Streams/NATS)读取事件
2. 对每个事件,生成一个 Tokio 任务来处理它
## 我们的任务模式
---
在我们的工作负载中,每个事件最多包含 1000 个用户令牌。对于每个用户令牌,我们需要发起一个出站调用,等待 I/O 并收集所有属于该事件的响应。以下是一个简化版本的代码,我们对用户令牌进行扇出(fan-out)Tokio 任务,扇入(fan-in)响应,并生成一个响应事件。
```rust
struct Event {
payload: Bytes, // ~4KB
user_tokens: Vec<UserToken>, // 最多 1000 个令牌
// 其他字段
}
// ... 在主循环中
{
let event: Event = fetch_next_event().await;
tokio::spawn(async move {
let data = event.payload.clone();
let mut tasks = JoinSet::new();
for token in &event.user_tokens {
let token = token.clone();
let data = data.clone();
tasks.spawn(async move {
// 调用一个出站 API 并返回响应
process(token, data).await
});
}
let mut responses = Vec::with_capacity(event.user_tokens.len());
while let Some(res) = tasks.join_next().await {
responses.push(res)
}
generate_response_event(event, responses);
});
}
```
我们主要关注的是吞吐量。虽然上面的扇出是无限制的,但我认为实际上不会成为问题。这些 Tokio 任务的生命周期很短,它们会在收到响应后立即完成并释放,通常只需几毫秒。所以我假设早期生成的任务也会早期完成。尽管单个出站调用可能无序完成,但我期望整体上早期的事件会先完成,即使新事件还在不断到来。我们的要求只是让每个出站调用尽可能快,以便突发流量能在预期时间内完成。
## 日志看起来什么样
---
以下是在突发 1000 个事件(每个事件约 1000 个用户令牌,总共生成约 100 万个任务)时的日志样子:
```
started: event 1, user 5
started: event 1, user 8
started: event 1, user 2
started: event 2, user 6
started: event 2, user 4
started: event 3, user 8
...
finished: event 779
finished: event 976
started: event 900, user 42
started: event 900, user 261
started: event 1, user 974
started: event 1, user 831
...
finished: event 5
finished: event 3
```
虽然突发流量在预期时间内完成,但日志显示来自早期事件的令牌任务被启动的时间要晚得多。我并没有要求任务之间有严格的顺序。它们可以以任何顺序完成,这没关系。让我惊讶的是,一些较早的任务在首次被轮询之前,偏离其提交顺序竟然如此之远。
## Tokio 调度器内部
---
Tokio 的多线程运行时具有固定数量的工作者线程、每个工作者一个本地队列,以及所有工作者共享的全局队列。每个工作者有一个容量为 256 个任务的本地队列,当溢出时,它会将一半的任务移动到全局队列。工作者优先从自己的本地队列拉取任务,偶尔检查全局队列,并在空闲时从其他工作者那里窃取任务。
当任务变得可独立调度时,Tokio 不再知道它们是由哪个事件创建的。现在它们只是可运行的任务,相互竞争以被轮询。一个简化的图如下所示:
```
全局队列
+------------------------------+
| * 从本地队列溢出的任务 |
| * 远程调度的任务 |
| * 新旧工作混合 |
+--------------+---------------+
|
+-------------------+-------------------+
| | |
v v v
worker 0 worker 1 worker 2
+-------------+ +-------------+ +-------------+
| 本地队列 | | 本地队列 | | 本地队列 |
| 最多 256 | | 最多 256 | | 最多 256 |
+------+------+ +------+------+ +------+------+
| | |
v v v
轮询任务 轮询任务 轮询任务
```
在这种架构下,一旦我们的事件扇出为 1000 个用户令牌任务,这些任务就会与其他事件的令牌任务、等待 `JoinSet` 的父事件任务以及从 I/O 就绪状态唤醒的任务混合在一起。由于同时提交了这么多任务,Tokio 不保证较早的任务会先被轮询。来自不同事件的任务竞争可用的工作者队列,而队列溢出和工作窃取等调度决策会改变任务的拾取顺序。结果,一些来自早期事件的任务首次轮询的时间要晚得多。
核心区别在于:
```
任务创建 != 任务轮询 != 任务完成
```
虽然 Tokio 会持续推进每一个任务,但内存受同时存活的线程数量影响。每个 Tokio 任务都带有一些状态,虽然单个状态可能不大,但成千上万个任务同时存活时,这些状态就会累积成很高的峰值内存使用。而且,由于在我们工作负载中,来自早期事件的一些任务一直存活到突发流量结束,它们的父事件任务也就保持存活,从而让事件状态保留在内存中,进而增加了我们的峰值内存使用。
## 任务创建需要限界
---
Tokio 保证(https://docs.rs/tokio/latest/tokio/runtime/index.html#detailed-runtime-behavior)在任务数量有界且假设没有任务阻塞工作者线程的情况下实现公平调度。我们的代码最初没有对任务创建施加限制:
```
尽可能快地读取事件
└── 为每个事件生成一个事件任务
└── 为每个事件生成最多 1000 个令牌任务
```
你交给 Tokio 多少任务,它就会接收多少,并持续推进它们。Tokio 为实现公平性所假定的有界性必须来自应用程序。这些界限因工作负载而异。在我们的案例中,我们想要的是事件级别的公平性。我们希望属于单个事件的所有令牌任务能够在其提交时间附近被轮询并完成。为了实现这一点,我们使用了一个 `Semaphore` 来限制同时处理的事件数量。
```
只允许 N 个事件同时进入处理
└── 为每个事件生成一个事件任务
└── 为每个事件生成最多 1000 个令牌任务
```
Tokio 同时能看到的任务的最大数量(这影响它能够多公平)也取决于每个任务在被轮询时返回所需的时间。我们通过一些试错找到了对我们有效的 `Semaphore` 计数。一个重要的细节是,即使有了这种有界结构,突发流量仍然在预期时间内完成,同时显著降低了峰值内存。一般来说,添加 `Semaphore` 似乎会降低服务的吞吐量,但对我们来说并非如此。
## 这不是 Tokio 的问题
---
这并非 Tokio 的问题。即使有 100 万个存活的任务,Tokio 仍然持续推进,直到我们的突发流量按时完成。区别在于我的假设(如果事件开始得早,整个事件应该早完成)与 Tokio 运行时看到的(任务就是任务,不管它来自哪个事件)之间的不同。对我来说很容易忽略这一点,因为没有正确性问题。没有任务中途丢弃,服务也没有明显延迟,吞吐量需求也得到了满足。它只是显示为比预期更高的内存峰值,原因在应用程序逻辑中无处可见。
## 结论
---
开始得早并不意味着轮询得早。轮询得早也不意味着完成得早。如果你期望 Tokio 尊重某个应用程序级别的公平单元(例如我们案例中的事件),运行时并不知道它的存在。界限应由应用程序添加。在生成一个 Tokio 任务之前,问问自己同时存活的任务的最大数量可能是多少,以及这可能会如何影响内存。在某些情况下,它还可能影响任务存活时持有的其他资源。
相似文章
Tokio/Rayon 陷阱:为何 async/await 在并发中失灵
本文探讨了 async/await 语法虽然易于编写,但在生产环境中因将异步与并发混为一谈而引入了巨大复杂性,常需手动将 I/O 与计算任务分配到 Tokio 和 Rayon 等不同运行时,导致延迟飙升与系统不稳定。
Cargo 的调度器能否改进?
本文通过分析 Rust 构建任务的基准测试来研究 Cargo 调度器的行为,并探讨替代调度算法。
mpsc 通道的隐藏成本
本文分析了 Rust 中 Tokio 的 mpsc 通道中意想不到的内存分配开销,揭示了由于内部块大小导致的每个通道的固定开销。文章展示了这一开销如何影响诸如 Agent Gateway 这样的大规模应用程序,并建议采用 futures-channel 等替代方案以提高内存效率。
管理无电池物联网中未知工作负载的任务执行:一种硬件无关的评估
本文提出了两种硬件无关的动态调度策略(一种无模型强化学习代理和一种即时近似预测方法),用于管理具有未知工作负载的无电池物联网设备中的任务执行,并使用真实太阳能数据的模拟框架对它们与现有方法进行了评估。
谁在运行你的 Rust Future?动手实践入门异步 Rust
这是一套动手实践教程系列,旨在弥合理解异步 Rust 内部机制(Future、poll、Pin、执行器)与使用 Tokio 部署实际异步代码之间的差距,面向熟悉 JavaScript 异步和基础 Rust 的开发者。