如何在Rust的FuturesUnordered中忽略错误适配bar函数
如何忽略FuturesUnordered流中的错误并调用期望成功值的函数
我有如下代码,futures是FuturesUnordered类型,每个内部future返回std::io::Result<i32>。foo函数可以正常接收它,但bar函数期望的是内部future返回值为非Result类型的FuturesUnordered,直接传入会编译失败。我无法修改futures的初始化方式,请问如何忽略流中的错误,只保留成功的元素来调用bar函数?
原始代码
use futures::{ future, stream::{self, Stream, FuturesUnordered}, }; use tokio; use rand::Rng; fn foo(futures: FuturesUnordered<impl futures::Future<Output = std::io::Result<impl std::fmt::Binary>>>) {} fn bar(futures: FuturesUnordered<impl futures::Future<Output = impl std::fmt::Binary>>) {} #[tokio::main] async fn main() { let futures: FuturesUnordered<_> = (0..10).map(move |i| async move { let delay = core::time::Duration::from_secs(rand::thread_rng().gen_range(1..3)); tokio::time::sleep(delay).await; Ok::<i32, std::io::Error>(i) // 这一行不能修改 }).collect(); // 可以正常编译 foo(futures); // 编译失败 bar(futures); }
我的尝试(编译失败)
我参考了类似问题,但其中的stream::iter_ok已被弃用,自己尝试的代码无法正常运行:
use futures::{ future, stream::{self, Stream, FuturesUnordered}, StreamExt, }; use tokio; use rand::Rng; fn foo(futures: FuturesUnordered<impl futures::Future<Output = std::io::Result<impl std::fmt::Binary>>>) {} async fn bar(futures: FuturesUnordered<impl futures::Future<Output = impl std::fmt::Binary>>) { futures.for_each(|n| { async move { println!("Success on {:b}", n); } }).await } #[tokio::main] async fn main() { let futures: FuturesUnordered<_> = (0..10).map(move |i| async move { let delay = core::time::Duration::from_secs(rand::thread_rng().gen_range(1..3)); tokio::time::sleep(delay).await; Ok::<i32, std::io::Error>(i) }).collect(); let futures = futures .then(|r| future::ok(iter_ok::<_, ()>(r))) .flatten(); bar(futures).await; }
解决方案
要实现忽略错误、仅保留成功元素并调用bar的需求,我们可以通过以下步骤处理原FuturesUnordered:
- 利用
StreamExt提供的filter_map方法,过滤掉Result中的错误值,只保留成功的结果; - 将每个成功的结果包装成一个立即完成的future(使用
future::ready); - 将处理后的future收集到新的
FuturesUnordered中,即可传入bar函数。
正确代码示例
use futures::{ future, stream::{Stream, FuturesUnordered}, StreamExt, }; use tokio; use rand::Rng; fn foo(futures: FuturesUnordered<impl futures::Future<Output = std::io::Result<impl std::fmt::Binary>>>) {} async fn bar(futures: FuturesUnordered<impl futures::Future<Output = impl std::fmt::Binary>>) { futures.for_each(|n| async move { println!("Success on {:b}", n); }).await; } #[tokio::main] async fn main() { let futures: FuturesUnordered<_> = (0..10).map(move |i| async move { let delay = core::time::Duration::from_secs(rand::thread_rng().gen_range(1..3)); tokio::time::sleep(delay).await; Ok::<i32, std::io::Error>(i) // 这一行不能修改 }).collect(); // 处理错误,仅保留成功的元素并转换为符合bar要求的FuturesUnordered let processed_futures: FuturesUnordered<_> = futures .filter_map(|result| async move { // 过滤掉错误,只保留Ok的值 result.ok() }) .map(|val| { // 将成功的值包装为立即完成的future future::ready(val) }) .collect(); // 现在可以正常调用bar bar(processed_futures).await; }
代码说明
filter_map方法会异步处理每个Result:如果是Ok(val)则返回Some(val),如果是Err(_)则返回None,None会被自动过滤掉;map(future::ready)将成功的值转换为一个立即完成的future,确保新的FuturesUnordered中的每个元素都符合bar函数要求的返回类型;- 最终收集得到的
processed_futures可以直接传入bar函数,实现仅处理成功元素的需求。
内容的提问来源于stack exchange,提问作者Kevin
相关产品推荐
相关产品推荐

