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

如何在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());
}

关键优化点说明

  1. 并发请求控制:将CONCURRENT_REQUESTS从usize::MAX改为合理值(如32),避免发起过多请求触发目标网站的限流或反爬机制。
  2. 重叠执行逻辑:每个页面拉取完成后立即启动处理任务,彻底消除"等待所有请求完成再处理"的时间浪费,总耗时接近"最长请求时间 + 总处理时间的并行化耗时"。
  3. 线程资源隔离:用spawn_blocking将CPU密集任务从Tokio的异步IO线程转移到阻塞线程池,再通过Rayon的自定义线程池执行处理逻辑,避免IO线程被阻塞,同时精准控制处理任务的并发度。
  4. 错误处理:保留了对请求失败、处理失败的过滤逻辑,确保最终结果仅包含成功处理的数据。

内容的提问来源于stack exchange,提问作者ahenshaw

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 18:18:50