我有一个Vec我想同时执行的 future(但不一定是并行的)。基本上,我正在寻找某种select类似于tokio::select!但需要 future 集合的函数,或者相反,一种类似于futures::join_all但在第一个 future 完成后返回的函数。
另一个要求是,一旦 future 完成,我可能想向Vec.
有了这样的函数,我的代码大致如下所示:
use std::future::Future;
use std::time::Duration;
use tokio::time::sleep;
async fn wait(millis: u64) -> u64 {
sleep(Duration::from_millis(millis)).await;
millis
}
// This pseudo-implementation simply removes the last
// future and awaits it. I'm looking for something that
// instead polls all futures until one is finished, then
// removes that future from the Vec and returns it.
async fn select<F, O>(futures: &mut Vec<F>) -> O
where
F: Future<Output=O>
{
let future = futures.pop().unwrap();
future.await
}
#[tokio::main]
async fn main() {
let mut futures = vec![
wait(500),
wait(300),
wait(100),
wait(200),
];
while !futures.is_empty() {
let finished = select(&mut futures).await;
println!("Waited {}ms", finished);
if some_condition() {
futures.push(wait(200));
}
}
}
Run Code Online (Sandbox Code Playgroud)
Flo*_*ker 20
这正是它futures::stream::FuturesUnordered的用途(我通过查看 的来源发现了这一点StreamExt::for_each_concurrent):
use futures::{stream::FuturesUnordered, StreamExt};
use std::time::Duration;
use tokio::time::{sleep, Instant};
async fn wait(millis: u64) -> u64 {
sleep(Duration::from_millis(millis)).await;
millis
}
#[tokio::main]
async fn main() {
let mut futures = FuturesUnordered::new();
futures.push(wait(500));
futures.push(wait(300));
futures.push(wait(100));
futures.push(wait(200));
let start_time = Instant::now();
let mut num_added = 0;
while let Some(wait_time) = futures.next().await {
println!("Waited {}ms", wait_time);
if num_added < 3 {
num_added += 1;
futures.push(wait(200));
}
}
println!("Completed all work in {}ms", start_time.elapsed().as_millis());
}
Run Code Online (Sandbox Code Playgroud)
(游乐场)
如果您使用 Tokio,请注意:正如@Bryan Larsen在评论中指出的那样,与 Tokio 结合时存在性能问题的风险FuturesUnordered。本文包含更多详细信息,并表示该问题应该在futures板条箱的最新版本(0.3.19及更高版本)中得到解决。尽管如此,Tokio 的用户还是使用 Tokio 的JoinSet. 与上面相同的示例如下所示:
use std::time::Duration;
use tokio::task::JoinSet;
use tokio::time::{sleep, Instant};
async fn wait(millis: u64) -> u64 {
sleep(Duration::from_millis(millis)).await;
millis
}
#[tokio::main]
async fn main() {
let mut futures = JoinSet::new();
futures.spawn(wait(500));
futures.spawn(wait(300));
futures.spawn(wait(100));
futures.spawn(wait(200));
let start_time = Instant::now();
let mut num_added = 0;
while let Some(result) = futures.join_next().await {
let wait_time = result.unwrap();
println!("Waited {}ms", wait_time);
if num_added < 3 {
num_added += 1;
futures.spawn(wait(200));
}
}
println!(
"Completed all work in {}ms",
start_time.elapsed().as_millis()
);
}
Run Code Online (Sandbox Code Playgroud)
(游乐场)
| 归档时间: |
|
| 查看次数: |
8442 次 |
| 最近记录: |