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

如何在Rust中高效将十万级RPC调用结果写入SQLite数据库

Rust异步批量RPC+SQLite并行读写优化方案

首先,并行RPC调用+批量写入的方案是可行的,但得避开SQLite的锁机制坑。SQLite本身是文件级锁,默认模式下写操作会独占锁,要是多个线程同时写,会触发频繁的锁竞争,反而拖慢速度。所以正确的姿势是:多线程并行发RPC请求,单线程负责批量写入数据库,这样既利用了并行提升RPC效率,又避免了数据库的竞态问题。

要让初始化写入的同时服务器能正常读数据,必须开启SQLite的WAL(Write-Ahead Log)模式。开启WAL后,写操作不会阻塞读操作,读请求可以直接读取旧版本的快照数据,完美解决你担心的读写冲突问题。

核心实现步骤

1. 并行发起RPC请求,控制并发数

用tokio的buffer_unordered控制并发请求数量,避免压垮目标RPC服务;同时复用reqwest的Client,减少连接建立的开销。

2. 用通道收集结果,单线程批量写入

用tokio的mpsc通道传递RPC结果,单独开一个异步任务从通道取数据,攒够100条就用事务批量写入——事务能大幅提升批量写入效率,比单条插入快得多。

3. 开启SQLite WAL模式

初始化数据库连接时执行PRAGMA journal_mode=WAL;,这是实现读写并发的关键。

代码示例

use rusqlite::{Connection, Result};
use tokio::sync::mpsc;
use reqwest::Client;
use futures::stream::StreamExt;

// 替换成你的RPC结果类型
#[derive(Debug)]
struct RpcResult {
    col1: String,
    col2: i32,
}

async fn fetch_rpc(client: &Client) -> Result<RpcResult> {
    // 这里写你的RPC请求逻辑
    let res = client.get("https://your-rpc-endpoint.com/api")
        .send()
        .await?
        .json()
        .await?;
    Ok(res)
}

fn batch_insert(conn: &Connection, batch: &[RpcResult]) -> Result<()> {
    // 用事务批量插入,提升效率
    let tx = conn.transaction()?;
    for item in batch {
        tx.execute(
            "INSERT INTO rpc_results (col1, col2) VALUES (?1, ?2)",
            (&item.col1, &item.col2),
        )?;
    }
    tx.commit()?;
    Ok(())
}

async fn run_initialization() -> Result<()> {
    // 初始化数据库,开启WAL模式
    let mut db_conn = Connection::open("rpc_results.db")?;
    db_conn.execute("PRAGMA journal_mode=WAL;", [])?;
    // 确保表存在(如果需要的话)
    db_conn.execute(
        "CREATE TABLE IF NOT EXISTS rpc_results (col1 TEXT, col2 INTEGER)",
        [],
    )?;

    // 创建通道,缓冲区设为1000,避免RPC任务阻塞
    let (tx, mut rx) = mpsc::channel(1000);
    let rpc_client = Client::new();

    // 并行发起200000次RPC请求,控制并发数为50
    tokio::spawn(async move {
        let rpc_stream = tokio_stream::iter(0..200000)
            .map(|_| {
                let client = rpc_client.clone();
                let tx = tx.clone();
                async move {
                    match fetch_rpc(&client).await {
                        Ok(result) => {
                            if let Err(e) = tx.send(result).await {
                                eprintln!("Failed to send RPC result: {}", e);
                            }
                        }
                        Err(e) => eprintln!("RPC request failed: {}", e),
                    }
                }
            })
            .buffer_unordered(50); // 并发数可根据实际调整

        rpc_stream.collect::<()>().await;
    });

    // 批量写入任务
    let mut batch = Vec::with_capacity(100);
    while let Some(result) = rx.recv().await {
        batch.push(result);
        if batch.len() >= 100 {
            if let Err(e) = batch_insert(&db_conn, &batch) {
                eprintln!("Batch insert failed: {}", e);
            }
            batch.clear();
        }
    }

    // 处理剩余的不足100条的数据
    if !batch.is_empty() {
        batch_insert(&db_conn, &batch)?;
    }

    Ok(())
}

// 在服务器启动时启动初始化任务
// let init_handle = tokio::task::spawn(run_initialization());

关键注意事项

  • WAL模式必须开:不开的话,写操作会阻塞所有读请求,服务器在初始化期间无法正常读取数据。
  • 并发数要合理:RPC并发数不能太高,否则目标服务可能限流或崩溃,根据对方API文档调整。
  • 事务不能少:批量写入时一定要用事务,否则每次插入都提交,效率和单条插入没区别。
  • 错误处理要到位:RPC请求、通道发送、数据库写入都可能失败,要捕获错误并打印,避免整个初始化任务直接panic。

推荐的库组合

  • reqwest:异步HTTP客户端,成熟稳定,适合RPC请求。
  • tokio:异步运行时,提供任务调度、通道等核心工具。
  • rusqlite:SQLite的Rust绑定,同步操作足够高效;若需异步数据库操作,可使用tokio-rusqlite。
  • futures:配合tokio处理流式任务,比如buffer_unordered控制并发。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 22:20:36