如何在Tokio多线程中控制仅指定数量的任务并行运行
如何在Tokio中限制任务的最大并发数
你不需要刻意创建任务却不启动,Tokio生态中有两种简洁的方式实现“最多10个任务并行,完成一个再启动下一个”的需求:
方法一:使用Tokio内置的信号量(Semaphore)
信号量是控制并发数的经典方案,通过许可机制限制同时运行的任务数量:
use tokio::sync::Semaphore; use std::sync::Arc; #[tokio::main(flavor = "multi_thread", worker_threads = 10)] async fn main() { // 初始化信号量,设置最大并发许可数为10 let semaphore = Arc::new(Semaphore::new(10)); let mut tasks = vec![]; for i in 0..100 { // 获取一个并发许可,若无可用许可则阻塞等待 let permit = semaphore.clone().acquire_owned().await.unwrap(); let task = tokio::spawn(async move { println!("This process: {}", i); // 任务完成后自动释放许可(也可手动drop提前释放) drop(permit); i }); tasks.push(task); } let mut results = vec![]; for task in tasks { let result = task.await.unwrap(); results.push(result); } dbg!(&results); }
原理说明
- 信号量初始化时拥有10个可用许可,每个任务必须先获取许可才能执行
- 当10个许可被全部占用时,后续任务会进入等待状态,直到某个任务完成并释放许可
- 许可通过
Arc共享给所有任务,任务持有许可的所有权,任务结束后许可自动归还信号量
方法二:使用futures库的buffer_unordered
如果希望代码更简洁,可以借助futures库的流处理能力,直接控制并发任务数:
首先在Cargo.toml中添加依赖:
[dependencies] tokio = { version = "1.0", features = ["full"] } futures = "0.3"
然后编写代码:
use futures::stream::{StreamExt, iter}; #[tokio::main(flavor = "multi_thread", worker_threads = 10)] async fn main() { // 将任务迭代器转为异步流,限制并发数为10 let results = iter(0..100) .map(|i| async move { println!("This process: {}", i); i }) .buffer_unordered(10) .collect::<Vec<_>>() .await; dbg!(&results); }
原理说明
iter(0..100)将数字范围转为异步流map将每个元素映射为一个异步任务buffer_unordered(10)会同时运行最多10个任务,当某个任务完成后,自动启动流中的下一个任务- 最终通过
collect收集所有任务的结果
为什么不推荐“创建任务不启动”?
Tokio的spawn设计就是将任务立即提交到调度器,手动控制任务启动时机不符合Tokio的异步调度模型。上述两种方法都是在异步调度框架内实现并发限制,既简洁又能保证任务调度的高效性。
内容的提问来源于stack exchange,提问作者Kingfisher Phuoc
相关产品推荐
相关产品推荐

