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

Rust+Datafusion:将DataFrame转换为JSON的类型兼容问题求助

问题:将Datafusion的DataFrame转换为JSON格式

WIP代码仓库:rust-datafusion-csv-processing

刚接触Rust编程2天,尝试Rust约3小时后便开始解决此问题,至今未果,恳请各位帮忙。

我的目标是将Datafusion中的DataFrame转为JSON格式(最终用于API的HTTP响应)。

调用DataFrame的collect方法后会得到datafusion::arrow::record_batch::RecordBatch类型,我在转换该类型时遇到困难。

我尝试过两种方法:

  • 使用Arrow库的json::writer::record_batches_to_json_rows,但报错提示datafusion::arrow::record_batch::RecordBatch与arrow::record_batch::RecordBatch是名称相似但实际不同的类型,无法成功转换类型解决此问题。
  • 尝试将RecordBatch转为向量,分别提取表头和值,表头提取成功,但值提取失败。
let mut header = Vec::new();
// let mut rows = Vec::new();

for record_batch in data_vec {
    // get data
    println!("record_batch.columns: : {:?}", record_batch.columns());
    for col in record_batch.columns() {
        for row in 0..col.len() {
            // println!("Cow: {:?}", col);
            // println!("Row: {:?}", row);
            // let value = col.as_any().downcast_ref::<StringArray>().unwrap().value(row);
            // rows.push(value);
        }
    }
    // get headers
    for field in record_batch.schema().fields() {
        header.push(field.name().to_string());
    }
};

有没有人知道如何完成这个转换?

完整脚本如下:

// datafusion examples: https://github.com/apache/arrow-datafusion/tree/master/datafusion-examples/examples
// datafusion docs: https://arrow.apache.org/datafusion/
use datafusion::prelude::*;
use datafusion::arrow::datatypes::{Schema};

use arrow::json;

// use serde::{ Deserialize };
use serde_json::to_string;

use std::sync::Arc;
use std::str;
use std::fs;
use std::ops::Deref;

type DFResult = Result<Arc<DataFrame>, datafusion::error::DataFusionError>;

struct FinalObject {
    schema: Schema,
    // columns: Vec<Column>,
    num_rows: usize,
    num_columns: usize,
}

// to allow debug logging for FinalObject
impl std::fmt::Debug for FinalObject {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        // write!(f, "FinalObject {{ schema: {:?}, columns: {:?}, num_rows: {:?}, num_columns: {:?} }}",
        write!(f, "FinalObject {{ schema: {:?}, num_rows: {:?}, num_columns: {:?} }}",
                // self.schema,  self.columns, self.num_columns, self.num_rows)
                self.schema, self.num_columns, self.num_rows)
    }
}

fn create_or_delete_csv_file(path: String, content: Option<String>, operation: &str) {
    match operation {
        "create" => {
            match content {
                Some(c) => fs::write(path, c.as_bytes()).expect("Problem with writing file!"),
                None => println!("The content is None, no file will be created"),
            }
        }
        "delete" => {
            // Delete the csv file
            fs::remove_file(path).expect("Problem with deleting file!"),
        }
        _ => println!("Invalid operation"),
    }
}


async fn read_csv_file_with_inferred_schema(file_name_string: String) -> DFResult {
    // create string csv data
    let csv_data_string = "heading,value\nbasic,1\ncsv,2\nhere,3".to_string();

    // Create a temporary file
    create_or_delete_csv_file(file_name_string.clone(), Some(csv_data_string), "create");

    // Create a session context
    let ctx = SessionContext::new();

    // Register a lazy DataFrame using the context
    let df = ctx.read_csv(file_name_string.clone(), CsvReadOptions::default()).await.expect("An error occurred while reading the CSV string");

    // return the dataframe
    Ok(Arc::new(df))
}

#[tokio::main]
async fn main() {

    let file_name_string = "temp_file.csv".to_string();

    let arc_csv_df = read_csv_file_with_inferred_schema(file_name_string.clone()).await.expect("An error occurred while reading the CSV string (funct: read_csv_file_with_inferred_schema)");

    // have to use ".clone()" each time I want to use this ref
    let deref_df = arc_csv_df.deref();

    // print to console
    deref_df.clone().show().await.expect("An error occurred while showing the CSV DataFrame");

    // collect to vec
    let record_batches = deref_df.clone().collect().await.expect("An error occurred while collecting the CSV DataFrame");
    // println!("Data: {:?}", data);
    
    // record_batches == <Vec<RecordBatch>>. Convert to RecordBatch
    let record_batch = record_batches[0].clone();

    // let json_string = to_string(&record_batch).unwrap();

    // let mut writer = datafusion::json::writer::RecordBatchJsonWriter::new(vec![]);
    // writer.write(&record_batch).unwrap();
    // let json_rows = writer.finish();
    
    let json_rows = json::writer::record_batches_to_json_rows(&[record_batch]);

    println!("JSON: {:?}", json_rows);

    // get final values from recordbatch
    // https://docs.rs/arrow/latest/arrow/record_batch/struct.RecordBatch.html
    // https://users.rust-lang.org/t/how-to-use-recordbatch-in-arrow-when-using-datafusion/70057/2
    // https://github.com/apache/arrow-rs/blob/6.5.0/arrow/src/util/pretty.rs
    // let record_batches_vec = record_batches.to_vec();

    let mut header = Vec::new();
    // let mut rows = Vec::new();

    for record_batch in data_vec {
        // get data
        println!("record_batch.columns: : {:?}", record_batch.columns());
        for col in record_batch.columns() {
            for row in 0..col.len() {
                // println!("Cow: {:?}", col);
                // println!("Row: {:?}", row);
                // let value = col.as_any().downcast_ref::<StringArray>().unwrap().value(row);
                // rows.push(value);
            }
        }
        // get headers
        for field in record_batch.schema().fields() {
            header.push(field.name().to_string());
        }
    };

    // println!("Header: {:?}", header);
    
    // Delete temp csv
    create_or_delete_csv_file(file_name_string.clone(), None, "delete");
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 00:35:28