You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 10:24:04