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

如何通过数据库流式处理将数据写入指定结构的JSON文件?

流式生成大JSON文件(结合SeaORM流式查询)

要实现流式写入目标JSON结构,你不能直接用serde_json::to_string序列化整个大结构体,得手动拼接JSON框架,逐个流式写入数组元素,配合SeaORM的streaming查询分批次处理数据。以下是具体实现步骤和代码:

核心思路

  1. 先手动写入JSON的头部和第一个数组的起始标记
  2. 流式读取数据库中的measurements数据,逐个序列化并写入文件,注意元素间的逗号分隔
  3. 写入第一个数组的结束标记和第二个数组的起始标记
  4. 同样流式读取media数据并逐个写入
  5. 最后写入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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 09:18:32