Rust异步代码报错:future无法安全跨线程发送的原因与解决
问题:Tokio异步函数中
await调用报错,block_on却可正常运行 我基于Tokio运行时编写了包含循环的异步函数,在collect方法中用await调用get_single_item时触发编译错误,但替换为executor::block_on就能正常编译。想了解两种写法的差异,以及如何用await实现正常运行。
代码示例
use std::time::Duration; use async_trait::async_trait; use scraper::{ Selector, Html }; use urlencoding::encode; use crate::{ currency::{ Currency, CurrencyConverter }, comparer::{ executor::{ BaseSiteInformation, Executor }, model::{ Item, ItemPrice, Response, SiteResponse }, error::Error, }, }; use super::common; pub struct Site { pub info: BaseSiteInformation, } impl Site{ async fn get_single_item(&self, url: String, currency: Currency) -> Item { let client = reqwest::Client ::builder() .cookie_store(true) .timeout(Duration::from_secs(20)) .build() .unwrap(); let inital_response = client .get(url) .header( "User-Agent", "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/119.0.0.0 Safari/537.36" ) .header("Sec-Ch-Ua-Platform", "\"Windows\"") .header("Sec-Ch-Ua-Mobile", "?0") .header( "Sec-Ch-Ua", "\"Google Chrome\";v=\"119\", \"Chromium\";v=\"119\", \"Not?A_Brand\";v=\"24" ) .send().await .unwrap(); let inital_response = Html::parse_document(&inital_response.text().await.unwrap()); let selector = Selector::parse("div#list").unwrap(); let item_list = inital_response.select(&selector); let mut extracted_prices: Vec<String> = vec![]; let price_selector = Selector::parse("span").unwrap(); let prices = item_list.last().unwrap().select(&price_selector); for price in prices { if price.attr("divprop").is_some() { let price = price.attr("content").unwrap().to_string(); let value = CurrencyConverter::convert( price.parse::<f64>().unwrap(), currency, self.info.currency ); extracted_prices.push(value.to_string()); } } let link_selector = Selector::parse("a.overlay").unwrap(); let item_list = inital_response.select(&selector); let links = item_list.last().unwrap().select(&link_selector); let mut collected_links = vec![]; for link in links { if link.attr("href").is_some() { let link = link.attr("href").unwrap().to_string(); let actual_url = common::get_resolved_url(link, &client).await; collected_links.push(actual_url); } } let name_selector = Selector::parse("div.top h1.visible-xs").unwrap(); let name = inital_response .select(&name_selector) .last() .unwrap() .inner_html() .to_string() .trim() .to_string(); let mut item_prices = vec![]; let mut i = 0; while i < collected_links.len() { item_prices.push(ItemPrice { link: collected_links.get(i).unwrap().to_string(), price: extracted_prices.get(i).unwrap().to_string(), }); i += 1; } Item { item: name, prices: item_prices } } } #[async_trait] impl Executor for Site { async fn collect(&self, prompt: String, currency: Currency) -> Result<Response, Error> { let client = reqwest::Client::builder().timeout(Duration::from_secs(20)).build().unwrap(); let inital_response = client .get(self.info.base_url.clone().replace("{PROMPT}", &encode(&prompt))) .header( "User-Agent", "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/119.0.0.0 Safari/537.36" ) .header("Sec-Ch-Ua-Platform", "\"Windows\"") .header("Sec-Ch-Ua-Mobile", "?0") .header( "Sec-Ch-Ua", "\"Google Chrome\";v=\"119\", \"Chromium\";v=\"119\", \"Not?A_Brand\";v=\"24" ) .send().await?; let inital_response = inital_response.text().await?; let inital_response = serde_json::from_str::<SiteResponse>(&inital_response)?; let initial_items = Html::parse_document(&inital_response.dom); let selector = Selector::parse("ul.products li a").unwrap(); let initial_items = initial_items.select(&selector); let mut items = vec![]; for item in initial_items { let local_item = item.clone(); let item_link = local_item.attr("href").unwrap().to_string(); let collected = self.get_single_item(item_link.clone(), currency).await; items.push(collected); } Ok(Response { items }) } }
报错信息
future cannot be sent between threads safely within `ego_tree::Node<Node>`, the trait `Sync` is not implemented for `Cell<NonZeroUsize>` if you want to do aliasing and mutation between multiple threads, use `std::sync::RwLock` required for the cast from `Pin<Box<{async block@src\carer\executor.rs:138:92: 178:6}>>` to `Pin<Box<dyn futures::Future<Output = Result<model::Response, comparer::error::Error>> + std::marker::Send>>`
原因分析
- 核心问题:Future不满足
Sendtraitasync_trait默认要求异步方法返回的Future实现Send(因为它会将Future装箱为Pin<Box<dyn Future + Send>>)。而scraper库的Html迭代器(如Select类型)内部包含Cell,Cell未实现Sync,导致持有该迭代器的整个Future无法安全在线程间传递。 await与block_on的差异await:异步等待时,Tokio可能会将当前任务切换到其他线程执行,因此要求Future必须是Send的。block_on:同步阻塞执行,直接在当前线程运行Future,无需线程切换,因此不要求Future满足Send,但会阻塞线程,失去异步编程的并发优势。
解决方案
方案1:提前消费迭代器,消除非Send状态
在collect方法中,先把所有商品链接收集到Vec中,完全消费掉scraper的迭代器,避免其非Send状态跨await点存在:
#[async_trait] impl Executor for Site { async fn collect(&self, prompt: String, currency: Currency) -> Result<Response, Error> { let client = reqwest::Client::builder().timeout(Duration::from_secs(20)).build().unwrap(); let initial_response = client .get(self.info.base_url.clone().replace("{PROMPT}", &encode(&prompt))) .header( "User-Agent", "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/119.0.0.0 Safari/537.36" ) .header("Sec-Ch-Ua-Platform", "\"Windows\"") .header("Sec-Ch-Ua-Mobile", "?0") .header( "Sec-Ch-Ua", "\"Google Chrome\";v=\"119\", \"Chromium\";v=\"119\", \"Not?A_Brand\";v=\"24" ) .send().await?; let initial_response_text = initial_response.text().await?; let site_response = serde_json::from_str::<SiteResponse>(&initial_response_text)?; let initial_items = Html::parse_document(&site_response.dom); let selector = Selector::parse("ul.products li a").unwrap(); // 提前收集所有链接,消费迭代器,避免持有非Send状态 let item_links: Vec<String> = initial_items .select(&selector) .filter_map(|item| item.attr("href").map(String::from)) .collect(); let mut items = vec![]; for link in item_links { let collected = self.get_single_item(link, currency).await; items.push(collected); } Ok(Response { items }) } }
方案2:使用非Send版本的async_trait
在#[async_trait]后添加(?Send),允许方法返回非Send的Future,但会限制Future只能在当前线程执行,可能影响调度灵活性:
#[async_trait(?Send)] impl Executor for Site { // 其余代码与原实现一致 }
优先推荐方案1,方案2仅适合无需跨线程调度的场景。
内容的提问来源于stack exchange,提问作者geo10
相关产品推荐
相关产品推荐

