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

在Rust中使用SQLX实现数据库异步IO是否可行?性能优化咨询

问题:Rust中用SQLX实现数据库异步IO及性能优化

我有一段将数据从其他存储迁移到PostgreSQL的代码,当前处理速率约200条/秒,性能很差。一开始用Python写,后来换成Rust,本以为瓶颈在Python,但换完后问题依旧。以下是简化后的示例代码:

use sqlx;
use chrono;
use futures;

fn main() -> Result<(), ()> {

    println!("connecting to Postgres...");

    let pool =
        futures::executor::block_on(
            sqlx::postgres::PgPoolOptions::new()
            .max_connections(5)
            .connect("postgresql://postgres:password@192.168.0.x/postgres")
        ).expect("failed to connect to Postgres");

    println!("connected to Postgres");

    
    let html_data = "<div>some long string read from file, ideally a few kb in size</div>";

    let datetime_start = chrono::offset::Utc::now();

    // insert html_data row
    let query_result =
            sqlx::query(
                "insert into html_data (html_data) values ($1)"
            )
            .bind(html_data)
            .fetch_optional(&pool);

    let _optional_row = futures::executor::block_on(query_result).expect("block_on failed");

    let datetime_stop = chrono::offset::Utc::now();
    let duration = datetime_stop - datetime_start;
    let duration = 1.0e+3 * duration.num_seconds() as f64 + 1.0e-9 * duration.subsec_nanos() as f64;
    println!("insert html_data: {duration} ms");
    
    Ok(())
}

这段代码的功能:

  • 建立PostgreSQL连接
  • 向包含主键id列和varchar类型html_data列的表插入单行数据
  • 所有数据库操作通过futures::executor::block_on转为同步执行

目前单条插入往返时间约10ms,导致速率上不去。我想改成异步执行,也就是发起请求后不等待响应,后续再处理结果,但不知道怎么处理sqlx::query返回的Future类型。请问在Rust中用SQLX实现数据库异步IO可行吗?该怎么实现?


解决方案

1. 异步IO完全可行,先改用异步运行时

SQLX本身就是为异步设计的,你当前用block_on把异步接口强行同步调用,完全浪费了异步的优势。首先要改用Rust的异步运行时,比如tokio(SQLX官方推荐)。

首先在Cargo.toml中添加依赖:

[dependencies]
sqlx = { version = "0.7", features = ["postgres", "tokio-rustls"] }
tokio = { version = "1.0", features = ["full"] }
chrono = "0.4"

然后把main函数改成异步的,用tokio的#[tokio::main]宏标记:

use sqlx::postgres::PgPool;
use chrono::Utc;
use tokio;

#[tokio::main]
async fn main() -> Result<(), sqlx::Error> {
    println!("connecting to Postgres...");

    let pool = PgPool::connect("postgresql://postgres:password@192.168.0.x/postgres")
        .await?;

    println!("connected to Postgres");

    let html_data = "<div>some long string read from file, ideally a few kb in size</div>";

    // 异步执行插入,无需block_on
    let start = Utc::now();
    sqlx::query("insert into html_data (html_data) values ($1)")
        .bind(html_data)
        .execute(&pool)
        .await?;
    let duration = Utc::now() - start;
    println!("insert html_data: {} ms", duration.num_milliseconds());

    Ok(())
}

2. 批量异步插入提升吞吐量

单条插入的往返延迟是硬伤,哪怕异步单条插,提升也有限。真正能大幅提升速率的是批量插入,结合异步并发处理:

方式一:用SQLX的query_batch

#[tokio::main]
async fn main() -> Result<(), sqlx::Error> {
    let pool = PgPool::connect("postgresql://postgres:password@192.168.0.x/postgres")
        .await?;

    // 模拟批量数据
    let batch_data: Vec<&str> = vec![
        "<div>data 1</div>",
        "<div>data 2</div>",
        // ... 更多数据,比如每次批量100条
    ];

    let start = Utc::now();
    // 批量插入,自动处理参数绑定
    sqlx::query("insert into html_data (html_data) values ($1)")
        .bind_all(batch_data)
        .execute(&pool)
        .await?;
    println!("batch insert took {} ms", Utc::now().signed_duration_since(start).num_milliseconds());

    Ok(())
}

方式二:并发异步插入(配合连接池)

如果数据来源是流式的,可以同时发起多个异步插入请求,利用连接池的多连接特性:

#[tokio::main]
async fn main() -> Result<(), sqlx::Error> {
    // 增大连接池大小,根据数据库配置调整,比如设为20
    let pool = sqlx::postgres::PgPoolOptions::new()
        .max_connections(20)
        .connect("postgresql://postgres:password@192.168.0.x/postgres")
        .await?;

    let data_list: Vec<&str> = vec![/* 你的数据列表 */];

    let start = Utc::now();
    // 用tokio的join_all并发执行所有插入请求
    let results = tokio::join_all(data_list.into_iter().map(|data| {
        let pool = pool.clone();
        async move {
            sqlx::query("insert into html_data (html_data) values ($1)")
                .bind(data)
                .execute(&pool)
                .await
        }
    })).await;

    // 检查所有结果
    for res in results {
        res?;
    }

    println!("concurrent inserts took {} ms", Utc::now().signed_duration_since(start).num_milliseconds());
    Ok(())
}

3. 其他性能优化点

  • 调整连接池大小:当前代码只设了5个连接,远远不够并发用。根据数据库的max_connections配置(Postgres默认是100),可以把连接池调到20-50之间。
  • 使用execute而非fetch_optional:插入操作不需要返回行,用execute更轻量,减少不必要的网络传输。
  • 数据库端优化:比如禁用表的自动提交、创建索引放到迁移完成后、开启wal_buffers优化等,这些能进一步降低插入延迟。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 08:38:09