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

Rust多线程并行处理百万JSON对象的性能优化求助

百万级交易数据异步处理优化:解决CPU近乎空闲、耗时过长问题

我有一份包含100万个JSON对象的交易数据,需要尽可能高效地处理。目前已将数据拆分为每份1000条的块,verify_transactions是异步函数(因其他场景也需调用),仅在检查是否需提前终止时存在一处阻塞逻辑,其余逻辑均为非阻塞。交易处理本身逻辑简单:每个块拆分为单条交易,执行哈希校验、时间戳校验等操作。

但当前代码运行耗时超16小时,CPU绝大多数时间处于空闲状态,使用率从未超过2%。百万级数据量虽大,但这种CPU空闲却耗时极长的情况明显不符合预期,我刚接触Rust,不知道该如何优化来充分利用CPU资源。

原处理代码

let full_results: Vec<String> = (0..num_threads)
    .into_par_iter()
    .flat_map(|thread_id| {
        let mut local_results: Vec<String> = Vec::new();
        rt.block_on(async {
            let chunk_size = 1000;
            let total_chunks = transactions.len() as u32 / chunk_size;
            let chunks_per_thread = total_chunks / num_threads;
            let extra_chunks = total_chunks % num_threads;
            let start_index = thread_id * chunks_per_thread + std::cmp::min(thread_id, extra_chunks);
            let end_index = (thread_id + 1) * chunks_per_thread + std::cmp::min(thread_id + 1, extra_chunks);
            for chunk_index in start_index..end_index {
                let chunk_start = (chunk_index * chunk_size) as usize;
                let chunk_end = ((chunk_index + 1) * chunk_size) as usize;
                let chunk_end = std::cmp::min(((chunk_index + 1) * chunk_size) as usize, transactions.len());
                if let Some(chunk) = transactions.get(chunk_start..chunk_end) {
                    let owned_chunk: Vec<Value> = chunk.iter().cloned().collect();
                    let result = verify_transactions(owned_chunk, db.clone(), stop_flag.clone()).await;
                    for individual_string in result {
                        println!("{:?}", individual_string);
                        local_results.extend(individual_string);
                    }
                }
            };
        });
        local_results
    }).collect();
    verified_transactions.blocking_send(full_results);
});

verify_transactions函数实现

// 交易校验
async fn verify_transactions(parsed_data: Vec<Value>, db: Db, stop_flag: Arc<Mutex<bool>>) -> Result<Vec<String>, String> {
    let mut results: Vec<String> = Vec::new();
    let mut error_flag = false;
    if parsed_data.len() == 1 {
        let transaction = &parsed_data[0];
        if *stop_flag.lock().await {
            return Err("收到终止信号.".to_string());
        }
        // 单条交易直接处理
        match verify_transaction(transaction.clone(), &db).await {
            Ok(result) => {
                // 存入verify_transaction返回的结果
                results.push(result);
            }
            Err(error) => {
                // 处理错误,设置标志并退出
                println!("错误: {}", error);
                *stop_flag.lock().await = true;
                return Err("交易处理过程中发生错误.".to_string());
            }
        }
    } else {
        for transaction in &parsed_data {
            if *stop_flag.lock().await {
                return Err("收到终止信号.".to_string());
            }
            match verify_transaction(transaction.clone(), &db).await {
                Ok(result) => {
                    println!("结果: {}", result);
                    // 存入verify_transaction返回的结果
                    results.push(result);
                }
                Err(error) => {
                    // 处理错误,设置标志并退出
                    println!("错误: {}", error);
                    *stop_flag.lock().await = true;
                    return Err("交易处理过程中发生错误.".to_string());
                }
            }
        }
    }
    Ok(results)
}

问题分析与优化方案

核心问题点

  1. 异步代码被同步阻塞,并发能力完全浪费
    用rayon并行迭代器创建多线程后,每个线程又通过rt.block_on将异步代码转为同步阻塞执行,且verify_transactions内部串行处理单条交易,相当于每个线程仅处理一个异步任务,CPU无法并行利用。
  2. 锁的滥用导致额外阻塞
    stop_flag使用Arc<Mutex<bool>>,且每次处理单条交易都要加锁检查,频繁的锁竞争会拖慢线程执行速度。
  3. 不必要的克隆操作消耗资源
    代码中多次出现transaction.clone()、db.clone()等操作,若Value或Db是大对象,会带来大量内存和CPU开销。
  4. 块内交易串行处理
    verify_transactions对块内交易逐个串行处理,未利用异步runtime的并发调度能力。
  5. 频繁打印拖慢速度
    大量println!操作会带来IO开销,严重影响处理效率。

具体优化步骤

1. 用异步任务并发替代线程+同步阻塞

放弃rayon并行迭代器,直接用异步runtime(如tokio)调度并发任务,让runtime自主管理线程和任务分配:

use tokio::task;

// 拆分交易为1000条的块切片
let chunks: Vec<&[Value]> = transactions.chunks(1000).collect();
let mut tasks = Vec::new();

// 为每个块创建异步任务
for chunk in chunks {
    let chunk_owned = chunk.to_vec();
    let db_clone = db.clone();
    let stop_flag_clone = stop_flag.clone();
    
    tasks.push(task::spawn(async move {
        verify_transactions(chunk_owned, db_clone, stop_flag_clone).await
    }));
}

// 收集所有任务结果
let mut full_results = Vec::new();
for task in tasks {
    match task.await.unwrap() {
        Ok(mut results) => full_results.append(&mut results),
        Err(e) => {
            *stop_flag.lock().await = true;
            eprintln!("任务失败: {}", e);
            break;
        }
    }
}

verified_transactions.send(full_results).await.unwrap();

2. 优化verify_transactions,并发处理块内交易

将块内单条交易改为异步并发处理,充分利用异步runtime的调度能力:

async fn verify_transactions(parsed_data: Vec<Value>, db: Db, stop_flag: Arc<Mutex<bool>>) -> Result<Vec<String>, String> {
    let mut tasks = Vec::new();
    let db_arc = Arc::new(db); // 用Arc包裹Db,避免频繁克隆

    for transaction in parsed_data {
        if *stop_flag.lock().await {
            return Err("收到终止信号.".to_string());
        }
        
        let db_clone = Arc::clone(&db_arc);
        let stop_flag_clone = Arc::clone(&stop_flag);
        
        tasks.push(task::spawn(async move {
            if *stop_flag_clone.lock().await {
                return Err("收到终止信号.".to_string());
            }
            verify_transaction(transaction, &db_clone).await
        }));
    }

    let mut results = Vec::new();
    for task in tasks {
        match task.await.unwrap() {
            Ok(result) => results.push(result),
            Err(error) => {
                eprintln!("错误: {}", error);
                *stop_flag.lock().await = true;
                return Err("交易处理过程中发生错误.".to_string());
            }
        }
    }

    Ok(results)
}

3. 用原子类型替换锁,减少竞争

将Arc<Mutex<bool>>替换为Arc<AtomicBool>,原子操作的开销远低于锁:

use std::sync::atomic::{AtomicBool, Ordering};

// 初始化终止标志
let stop_flag = Arc::new(AtomicBool::new(false));

// 检查终止标志
if stop_flag.load(Ordering::SeqCst) {
    return Err("收到终止信号.".to_string());
}

// 设置终止标志
stop_flag.store(true, Ordering::SeqCst);

4. 减少不必要的克隆

  • 修改verify_transaction参数为&Value,避免克隆交易对象;
  • 用Arc<Db>共享数据库连接,避免频繁克隆Db;
  • 尽量使用切片引用而非提前克隆整个数据块。

5. 移除或控制打印操作

注释掉所有println!,或改用日志库(如tracing)并设置合适的日志级别,避免IO拖慢处理速度。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 19:40:14