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

如何在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:

  1. 利用StreamExt提供的filter_map方法,过滤掉Result中的错误值,只保留成功的结果;
  2. 将每个成功的结果包装成一个立即完成的future(使用future::ready);
  3. 将处理后的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 02:30:51