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

Rust如何从并发运行的for循环中收集多个返回结果?

解决方案

错误原因

你之前的写法无法运行主要有两个核心问题:

  1. 多个并发异步任务同时尝试可变借用外部的return_values向量,违反了Rust的所有权规则,编译器会直接报错。
  2. 代码中存在变量名笔误,把input写成了未定义的port。

方案1:全量收集结果后统一写入

适合返回结果量不大,不会占用过多内存的场景,实现最简单:

首先修改do_the_hard_job的函数签名,将符合条件的结果返回而非直接打印:

async fn do_the_hard_job(my_input: u16) -> Option<u16> {
    // 原有耗时执行逻辑
    if my_condition {
        Some(my_input)
    } else {
        None
    }
}

主逻辑代码:

use futures::StreamExt;
use tokio::fs::File;
use tokio::io::AsyncWriteExt;

#[tokio::main]
async fn main() -> std::io::Result<()> {
    // 生成输入流
    let inputs = futures::stream::iter(1..=u16::MAX);

    // 并发执行所有任务,筛选符合条件的结果
    let results: Vec<u16> = inputs
        .map(|input| async move { do_the_hard_job(input).await })
        .buffer_unordered(0) // 0表示不限制并发数,可根据硬件调整为固定值如1024
        .filter_map(|res| async move { res })
        .collect()
        .await;

    // 统一写入目标文件
    let mut file = File::create("output.txt").await?;
    for val in results {
        writeln!(&mut file, "{}", val)?;
    }
    file.flush().await?;

    Ok(())
}

方案2:通道异步写入(内存更友好)

如果结果量极大,不适合全量存放在内存中,可以用多生产者单消费者通道,所有并发任务只负责发送结果,单独开一个任务独占写文件,完全避免并发写冲突:

use futures::StreamExt;
use tokio::sync::mpsc;
use tokio::fs::File;
use tokio::io::AsyncWriteExt;

#[tokio::main]
async fn main() -> std::io::Result<()> {
    let inputs = futures::stream::iter(1..=u16::MAX);
    // 创建带缓存的通道,缓存大小可按需调整
    let (tx, mut rx) = mpsc::channel::<u16>(1024);

    // 启动独立的写文件任务,全程独占文件对象
    let write_task = tokio::spawn(async move {
        let mut file = File::create("output.txt").await?;
        while let Some(val) = rx.recv().await {
            writeln!(&mut file, "{}", val)?;
        }
        file.flush().await?;
        Ok::<(), std::io::Error>(())
    });

    // 并发执行所有任务,符合条件的结果发送到通道
    inputs.for_each_concurrent(0, |input| {
        let tx = tx.clone();
        async move {
            if let Some(res) = do_the_hard_job(input).await {
                let _ = tx.send(res).await;
            }
        }
    }).await;

    // 所有任务执行完成,关闭发送端触发写任务退出
    drop(tx);
    // 等待写任务执行完成
    write_task.await??;

    Ok(())
}

注意事项

  • 如果使用async-std运行时,替换对应异步IO、通道的API即可,核心逻辑完全一致。
  • 并发数不要无脑设置为无限制,CPU密集型任务建议设置为CPU核心数的1~2倍,IO密集型任务可以适当调大。
  • 如果需要输出顺序和输入顺序完全一致,将buffer_unordered替换为buffered即可。

内容的提问来源于stack exchange,提问作者Kürşat Kobya

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 18:27:07