使用reqwest与Tokio实现批量URL异步请求时未并发执行的问题解决咨询
问题:Rust异步请求未并发执行,buffer_unordered不生效
我正在编写一个程序,该程序接收URL列表并使用异步请求获取每个URL的StatusCode和响应体。当前代码如下:
extern crate futures; use futures::{stream, StreamExt}; use reqwest::{Client as http, StatusCode}; use std::time::Duration; #[tokio::main] async fn main() { let client_builder = http::builder() .connect_timeout(Duration::from_secs(5)) .danger_accept_invalid_certs(true) .redirect(reqwest::redirect::Policy::none()) .timeout(Duration::from_secs(5)) .user_agent("Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/87.0.4280.88 Safari/537.36"); let client = client_builder.build().unwrap(); let mut urls = vec![]; for _ in 0 .. 100 { urls.push("https://google.com:443".to_string()); } let results = stream::iter(urls) .filter_map(|url| async { let response = (&client) .get(&url) .send() .await .ok()?; // Result<T, E> => Option<T> => T Some((url, response)) }) .filter_map(|(url, response)| async { let status = response.status(); let body = response .text() .await .ok()?; // Result<T, E> => Option<T> => T println!("{}", url); Some((url, status, body)) }) .map(|elem| futures::future::ready(elem)) .buffer_unordered(40) .collect::<Vec<(String, StatusCode, String)>>() .await; }
但实际运行时,请求并未以异步方式发送,而是逐个串行执行。将.buffer_unordered(40)修改为1后,程序运行速度没有差异,这进一步确认了请求未异步执行的问题。观察打印输出也能明显看到程序在逐个等待每个请求的响应。
我需要修改代码以实现请求的异步并发,同时要求能够通过.buffer_unordered()控制并发数,且最终结果需收集到Vector中,因此不能移除末尾的.collect()方法。请问需要如何修改代码?
解决方案:调整流处理逻辑,让
buffer_unordered作用于未执行的请求Future 你的问题根源在于当前的流处理链是串行执行异步逻辑,buffer_unordered并没有真正作用在请求的异步任务上。具体来说:
filter_map是StreamExt提供的方法,它会逐个处理流中的元素:等前一个元素的异步逻辑执行完成后,才会处理下一个URL。- 后面的
map(|elem| futures::future::ready(elem))只是把已经完成的结果包装成一个立即就绪的Future,buffer_unordered处理这种Future完全起不到并发的效果。
要实现真正的并发,我们需要把每个URL的完整请求逻辑封装成一个独立的异步任务(Future),然后让流产生这些未执行的Future,再用buffer_unordered来控制同时执行的任务数量。
修改后的完整代码
extern crate futures; use futures::{stream, StreamExt}; use reqwest::{Client as http, StatusCode}; use std::time::Duration; // 把单个URL的请求逻辑封装成独立的异步函数 async fn fetch_single_url(client: &http, url: String) -> Option<(String, StatusCode, String)> { // 发送请求 let response = client.get(&url).send().await.ok()?; let status = response.status(); // 获取响应体 let body = response.text().await.ok()?; println!("{}", url); Some((url, status, body)) } #[tokio::main] async fn main() { let client_builder = http::builder() .connect_timeout(Duration::from_secs(5)) .danger_accept_invalid_certs(true) .redirect(reqwest::redirect::Policy::none()) .timeout(Duration::from_secs(5)) .user_agent("Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/87.0.4280.88 Safari/537.36"); let client = client_builder.build().unwrap(); let mut urls = vec![]; for _ in 0 .. 100 { urls.push("https://google.com:443".to_string()); } let results = stream::iter(urls) // 把每个URL映射为一个未执行的请求Future .map(move |url| fetch_single_url(&client, url)) // 控制并发数:同时最多执行40个请求任务 .buffer_unordered(40) // 过滤掉请求失败的结果(对应fetch_single_url返回的None) .filter_map(|result| futures::future::ready(result)) // 收集所有成功的结果到Vector中 .collect::<Vec<(String, StatusCode, String)>>() .await; }
关键修改点说明
- 封装请求逻辑为独立异步函数:
fetch_single_url把单个URL的发送请求、获取状态码和响应体的逻辑打包成一个异步任务,返回Option表示请求是否成功。 - 流产生未执行的Future:
stream::iter(urls).map(...)会为每个URL生成一个未执行的请求Future,而不是立刻执行请求逻辑。 buffer_unordered作用于请求Future:此时buffer_unordered(40)会同时启动最多40个请求任务,真正实现异步并发。- 过滤失败结果:
filter_map用来过滤掉请求失败的None结果,只保留成功的请求数据。
现在你运行代码就能看到请求是并发执行的,调整buffer_unordered的参数也能明显看到并发数变化带来的速度差异。
内容的提问来源于stack exchange,提问作者iustin
相关产品推荐
相关产品推荐

