如何在Rust中结合异步数据采集与Rayon线程池数据处理
优化异步数据采集与Rayon处理的重叠执行方案
我想通过让数据获取和处理重叠执行,来优化异步数据采集与Rayon数据处理的整合逻辑。目前我用常规异步代码从网站拉取大量页面,等全部拉取完成后,再用Rayon的par_iter执行CPU密集型任务。
由于每个拉取的页面彼此独立,不需要等所有页面都拉完就能开始转换处理,所以我觉得可以轻松实现处理与获取的重叠,避免等到最后一个页面完成才开始繁重的处理工作。
以下是我当前可运行的简化代码:
use rayon::prelude::*; use futures::{stream, StreamExt}; use reqwest::{Client, Result}; const CONCURRENT_REQUESTS: usize = usize::MAX; const MAX_PAGE: usize = 1000; #[tokio::main] async fn main() { // 获取服务器数据 let client = Client::new(); let bodies: Vec<Result<String>> = stream::iter(1..MAX_PAGE+1) .map(|page_number| { let client = &client; async move { client .get(format!("https://someurl?{page_number}")) .send() .await? .text() .await } }) .buffer_unordered(CONCURRENT_REQUESTS) .collect() .await; // 转换数据 let mut rows: Vec<MyRow> = bodies .par_iter() .filter_map(|body| body.as_ref().ok()) .map(|data| { let page = serde_json::from_str::<MyPage>(data).unwrap(); page.rows .iter() .map(|x| Row::new(x)) .collect::<Vec<MyRow>>() }) .flatten() .collect(); // 处理转换后的rows }
解决方案:实现拉取与处理的重叠执行
核心思路是每个页面异步拉取完成后,立即启动CPU密集型处理任务,而非等待所有请求结束再批量处理。结合Tokio的异步 runtime 和 Rayon 的线程池,可以高效实现这种重叠执行。
修改后的代码
use rayon::prelude::*; use futures::{stream, StreamExt}; use reqwest::{Client, Result}; use tokio::task::spawn_blocking; use num_cpus; // 调整并发请求数为合理值,避免被目标网站限流 const CONCURRENT_REQUESTS: usize = 32; const MAX_PAGE: usize = 1000; // 假设的结构体定义,根据实际业务替换 #[derive(Debug, serde::Deserialize)] struct MyPage { rows: Vec<RowData>, } #[derive(Debug, serde::Deserialize)] struct RowData { // 原始数据字段 } #[derive(Debug)] struct MyRow { // 转换后的业务字段 } impl MyRow { fn new(data: &RowData) -> Self { // 自定义数据转换逻辑 MyRow {} } } #[tokio::main] async fn main() { let client = Client::new(); // 初始化Rayon线程池,按CPU核心数设置线程数量 let rayon_pool = rayon::ThreadPoolBuilder::new() .num_threads(num_cpus::get()) .build() .unwrap(); // 异步拉取 + 并行处理,实现重叠执行 let rows: Vec<MyRow> = stream::iter(1..=MAX_PAGE) .map(|page_number| { let client = client.clone(); async move { // 异步拉取单页数据 let body = client .get(format!("https://someurl?page={page_number}")) .send() .await? .text() .await?; Ok(body) } }) .buffer_unordered(CONCURRENT_REQUESTS) // 过滤请求失败的结果 .filter_map(|result| async move { result.ok() }) // 对每个成功拉取的页面,提交到Rayon线程池处理 .map(move |body| { let pool = rayon_pool.clone(); // 用spawn_blocking包装Rayon任务,避免阻塞异步IO线程 spawn_blocking(move || { pool.install(|| { // CPU密集型处理:解析JSON + 转换行数据 let page: MyPage = serde_json::from_str(&body).unwrap(); page.rows.iter().map(MyRow::new).collect::<Vec<MyRow>>() }) }) }) // 控制同时处理的任务数,避免线程过载 .buffer_unordered(CONCURRENT_REQUESTS) // 过滤处理失败的结果 .filter_map(|result| async move { result.ok() }) // 合并所有处理后的行数据 .flatten() .collect() .await; // 后续业务逻辑处理 println!("共处理完成 {} 条数据", rows.len()); }
关键优化点说明
- 并发请求控制:将
CONCURRENT_REQUESTS从usize::MAX改为合理值(如32),避免发起过多请求触发目标网站的限流或反爬机制。 - 重叠执行逻辑:每个页面拉取完成后立即启动处理任务,彻底消除"等待所有请求完成再处理"的时间浪费,总耗时接近"最长请求时间 + 总处理时间的并行化耗时"。
- 线程资源隔离:用
spawn_blocking将CPU密集任务从Tokio的异步IO线程转移到阻塞线程池,再通过Rayon的自定义线程池执行处理逻辑,避免IO线程被阻塞,同时精准控制处理任务的并发度。 - 错误处理:保留了对请求失败、处理失败的过滤逻辑,确保最终结果仅包含成功处理的数据。
内容的提问来源于stack exchange,提问作者ahenshaw
相关产品推荐
相关产品推荐

