如何在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
相关产品推荐
相关产品推荐

