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

如何在Rust中将内存数据转为DataFusion可用的DataFrame?

解决方法

不用把数据写入磁盘,你可以直接基于内存中的数据构建LogicalPlan进而创建DataFrame,核心是利用DataFusion提供的MemTable组件。具体步骤如下:

1. 准备Schema与内存数据格式转换

先确认已有明确的Schema(对应数据列的类型定义),再将内存数据转换成DataFusion的标准内存数据单元RecordBatch:

use arrow::datatypes::{DataType, Field, Schema};
use arrow::record_batch::RecordBatch;
use datafusion::datasource::MemTable;
use datafusion::execution::context::SessionState;
use datafusion::logical_expr::LogicalPlan;
use datafusion::dataframe::DataFrame;
use std::sync::Arc;

// 定义数据Schema
let schema = Schema::new(vec![
    Field::new("id", DataType::Int32, false),
    Field::new("name", DataType::Utf8, false),
    Field::new("age", DataType::Int32, true),
]);

// 将内存原始数据转为Arrow数组格式
let ids = arrow::array::Int32Array::from(vec![1, 2, 3]);
let names = arrow::array::StringArray::from(vec!["Alice", "Bob", "Charlie"]);
let ages = arrow::array::Int32Array::from(vec![Some(25), None, Some(30)]);

// 构建RecordBatch
let batch = RecordBatch::try_new(schema.clone(), vec![
    ids.into(),
    names.into(),
    ages.into(),
]).unwrap();

2. 用MemTable生成LogicalPlan

通过MemTable封装内存数据,再基于它生成LogicalPlan:

// 创建内存表实例
let mem_table = MemTable::try_new(schema, vec![vec![batch]]).unwrap();

// 生成扫描内存表的LogicalPlan
let plan = LogicalPlan::Scan(datafusion::logical_expr::Scan {
    table_name: "in_memory_data".to_string(),
    source: Arc::new(mem_table),
    projection: None, // 传None表示选择所有列,也可指定列索引列表
    filters: vec![],  // 可选:提前设置过滤条件
    fetch: None,      // 可选:限制返回的行数
});

3. 构建DataFrame并使用API处理

最后用SessionState和生成的LogicalPlan初始化DataFrame,即可调用DataFusion的DataFrame API进行处理:

// 创建默认会话状态
let session_state = SessionState::new();

// 初始化DataFrame
let df = DataFrame::new(session_state, plan);

// 示例:选择指定列并打印结果
df.select_columns(&["id", "name"]).unwrap().show().unwrap();

补充说明

  • 如果有多个RecordBatch,创建MemTable时可传入二维向量(外层对应数据分区,内层对应每个分区的Batch)。
  • 该方案全程在内存中操作,无磁盘IO开销,适合大数据量场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 10:15:03