
大数据数据分析后端【免费下载链接】datafusionApache DataFusion SQL Query Engine项目地址https://gitcode.com/gh_mirrors/datafu/datafusion点击查看免费下载导读Apache DataFusion 将 Apache Arrow 作为其原生内存数据格式因此任何使用 DataFusion 的开发者都不可避免地会与 Arrow 打交道。本文基于 DataFusion 官方用户指南docs/source/user-guide/arrow-introduction.md整理而成系统讲解 Arrow 的列式内存布局、RecordBatch与ArrayRef这两大核心数据结构、DataFusion 的流式拉取执行管线以及用 Rust 编程方式构造和操作 Arrow 数据的完整示例。读完本文你将掌握如何在 DataFusion 中看懂并编写操作 Arrow 数组与批次的代码并理解批大小、Schema 一致性等关键实践细节。Overview为什么 DataFusion 依赖 ArrowDataFusion 使用 Apache Arrow 作为其原生内存格式这意味着几乎所有与 DataFusion 的交互最终都会落在 Arrow 的数据结构上。Arrow 的核心价值有两点标准化的列式内存表示不同系统与语言例如 Rust 与 Python可以以零拷贝zero-copy方式共享数据省去序列化开销列式数据表示的最佳实践通过向量化vectorized执行实现高性能分析处理。列式布局让 CPU 可以一次性处理整列连续内存中的数据配合 SIMD 等向量化手段显著提升分析型工作负载的效率。这也是 DataFusion 把执行引擎建立在 Arrow 之上的根本原因。列式布局行存储与列存储的直观对比理解 Arrow 的第一步是理解它与传统行式存储的差异。下面这张 ASCII 示意图直观展示了二者的区别Traditional Row Storage: Arrow Columnar Storage: ┌──────────────────┐ ┌─────────┬─────────┬──────────┐ │ id │ name │ age │ │ id │ name │ age │ ├────┼──────┼──────┤ ├─────────┼─────────┼──────────┤ │ 1 │ A │ 30 │ │ [1,2,3] │ [A,B,C] │[30,25,35]│ │ 2 │ B │ 25 │ └─────────┴─────────┴──────────┘ │ 3 │ C │ 35 │ ↑ ↑ ↑ └──────────────────┘ Int32Array StringArray Int32Array (read entire rows) (process entire columns at once)传统行存储左侧数据按行连续存放查询时通常需要整行读取即使只需要其中的个别列Arrow 列存储右侧每一列如id、name、age各自是一段连续数组Int32Array、StringArray处理时一次处理整列天然契合向量化运算。需要更深入理解列式内存布局的读者可参阅 arrow2 guide。RecordBatchArrow 打包数据的标准单元两个视角内部列式外部行块Arrow 打包数据的标准单元是[RecordBatch]官方 API 参考。一个RecordBatch表示一张表的水平切片——一组等长的列式数组并遵循一个已定义的 Schema。其中每一列都是一个连续的 Arrow 数组所有列拥有相同的行数长度。可以把RecordBatch理解为两个视角的结合内部列式Columnar inside每一列如id、name、age都是连续数组为向量化操作而优化外部按行分块Row-chunked externally这个批次代表一个行块例如第 1~1000 行是流式传输中的可控单元。这种设计让 DataFusion 既可以按行块流式处理数据又能在每个块内部享受列式布局的最大性能收益。不可变性无需锁的并行安全RecordBatch是不可变快照——一旦创建便无法修改任何变换都会产生一个新的RecordBatch。这一特性使得多个并行任务可以安全地共享同一批次而无需加锁或协调开销。在 DataFusion 的执行引擎中这一不可变性是并行执行的基础上游算子产出的批次可以无顾虑地分发给多个下游线程。流式执行DataFusion 的拉取式管线DataFusion 以拉取式pull-based管线处理查询算子向它的输入请求批次。这种流式方法带来三重收益能够尽早产生结果无需等待整个输入读完约束内存占用仅在必要时将中间结果溢写spill到磁盘天然支持跨多核并行执行。以如下查询为例SELECT name FROM data.parquet WHERE id 10DataFusion 的物理执行管线如下┌─────────────┐ ┌──────────────┐ ┌────────────────┐ ┌──────────────────┐ ┌──────────┐ │ Parquet │───▶│ Scan │───▶│ Filter │───▶│ Projection │───▶│ Results │ │ File │ │ Operator │ │ Operator │ │ Operator │ │ │ └─────────────┘ └──────────────┘ └────────────────┘ └──────────────────┘ └──────────┘ (reads data) (id 10) (keeps name col) RecordBatch ───▶ RecordBatch ────▶ RecordBatch ────▶ RecordBatch在这个管线中RecordBatch是列式数据的“包裹”在查询执行的各个阶段之间流动。每个算子增量地处理批次从而在读取完整输入之前就能产出结果。源码佐证流接口与批大小从源码结构看DataFusion 的流式接口定义在 datafusion/execution/src/stream.rspub type SendableRecordBatchStream PinBoxdyn RecordBatchStream Send;每个RecordBatchStream返回的RecordBatch都必须与RecordBatchStream::schema()返回的 Schema 一致见 stream.rs这与后文“Schema 一致性”的注意事项相呼应。批大小batch size由配置项datafusion.execution.batch_size控制默认值为8192见 docs/source/user-guide/configs.md。正如官方指南中的常见陷阱所述“一个文件可能产出 8192 行的批次而另一个文件可能产出 1024 行的批次”——因此永远不要假设批次大小固定而应迭代直到流结束。创建 ArrayRef 与 RecordBatch有时你需要以编程方式创建 Arrow 数据而不是从文件读取。第一步是为每一列创建一个 Arrow 数组。arrow-rs 提供了数组构建器array builders以及从 Rust 向量直接构造数组的From实现。从 Rust 向量构造数组use arrow::array::{StringArray, Int32Array}; // 从一个 i32 向量创建 Int32Array let ids Int32Array::from(vec![1, 2, 3]); // 其他数组类型有类似的构造器例如 StringArray、Float64Array 等 let names StringArray::from(vec![Some(alice), None, Some(carol)]);注意 Arrow 数组中每个元素都可以是“null”即缺失。通常用OptionT值创建数组以表达可空性——上面的Some(alice)与None即分别代表“有值”与“缺失”。ArrayRefArc 包裹的数组你会频繁在 DataFusion 代码中看到Arc原子引用计数指针Arrow 数组被包裹在Arc中以在算子与任务之间实现廉价、线程安全的共享。ArrayRef只是Arcdyn Array的类型别名。要创建ArrayRef用Arc::new(...)包裹你的数组即可use std::sync::Arc; use arrow::array::{ArrayRef, Int32Array, StringArray}; // 要得到 ArrayRef将 Int32Array 包进 Arc // 注意通常你需要显式标注类型为 ArrayRef let arr: ArrayRef Arc::new(Int32Array::from(vec![1, 2, 3])); // 也可以把字符串等类型放进 ArrayRef let arr: ArrayRef Arc::new( StringArray::from(vec![Some(alice), None, Some(carol)]) );定义 Schema 并组装 RecordBatch要创建RecordBatch需要先定义它的Schema列名与类型然后把对应列作为ArrayRef提供给它use std::sync::Arc; use arrow_schema::{DataType, Field, Schema}; use arrow::array::{ArrayRef, Int32Array, StringArray, RecordBatch}; // 创建列Arrow 数组 let ids Int32Array::from(vec![1, 2, 3]); let names StringArray::from(vec![Some(alice), None, Some(carol)]); // 创建 Schema let schema Arc::new(Schema::new(vec![ Field::new(id, DataType::Int32, false), // false 表示不可空 Field::new(name, DataType::Utf8, true), // true 表示可空 ])); // 组装列 let cols: VecArrayRef vec![ Arc::new(ids), Arc::new(names) ]; // 最终创建 RecordBatch RecordBatch::try_new(schema, cols).expect(Failed to create RecordBatch);Field::new的三个参数依次是列名、数据类型与可空标志false/true。RecordBatch::try_new会校验列长度与 Schema 一致性失败时返回Err因此示例中用expect显式处理。在真实项目中同样的构造模式可见于 datafusion-examples/examples/dataframe/dataframe.rs 等示例这些示例用SessionContext配合read_parquet、read_csv读取文件并执行查询而read_memory一类功能正是把内存中的RecordBatch注册为可查询的表。DataFusion 的Cargo.toml也直接依赖 workspace 级的arrow与arrow-schemacrate见 datafusion/core/Cargo.toml。操作 ArrayRef 与 RecordBatchDataFusion 的大部分 API 都以ArrayRef和RecordBatch为操作单位。要访问底层数据通常需要把ArrayRef**向下转型downcast**为具体类型例如Int32Array。方式一as_any().downcast_ref::T()通过as_any().downcast_ref::T()方法可以拿到具体类型的引用use std::sync::Arc; use arrow::datatypes::{DataType, Int32Type}; use arrow::array::{AsArray, ArrayRef, Int32Array, RecordBatch}; let arr: ArrayRef Arc::new(Int32Array::from(vec![1, 2, 3])); // 先检查数组的数据类型 match arr.data_type() { DataType::Int32 { // 向下转型为 Int32Array let int_array arr.as_primitive::Int32Type(); // 现在可以访问 Int32Array 的方法 for i in 0..int_array.len() { println!(Value at index {}: {}, i, int_array.value(i)); } } _ { println!(Array is not of type Int32); } }方式二AsArraytrait 的as_::T()辅助方法也可以使用 AsArray trait 提供的as_::T()辅助方法。以下两种向下转型方式是等价的use std::sync::Arc; use arrow::datatypes::{DataType, Int32Type}; use arrow::array::{AsArray, ArrayRef, Int32Array, RecordBatch}; let arr: ArrayRef Arc::new(Int32Array::from(vec![1, 2, 3])); // 使用 as_any 向下转型为 Int32Array let int_array1 arr.as_any().downcast_ref::Int32Array().unwrap(); // 与使用 as_::T() 辅助方法相同 let int_array2 arr.as_primitive::Int32Type(); assert_eq!(int_array1, int_array2);常见陷阱在 DataFusion 中处理 Arrow 与RecordBatch时官方指南提醒注意以下常见问题Schema 一致性一个流中所有批次必须共享完全相同的Schema。例如你不能让一个批次的某列是Int32而下一个批次同一列变成Int64即使数值放得下也不行。这与源码中RecordBatchStream的契约datafusion/execution/src/stream.rs完全对应不可变性数组是不可变的——要“修改”数据必须构建新数组或新的RecordBatch。例如要修改数组中的某个值就创建一个携带更新值的新数组逐行处理Row by Row Processing尽量避免逐元素迭代数组优先使用 Arrow 内置的 compute kernels类型不匹配跨文件的混合输入类型可能需要显式 cast。例如来自 CSV 文件的字符串列123不会自动与来自 Parquet 文件的整数列123做连接join你需要把其中之一 cast 成另一个的类型。适当使用 Arrow 的castkernel批次大小假设不要假设某个固定的批次大小始终迭代直到流结束。一个文件可能产出 8192 行的批次另一个可能产出 1024 行的批次。批大小默认值即由datafusion.execution.batch_size默认 8192决定可通过 配置项 调整。进一步阅读Arrow 官方文档Arrow Format Introduction理解 Arrow 规范及为何它能实现零拷贝数据共享Arrow Columnar Format深入内存布局以做性能优化Arrow Rust DocumentationRust 实现的完整 API 参考。关键 API 参考RecordBatch列式数据表的切片的基础数据结构ArrayRef引用计数的 Arrow 数组单列DataType所有受支持 Arrow 数据类型的枚举例如 Int32、Utf8Schema描述 RecordBatch 的结构列名与类型。在 DataFusion 中继续深入实践时可以参考 DataFrame 与数据读写示例 了解read_parquet、read_csv的用法或阅读 库用户指南扩展点 了解如何基于RecordBatch流实现自定义TableProvider与算子。赞分享大数据数据分析后端【免费下载链接】datafusionApache DataFusion SQL Query Engine项目地址https://gitcode.com/gh_mirrors/datafu/datafusion点击查看免费下载相关推荐Apache Arrow DataFusion 配置参数详解Apache Arrow DataFusion 配置参数详解 概述 Apache Arrow DataFusion 是一个高性能、可扩展的查询引擎专为构建高质上一篇Reddit视频制作终极指南如何为AI配音添加专业级音频混响效果下一篇Native Client API完全手册spawn、exec、env等命令的实战应用创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考