Rust + Tokio:等待未知长度任务列表中首个完成任务的方法
Tokio动态任务列表首任务完成等待实现方案
要实现动态长度任务列表下等待首个完成任务、支持处理任务时追加新任务、最终实现N并发限流的需求,不要使用编译期固定分支的tokio::select!,直接使用FuturesUnordered即可,这是该场景下的标准实现方式。
核心特性匹配
FuturesUnordered是futures工具库提供的future集合,完全匹配需求点:
- 支持运行时随时向集合内追加新任务,没有固定长度限制
- 调用
next()方法时会自动等待集合内所有任务,返回最先执行完成的任务结果 - 任务完成后会自动从集合中移除,不会重复调度
- 可以随时查询集合当前长度,方便做并发数控制
实现逻辑
整体流程完全贴合需求:
- 初始化
FuturesUnordered实例作为运行中任务容器,额外维护一个待执行任务队列做限流缓冲,设置最大并发数N - 启动阶段先从待执行队列取最多N个任务,spawn到Tokio运行时后将JoinHandle推入运行中集合
- 循环调用运行中集合的
next()方法,阻塞等待首个完成的任务返回 - 拿到完成任务的结果后,执行自定义处理逻辑,处理过程中产生的新任务直接推入待执行队列即可
- 每处理完一个任务,就从待执行队列取任务补入运行中集合,保证运行中的任务数始终不超过N,直到待执行队列和运行中集合都为空,流程结束
可运行代码示例
use std::future::Future; use futures::stream::FuturesUnordered; use futures::StreamExt; use tokio::task::JoinHandle; /// N并发限流执行任务 /// # 参数 /// - tasks: 初始任务列表,每个任务是返回Future的闭包 /// - concurrency: 最大并发数 async fn concurrency_limited_run<F, Fut>(mut tasks: Vec<F>, concurrency: usize) where F: FnOnce() -> Fut + Send + 'static, Fut: Future<Output = ()> + Send + 'static, { let mut running_tasks = FuturesUnordered::new(); // 初始化填充第一批任务,打到并发上限 for _ in 0..concurrency.min(tasks.len()) { let task = tasks.remove(0); running_tasks.push(tokio::spawn(task())); } // 循环处理完成的任务 while let Some(join_result) = running_tasks.next().await { // 处理任务执行结果,可自定义错误处理、结果消费逻辑 match join_result { Ok(_) => { // 任务正常完成的处理逻辑 } Err(join_err) => { eprintln!("task exited unexpectedly: {}", join_err); } } // -------------------------- // 这里可以写任务完成后的业务逻辑 // 逻辑中产生的新任务直接push到tasks队列即可 // 示例:tasks.push(|| async { /* 新任务逻辑 */ }); // -------------------------- // 补位新任务,始终维持并发数不超过上限 while running_tasks.len() < concurrency && !tasks.is_empty() { let next_task = tasks.remove(0); running_tasks.push(tokio::spawn(next_task())); } } }
注意事项
- 不要尝试用
Vec存储JoinHandle配合循环轮询实现,会产生大量无效唤醒,CPU占用高且延迟大 - 单个任务需要超时控制的话,直接在spawn时给future包裹
tokio::time::timeout即可,不需要修改整体调度逻辑 FuturesUnordered采用公平调度机制,不会出现任务长期得不到执行的饥饿问题- 如果需要收集所有任务的返回值,在处理完成任务的逻辑中把返回值存入指定集合即可
内容的提问来源于stack exchange,提问作者Test
相关产品推荐
相关产品推荐

