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

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>>`

原因分析

  1. 核心问题:Future不满足Send trait
    async_trait默认要求异步方法返回的Future实现Send(因为它会将Future装箱为Pin<Box<dyn Future + Send>>)。而scraper库的Html迭代器(如Select类型)内部包含Cell,Cell未实现Sync,导致持有该迭代器的整个Future无法安全在线程间传递。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 08:43:17