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

Rust中如何将Polars.DataFrame存入PostgreSQL等数据库?

在Rust中把Polars DataFrame写入PostgreSQL的可行方案

方法1:利用Polars的SQLx集成直接写入

这是最简洁的方式,Polars通过sqlx提供了直接写入数据库的API,无需手动处理数据转换。

步骤1:配置依赖

在Cargo.toml中添加所需依赖:

[dependencies]
polars = { version = "0.35", features = ["sqlx", "postgres", "tokio", "io"] }
sqlx = { version = "0.7", features = ["postgres", "runtime-tokio-native-tls"] }
tokio = { version = "1.0", features = ["full"] }

步骤2:编写写入代码

use polars::prelude::*;
use sqlx::postgres::PgPoolOptions;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // 构造示例DataFrame
    let df = df!(
        "id" => [1, 2, 3, 4],
        "name" => ["Alice", "Bob", "Charlie", "Diana"],
        "age" => [25, 30, 35, 40]
    )?;

    // 建立PostgreSQL连接池
    let pool = PgPoolOptions::new()
        .max_connections(5)
        .connect("postgres://username:password@localhost:5432/db_name")
        .await?;

    // 创建目标表(如果不存在)
    sqlx::query(
        r#"
        CREATE TABLE IF NOT EXISTS users (
            id INTEGER PRIMARY KEY,
            name VARCHAR(50) NOT NULL,
            age INTEGER NOT NULL
        )
        "#,
    )
    .execute(&pool)
    .await?;

    // 将DataFrame写入数据库
    df.write_database(
        "users",
        "postgres://username:password@localhost:5432/db_name",
        WriteDatabaseOptions::default(),
    )
    .await?;

    Ok(())
}

优势:代码简洁,Polars自动处理类型映射;局限:对复杂类型(如JSON、数组)的支持依赖Polars的内置映射,自定义逻辑空间有限。

方法2:手动批量插入(基于tokio-postgres)

适合需要精细控制插入逻辑、处理特殊类型或分批次导入大数据的场景。

步骤1:配置依赖

[dependencies]
polars = { version = "0.35", features = ["tokio"] }
tokio-postgres = { version = "0.7", features = ["native-tls"] }
tokio = { version = "1.0", features = ["full"] }

步骤2:编写批量插入代码

use polars::prelude::*;
use tokio_postgres::{Client, NoTls};

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // 示例DataFrame
    let df = df!(
        "id" => [1, 2, 3, 4],
        "name" => ["Alice", "Bob", "Charlie", "Diana"],
        "age" => [25, 30, 35, 40]
    )?;

    // 连接PostgreSQL
    let (client, connection) = tokio_postgres::connect(
        "postgres://username:password@localhost:5432/db_name", 
        NoTls
    ).await?;

    // 后台维护连接
    tokio::spawn(async move {
        if let Err(e) = connection.await {
            eprintln!("connection error: {}", e);
        }
    });

    // 创建表(如果不存在)
    client.execute(
        r#"
        CREATE TABLE IF NOT EXISTS users (
            id INTEGER PRIMARY KEY,
            name VARCHAR(50) NOT NULL,
            age INTEGER NOT NULL
        )
        "#,
        &[],
    ).await?;

    // 提取DataFrame列数据
    let ids = df.column("id")?.i32()?.into_no_null_iter().collect::<Vec<_>>();
    let names = df.column("name")?.utf8()?.into_no_null_iter()
        .map(|s| s.to_string())
        .collect::<Vec<_>>();
    let ages = df.column("age")?.i32()?.into_no_null_iter().collect::<Vec<_>>();

    // 开启事务批量插入
    let transaction = client.transaction().await?;
    let stmt = transaction.prepare(
        "INSERT INTO users (id, name, age) VALUES ($1, $2, $3)"
    ).await?;

    for i in 0..df.height() {
        transaction.execute(&stmt, &[&ids[i], &names[i], &ages[i]]).await?;
    }

    transaction.commit().await?;

    Ok(())
}

优势:完全自定义插入逻辑,可处理空值、特殊类型转换;建议:针对超大数据集,可拆分数据为多个批次(如每1000条提交一次事务),避免内存溢出。

方法3:用PostgreSQL COPY FROM优化大数据写入

这是PostgreSQL中性能最高的批量导入方式,适合百万级以上的数据量。

代码示例

use polars::prelude::*;
use tokio_postgres::{Client, NoTls};

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let df = df!(
        "id" => [1, 2, 3, 4],
        "name" => ["Alice", "Bob", "Charlie", "Diana"],
        "age" => [25, 30, 35, 40]
    )?;

    let (client, connection) = tokio_postgres::connect(
        "postgres://username:password@localhost:5432/db_name", 
        NoTls
    ).await?;

    tokio::spawn(async move {
        if let Err(e) = connection.await {
            eprintln!("connection error: {}", e);
        }
    });

    // 创建表
    client.execute(
        r#"
        CREATE TABLE IF NOT EXISTS users (
            id INTEGER PRIMARY KEY,
            name VARCHAR(50) NOT NULL,
            age INTEGER NOT NULL
        )
        "#,
        &[],
    ).await?;

    // 将DataFrame转为CSV字节流
    let mut csv_buffer = Vec::new();
    df.write_csv(&mut csv_buffer)?;

    // 用COPY FROM导入
    client.execute(
        "COPY users (id, name, age) FROM STDIN WITH (FORMAT csv, HEADER true)",
        &[&csv_buffer],
    ).await?;

    Ok(())
}

优势:写入速度比单条插入快10-100倍,适合大规模数据导入;扩展:也可以将DataFrame转为Parquet格式,用COPY ... FROM PROGRAM或外部文件导入。

关键注意事项

  • 类型匹配:确保Polars列类型与PostgreSQL字段类型对应,例如Utf8→TEXT/VARCHAR、Int32→INTEGER、Float64→DOUBLE PRECISION。
  • 空值处理:若DataFrame含空值,需将表字段设为允许NULL,代码中用into_iter()替代into_no_null_iter()处理Option类型。
  • 连接配置:生产环境建议使用连接池(如sqlx的PgPool),避免频繁创建销毁连接。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 16:28:17