我正在嘗試撰寫一個程式,該程式采用 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;
}
修改 to 的值.buffer_unordered(40)不會1導致速度差異,這證實了請求不是異步的事實。很明顯,程式在觀察列印輸出時分別等待每個請求的回應。
為了使我的請求異步發送,我需要修改什么?對我來說,我可以控制并發性也很重要,.buffer_unordered()并且結果被收集到一個向量中,而不是列印出來,所以最后洗掉 .collect() 是行不通的。
uj5u.com熱心網友回復:
首先,這條線不好:
.map(|elem| futures::future::ready(elem))
這基本上是強迫,如果你想繼續,我們需要能夠得到一個elem: (url, status, body)結果。
有了這個,當 buffer_unordered() 想要一個一個地緩沖它們時,future 就會一個一個地完成。buffer_unordered 不能神奇地撤銷之前發生的事情。
其次,一旦你洗掉了這一行,你會看到 buffer_unordered 需要 Future 專案(能夠先對它們進行分組,然后并行運行),但你是(url, status, body)從 filter_map 給它的。
因此,您可以使用普通的“map”代替 filter_map 將 url 映射到異步塊(即期貨):
let results = stream::iter(urls)
.map(|url| async { // not run anything yet at this point,
// just convert url -> Future
let response = (&client)
.get(&url)
.send()
.await
.ok()?; // Result<T, E> => Option<T> => T
let status = response.status();
let body = response
.text()
.await
.ok()?; // Result<T, E> => Option<T> => T
println!("{} ready", url);
Some((url, status, body))
})
.buffer_unordered(40) // get up to 40 futures (async blocks)
// and run them in parallel
...
你可以讓這個更明確,首先使 Futures 急切,然后在流中懶惰地運行它們:
let futures = urls.into_iter().map(|url| async { ... })
let results = stream::iter(futures)
.buffer_unordered(40)
...
如果你想過濾掉錯誤,你可以在 buffer_unordered 之后這樣做:
.buffer_unordered(40)
.filter_map(|result| result)
...
轉載請註明出處,本文鏈接:https://www.uj5u.com/qukuanlian/417775.html
標籤:
上一篇:如何遍歷陣列并獲取每個元素的資料?-使用貓鼬和node.js
下一篇:角度異步請求未及時加載
