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

如何在Rust异步中等待Future池中的特定结果?

问题解答

1. 如何等待向量中的所有Future?

如果需要等待所有任务完成,可以使用futures::future::join_all,它接收一个Future迭代器,返回一个包含所有任务结果的Future:

use futures::future;
use rand::Rng;
use tokio;

async fn async_function(guess: u8) -> (u8, bool) {
    let random_wait = rand::thread_rng().gen_range(0..2);
    tokio::time::sleep(tokio::time::Duration::from_secs(random_wait)).await;
    println!("Running guess {guess}");
    (guess, guess == 231)
}

#[tokio::main]
async fn main() {
    let pool: Vec<_> = (0..=255).map(async_function).collect();
    let results = future::join_all(pool).await;
    // results 是 Vec<(u8, bool)>,包含所有guess和对应结果
}

但该方法会等待所有任务结束,不符合你提前终止的需求,更适合你的方案见下文。

2. 仅等待第一个返回true的Future,并获取对应的guess

要实现“找到第一个匹配的任务就终止”,核心是并发执行所有任务,一旦找到目标就停止剩余任务,以下是两种可行方案:

方案一:流处理 + 并发执行

通过stream模块实现并发任务调度,找到第一个匹配结果后立即返回:

use futures::stream::{self, StreamExt};
use rand::Rng;
use tokio;

async fn async_function(guess: u8) -> (u8, bool) {
    let random_wait = rand::thread_rng().gen_range(0..2);
    tokio::time::sleep(tokio::time::Duration::from_secs(random_wait)).await;
    println!("Running guess {guess}");
    (guess, guess == 231)
}

#[tokio::main]
async fn main() {
    let found_guess = stream::iter(0..=255)
        .map(|guess| async move {
            let (g, is_match) = async_function(guess).await;
            is_match.then_some(g)
        })
        .buffer_unordered(255) // 并发执行所有任务,参数为并发数
        .find_map(|opt| async move { opt }) // 捕获第一个非None的结果
        .await;

    match found_guess {
        Some(guess) => println!("找到匹配的guess: {guess}"),
        None => println!("没有找到匹配的guess"),
    }
}

buffer_unordered会同时启动所有任务,find_map在找到第一个匹配结果后立即终止,剩余未完成的任务会被自动取消。

方案二:循环使用select_all

select_all可以从一组Future中选出第一个完成的任务,循环判断直到找到目标:

use futures::future::{select_all, FutureExt};
use rand::Rng;
use tokio;

async fn async_function(guess: u8) -> (u8, bool) {
    let random_wait = rand::thread_rng().gen_range(0..2);
    tokio::time::sleep(tokio::time::Duration::from_secs(random_wait)).await;
    println!("Running guess {guess}");
    (guess, guess == 231)
}

#[tokio::main]
async fn main() {
    let mut futures: Vec<_> = (0..=255)
        .map(|guess| async_function(guess).boxed())
        .collect();

    let mut found_guess = None;
    while !futures.is_empty() {
        let ((g, is_match), _, remaining) = select_all(futures).await;
        if is_match {
            found_guess = Some(g);
            break;
        }
        futures = remaining;
    }

    match found_guess {
        Some(guess) => println!("找到匹配的guess: {guess}"),
        None => println!("没有找到匹配的guess"),
    }
}

每次select_all返回第一个完成的任务结果,判断是否匹配,匹配则记录结果并退出循环,否则继续处理剩余任务。

关键注意事项

原代码中使用std::thread::sleep会阻塞整个线程,导致异步任务无法真正并发,必须替换为异步睡眠(如tokio::time::sleep或async_std::task::sleep),才能让多个任务同时执行。

内容的提问来源于stack exchange,提问作者dangarbri

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 07:15:33