【问题标题】:How to avoid incurring in lifetime errors while filter with an `async` predicate?如何避免在使用“异步”谓词进行过滤时发生生命周期错误?
【发布时间】:2022-02-04 19:25:04
【问题描述】:

使用 async 谓词过滤值列表会使 Rust 抱怨生命周期。即使集合是awaited,这意味着谓词不会超过过滤值,Rust 仍然持怀疑态度。

下面是playground here 的完整复制。请注意,它会过滤我们宁愿通过引用传递的非复制结构,而不是我们可以复制并忘记而不会产生开销的简单值。

use futures::stream::iter;
use futures::StreamExt;

#[derive(Debug)]
struct Foo {
    bar: usize
}

impl Foo {
    fn new(bar: usize) -> Self {
        Self {
            bar
        }
    }
}

#[tokio::main]
async fn main() {
    let arr = vec![
      Foo::new(0),
      Foo::new(1),
      Foo::new(2)
    ];
    
    let filtered = iter(arr)
      .filter(|f| async {compute_baz(f).await > 0})
      .collect::<Vec<_>>()
      .await; 

    // should print Foo{bar:1} and Foo{bar:2}
    println!("{:?}", filtered) 
}

async fn compute_baz(foo: &Foo) -> usize {
    // ...do lengthy task...
    
    foo.bar
}

更新

正如@Ceasar 在下面指出的,异步函数不是并行运行的,可以这样做吗?

我正在尝试做类似的事情:

let filter_mask = join_all(items.map(predicate)); 
let filtered = items.filter(|i| filter_mask[i]).collect::<Vec<_>>();

没有杂乱。

【问题讨论】:

    标签: rust async-await


    【解决方案1】:

    一个简单的解决方法是避免关闭:

    let mut filtered = vec![];
    for f in arr.iter() {
        if compute_baz(f).await > 0 {
            filtered.push(f);
        }
    }
    

    【讨论】:

    • 这不是按顺序过滤的吗?我正在尝试并行运行.await,但现在我想到了stream 可能会转换为您的版本。我正在考虑做类似let filter_mask = join_all(items.map(predicate); let filtered = items.filter(|i| filter_mask[i].collect() 之类的事情
    • 我是这么认为的,但并没有真正考虑过。只是为了确保:Playground。如果你希望它是并行的,你需要buffer_unordered,我想。但是你又回到了你一生的问题。
    • @doplumi 您可能应该编辑您的问题以包含该要求,以便有人可以正确回答。因为我有gone astray...
    • 现在看,不知道buffer_unordered
    猜你喜欢
    • 1970-01-01
    • 2023-01-09
    • 1970-01-01
    • 2011-01-25
    • 1970-01-01
    • 2014-09-10
    • 1970-01-01
    • 1970-01-01
    • 2020-09-11
    相关资源
    最近更新 更多