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

如何在Rust中通过多线程高效写入数据库并处理锁冲突?

问题分析与解决方案

核心问题拆解

你的代码存在几个关键问题,直接导致多线程写入效率不升反降,还频繁出现锁panic:

  • 连接频繁创建销毁:每次循环都新建数据库连接,连接本身有初始化开销,且多线程下大量连接竞争数据库锁,反而拖慢整体速度。
  • 线程数量过载:创建1000个线程,远超过CPU核心数,线程上下文切换的开销会完全抵消多线程带来的并行收益。
  • 单条插入效率极低:单条插入数据库的IO开销占比极高,没有利用批量插入减少IO次数的优势。
  • 错误处理缺失:unwrap()遇到数据库锁失败会直接panic,且insert的错误被忽略,既导致程序崩溃又可能丢失数据。
  • 原子变量竞态:多次load原子变量可能导致多个线程重复处理同一计数,或者大量线程同时触发打印,造成输出混乱。

针对性优化方案

1. 复用数据库连接

不要每次循环创建新连接,每个线程初始化一个连接并复用整个任务周期,减少连接开销和锁竞争。如果有条件,使用连接池进一步优化连接复用。

2. 控制线程数量

线程数建议设置为CPU核心数的1-2倍(比如用num_cpus::get()获取核心数),避免过多上下文切换消耗资源。

3. 批量插入数据

将多条数据打包成批量插入,这是提升写入速度最关键的优化——批量操作能大幅减少数据库IO次数,把单条插入的高开销摊薄到多条数据上。

4. 处理锁冲突与错误

  • 替换unwrap()为显式错误处理,遇到锁冲突时用指数退避策略重试。
  • 不要忽略insert的错误,确保数据写入成功或记录失败情况以便后续处理。

5. 提前分配任务范围

避免多个线程竞争同一个原子计数,提前给每个线程分配固定的写入区间,减少原子操作的竞态开销。

修改后的示例代码

#[cfg(test)]
mod tests {
    use std::{
        thread,
        time::Duration,
    };
    use common::store::StoreConnection;
    use num_cpus;
    use reports::bat_osc::BatOscData;

    const TEST_DB: &str = "/tmp/test_database.db";
    const TOTAL_ENTRIES: u32 = 7_000_000;
    const BATCH_SIZE: usize = 1000; // 可根据数据库性能调整批量大小

    // 带重试的批量插入逻辑
    fn batch_insert(store: &mut StoreConnection, batch: &[BatOscData]) -> Result<(), Box<dyn std::error::Error>> {
        let mut retries = 3;
        while retries > 0 {
            match store.insert_batch(batch) { // 假设StoreConnection支持批量插入,若不支持可在连接内循环单条插入
                Ok(_) => return Ok(()),
                Err(e) if e.to_string().contains("locked") => {
                    retries -= 1;
                    // 指数退避:重试间隔逐步拉长,避免持续竞争锁
                    thread::sleep(Duration::from_millis(100 * (4 - retries)));
                }
                Err(e) => return Err(e.into()),
            }
        }
        Err("Batch insert failed after 3 retries".into())
    }

    fn write_entries(start: u32, end: u32) {
        // 每个线程初始化一个连接,全程复用
        let mut store = StoreConnection::new(TEST_DB, reports::bat_osc::BAT_OSC_TABLE)
            .expect("Failed to create store connection")
            .writer(2592000 * 2)
            .expect("Failed to create writer");

        let mut batch = Vec::with_capacity(BATCH_SIZE);
        for i in start..end {
            batch.push(BatOscData { ..Default::default() });
            
            // 达到批量大小或到末尾时执行插入
            if batch.len() >= BATCH_SIZE || i == end - 1 {
                if let Err(e) = batch_insert(&mut store, &batch) {
                    eprintln!("Batch insert failed at {}: {}", i, e);
                }
                batch.clear();
                
                // 每个线程单独打印进度,避免混乱
                if i % 100_000 == 0 {
                    eprintln!("Thread processed {} entries", i - start + 1);
                }
            }
        }
    }

    #[test]
    fn write_a_lot() {
        let num_threads = num_cpus::get(); // 使用CPU核心数作为线程数
        let entries_per_thread = TOTAL_ENTRIES / num_threads as u32;
        let mut handles = vec![];

        // 给每个线程分配固定的写入区间
        for i in 0..num_threads {
            let start = i as u32 * entries_per_thread;
            let end = if i == num_threads - 1 {
                TOTAL_ENTRIES // 最后一个线程处理剩余条目
            } else {
                (i + 1) as u32 * entries_per_thread
            };

            let handle = thread::spawn(move || {
                write_entries(start, end);
            });
            handles.push(handle);
        }

        // 等待所有线程完成
        for handle in handles {
            handle.join().expect("Thread panicked");
        }
    }
}

额外优化建议

  • 数据库配置调整:如果用SQLite,开启PRAGMA journal_mode=WAL模式,该模式支持多读者单写者,比默认的DELETE模式并发写入性能提升明显。
  • 连接池集成:如果你的StoreConnection不支持连接池,可以用r2d2这类库实现连接池,进一步减少连接初始化开销。
  • 批量大小调优:BATCH_SIZE需要根据数据库类型和硬件调整,过大可能导致内存占用过高或事务超时,过小则无法发挥批量优势。
  • 错误日志完善:增加更详细的错误日志,方便排查锁冲突、磁盘空间不足等深层问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 07:45:17