如何将future::join_all()的用法改造为使用FuturesUnordered?
Rust FuturesUnordered 批量异步操作改造方案
错误代码问题说明
自行编写的FuturesUnordered代码存在两处核心错误:
- FuturesUnordered没有接收多个参数的
collect静态方法 - 不需要手动调用
poll方法,FuturesUnordered实现了Stream特征,直接通过StreamExt提供的next方法遍历即可执行所有异步任务
前置依赖配置
首先在Cargo.toml中添加必要依赖:
[dependencies] futures = "0.3" # 若使用tokio作为异步运行时则添加 tokio = { version = "1.32", features = ["rt-multi-thread", "macros", "sync"] }
然后导入必要的特征和结构体:
use futures::stream::FuturesUnordered; use futures::StreamExt;
基础改造版本(对应原有join_all功能)
直接替换原有hold函数的逻辑即可:
async fn hold() { // 创建空的FuturesUnordered实例 let mut futures = FuturesUnordered::new(); // 推入需要执行的异步任务 futures.push(function1()); futures.push(function2()); // 若已有future列表,也可以直接通过迭代器生成 // let futures: FuturesUnordered<_> = vec![function1(), function2()].into_iter().collect(); // 遍历执行所有任务,收集返回结果 let mut _results = Vec::new(); while let Some(res) = futures.next().await { _results.push(res); } }
超大任务量限流优化
如果需要执行上千甚至上万条数据库插入任务,直接将所有future推入FuturesUnordered仍会导致并发过高,耗尽数据库连接或者内存,需要配合信号量限制最大并发数:
use tokio::sync::Semaphore; use std::sync::Arc; async fn hold() { // 根据数据库负载调整最大并发数,比如设置为100 let semaphore = Arc::new(Semaphore::new(100)); let mut futures = FuturesUnordered::new(); // 假设此处生成所有待执行的插入任务 let insert_tasks = generate_all_insert_tasks().await; for task in insert_tasks { let sem_clone = semaphore.clone(); futures.push(async move { // 获取执行许可,超过并发数时会等待 let _permit = sem_clone.acquire().await.expect("信号量获取失败"); task.await }); } // 执行所有任务并处理结果 while let Some(insert_res) = futures.next().await { match insert_res { Ok(_) => println!("数据插入成功"), Err(e) => eprintln!("数据插入失败: {}", e) } } }
原理解释
join_all需要等待所有异步任务全部完成后才会统一返回结果,任务量过大时会占用极高内存,同时没有并发限制很容易触发数据库限流导致程序卡住。FuturesUnordered是任务无序完成的流,每完成一个任务就返回一个结果,内存占用远低于join_all,配合限流逻辑可以稳定处理超大规模的异步批量操作。
内容的提问来源于stack exchange,提问作者Lee Alex
相关产品推荐
相关产品推荐

