如何通过数据库流式处理将数据写入指定结构的JSON文件?
流式生成大JSON文件(结合SeaORM流式查询)
要实现流式写入目标JSON结构,你不能直接用serde_json::to_string序列化整个大结构体,得手动拼接JSON框架,逐个流式写入数组元素,配合SeaORM的streaming查询分批次处理数据。以下是具体实现步骤和代码:
核心思路
- 先手动写入JSON的头部和第一个数组的起始标记
- 流式读取数据库中的
measurements数据,逐个序列化并写入文件,注意元素间的逗号分隔 - 写入第一个数组的结束标记和第二个数组的起始标记
- 同样流式读取
media数据并逐个写入 - 最后写入JSON的尾部标记
代码实现
1. 定义对应的结构体(需实现SeaORM的Model和Serde的Serialize)
use sea_orm::{entity::*, query::*, DbConn}; use serde::Serialize; use std::fs::File; use std::io::{BufWriter, Write}; use serde_json::Serializer; // Measurements 实体和序列化结构体 #[derive(Clone, Debug, PartialEq, DeriveEntityModel, Serialize)] #[sea_orm(table_name = "measurements")] pub struct Model { #[sea_orm(primary_key)] pub id: i32, pub timestamp: String, pub value: f64, } #[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)] pub enum Relation {} impl ActiveModelBehavior for ActiveModel {} // Media 实体和序列化结构体 #[derive(Clone, Debug, PartialEq, DeriveEntityModel, Serialize)] #[sea_orm(table_name = "media")] pub struct MediaModel { #[sea_orm(primary_key)] pub id: i32, pub name: String, pub author: String, } #[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)] pub enum MediaRelation {} impl ActiveModelBehavior for MediaActiveModel {}
2. 流式写入JSON的核心函数
pub async fn stream_export_json(conn: &DbConn, output_path: &str) -> anyhow::Result<()> { // 创建带缓冲的文件写入器,提升IO效率 let file = File::create(output_path)?; let mut writer = BufWriter::new(file); // 1. 写入JSON头部和measurements数组起始 writer.write_all(b"{\"measurements\": [")?; // 2. 流式处理measurements数据 let mut measurements_stream = Entity::find().stream(conn).await?; let mut first_measurement = true; while let Some(measurement) = measurements_stream.next().await { let measurement = measurement?; if !first_measurement { // 非第一个元素,先写逗号分隔 writer.write_all(b",")?; } else { first_measurement = false; } // 序列化当前条目并写入文件 serde_json::to_writer(&mut writer, &measurement)?; // 按需刷新缓冲区,平衡性能与实时写入 writer.flush()?; } // 3. 写入measurements结束和media数组起始 writer.write_all(b"], \"media\": [")?; // 4. 流式处理media数据 let mut media_stream = MediaEntity::find().stream(conn).await?; let mut first_media = true; while let Some(media) = media_stream.next().await { let media = media?; if !first_media { writer.write_all(b",")?; } else { first_media = false; } serde_json::to_writer(&mut writer, &media)?; writer.flush()?; } // 5. 写入JSON尾部 writer.write_all(b"]}")?; writer.flush()?; Ok(()) }
关键细节说明
- 缓冲写入:用
BufWriter包装File,避免频繁磁盘IO,大幅提升写入性能 - 逗号处理:每个数组的第一个元素无需前置逗号,后续元素需要,通过布尔标记控制逻辑
- 错误处理:所有IO和数据库操作的错误通过
?向上传递,可根据项目需求替换错误类型(比如用标准库std::result::Result) - SeaORM Streaming:通过
.stream(conn)获取异步迭代器,逐个拉取数据库数据,不会一次性加载所有数据到内存
内容的提问来源于stack exchange,提问作者IgnisDa
相关产品推荐
相关产品推荐

