Rust如何从并发运行的for循环中收集多个返回结果?
解决方案
错误原因
你之前的写法无法运行主要有两个核心问题:
- 多个并发异步任务同时尝试可变借用外部的
return_values向量,违反了Rust的所有权规则,编译器会直接报错。 - 代码中存在变量名笔误,把
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
相关产品推荐
相关产品推荐

