如何在Rust中从Vec<Struct>创建Datafusion DataFrame并保存为Parquet?
使用Datafusion将Vec保存为Parquet文件
不需要通过JSON中转,Datafusion基于Apache Arrow生态,直接将结构体转换为Arrow的内存结构即可写入Parquet,以下是两种可行方案:
方案一:直接用Parquet-Arrow写入器(轻量高效)
这种方式无需依赖Datafusion的查询引擎,适合单纯写入场景:
依赖配置(Cargo.toml)
[dependencies] arrow = "51.0.0" parquet = "51.0.0" serde = { version = "1.0", features = ["derive"] } # 仅保留原结构体的Serialize,非必须
实现代码
use arrow::array::{UInt32Array, RecordBatch}; use arrow::datatypes::{DataType, Field, Schema}; use parquet::arrow::arrow_writer::ArrowWriter; use std::fs::File; use std::sync::Arc; #[derive(Serialize)] struct Test { id: u32, amount: u32 } fn write_tests_to_parquet(tests: Vec<Test>, path: &str) -> Result<(), Box<dyn std::error::Error>> { // 1. 定义与结构体匹配的Arrow Schema let schema = Arc::new(Schema::new(vec![ Field::new("id", DataType::UInt32, false), Field::new("amount", DataType::UInt32, false), ])); // 2. 从结构体集合中提取字段数据 let ids: Vec<u32> = tests.iter().map(|t| t.id).collect(); let amounts: Vec<u32> = tests.iter().map(|t| t.amount).collect(); // 3. 构建Arrow强类型数组 let id_array = UInt32Array::from(ids); let amount_array = UInt32Array::from(amounts); // 4. 组合成Arrow RecordBatch(Parquet的写入单元) let batch = RecordBatch::try_new( schema.clone(), vec![Arc::new(id_array), Arc::new(amount_array)], )?; // 5. 写入Parquet文件 let file = File::create(path)?; let mut writer = ArrowWriter::try_new(file, schema, None)?; writer.write(&batch)?; writer.close()?; Ok(()) }
方案二:通过Datafusion DataFrame API(适合带数据处理的场景)
如果需要先对数据做过滤、聚合等操作,再写入Parquet,可以用Datafusion的上下文API:
依赖配置(Cargo.toml)
[dependencies] datafusion = "32.0.0" arrow = "51.0.0" serde = { version = "1.0", features = ["derive"] }
实现代码
use arrow::array::{UInt32Array, RecordBatch}; use arrow::datatypes::{DataType, Field, Schema}; use datafusion::prelude::*; use std::sync::Arc; #[derive(Serialize)] struct Test { id: u32, amount: u32 } async fn write_tests_via_datafusion(tests: Vec<Test>, path: &str) -> Result<(), Box<dyn std::error::Error>> { // 1. 构建Arrow Schema和RecordBatch(同方案一) let schema = Arc::new(Schema::new(vec![ Field::new("id", DataType::UInt32, false), Field::new("amount", DataType::UInt32, false), ])); let ids: Vec<u32> = tests.iter().map(|t| t.id).collect(); let amounts: Vec<u32> = tests.iter().map(|t| t.amount).collect(); let id_array = UInt32Array::from(ids); let amount_array = UInt32Array::from(amounts); let batch = RecordBatch::try_new( schema.clone(), vec![Arc::new(id_array), Arc::new(amount_array)], )?; // 2. 创建Datafusion会话上下文 let ctx = SessionContext::new(); // 3. 将RecordBatch注册为临时表 ctx.register_batch("test_table", batch)?; // 4. 执行查询并写入Parquet(可替换为带过滤/聚合的SQL) ctx.write_parquet( "SELECT * FROM test_table", path, ParquetWriteOptions::default(), ).await?; Ok(()) }
关键说明
- 避免JSON中转:Arrow的内存模型与Parquet原生兼容,直接转换结构体字段为Arrow数组,比JSON序列化效率高得多。
- Schema匹配:必须保证Arrow Schema的字段类型、名称与结构体完全对应,否则会写入失败。
- 批量写入:如果数据量很大,可以拆分多个RecordBatch分批写入,避免内存溢出。
内容的提问来源于stack exchange,提问作者dade
相关产品推荐
相关产品推荐

