Arrow 简要介绍
Apache Arrow 轻松入门
概述
DataFusion 使用 Apache Arrow 作为其原生内存格式,因此任何使用 DataFusion 的人都可能在某个时刻与 Arrow 打交道。本指南介绍有效使用 DataFusion 所需了解的 Arrow 核心概念。
Apache Arrow 为内存数据定义了标准化的列式表示。这使得不同的系统和语言(例如 Rust 和 Python)能够以零拷贝方式交换数据,从而避免序列化开销。除了零拷贝交换之外,Arrow 还标准化了列式数据表示的最佳实践,通过向量化执行实现高性能的分析处理。
列式布局
快速示意:行式存储(左)与 Arrow 的列式布局(右)。如需更深入的入门介绍,请参阅 arrow2 指南。
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)RecordBatch
Arrow 打包数据的标准单元是 RecordBatch。
RecordBatch 表示表的一个水平切片——一组长度相等、符合既定 schema 的列式数组。切片中的每一列都是一个连续的 Arrow 数组,且所有列的行数(长度)相同。这种分块的、不可变的单元使得高效流式处理和并行执行成为可能。
可以把它看作具有两个视角:
- 内部是列式的:每一列(
id、name、age)都是一个连续的数组,针对向量化操作进行了优化 - 外部是按行分块的:每个 batch 代表一批行(例如第 1 到 1000 行),因此它是流式处理中易于管理的单元
RecordBatch 是不可变的快照——一旦创建便无法修改。任何转换都会产生一个新的 RecordBatch,从而无需锁或协调开销即可实现安全的并行处理。
这一设计使 DataFusion 能够以基于行的分块流进行处理,同时从列式布局中获得最大性能。
流式穿越引擎
DataFusion 将查询处理为基于拉取(pull)的管道,其中算子向其输入请求 batch。这种流式方法可以尽早产出结果、限制内存使用(仅在必要时溢写到磁盘),并天然支持跨多个 CPU 核心的并行执行。
例如,给定以下查询:
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 是在查询执行各阶段之间流转的列式数据“包裹”。每个算子都以增量方式处理批次,使系统能够在读取完整输入之前就开始产出结果。
创建 ArrayRef 与 RecordBatch
有时你需要以编程方式构造 Arrow 数据,而不是从文件中读取。
首先需要为每一列创建一个 Arrow 数组。arrow-rs 提供了数组构建器(array builder)以及 From 实现,可用于从 Rust 向量创建数组。
use arrow::array::{StringArray, Int32Array};
// Create an Int32Array from a vector of i32 values
let ids = Int32Array::from(vec![1, 2, 3]);
// There are similar constructors for other array types, e.g., StringArray, Float64Array, etc.
let names = StringArray::from(vec![Some("alice"), None, Some("carol")]);Arrow 数组中的每个元素都可以为 “null”(即缺失值)。通常,数组由 Option<T> 值创建,以表示可空性(例如上面的 Some("alice") 与 None)。
注意:你会在代码中频繁看到 Arc——Arrow 数组被包装在 Arc(原子引用计数指针)中,以便在算子与任务之间实现低成本、线程安全的共享。ArrayRef 只是 Arc<dyn Array> 的类型别名。要创建 ArrayRef,只需像下面这样用 Arc::new(...) 包装你的数组。
use std::sync::Arc;
// To get an ArrayRef, wrap the Int32Array in an Arc.
// (note you will often have to explicitly type annotate to ArrayRef)
let arr: ArrayRef = Arc::new(Int32Array::from(vec![1, 2, 3]));
// you can also store Strings and other types in ArrayRefs
let arr: ArrayRef = Arc::new(
StringArray::from(vec![Some("alice"), None, Some("carol")])
);要创建一个 RecordBatch,你需要定义它的 Schema(即列名和类型),并如下面所示以 ArrayRef 的形式提供相应的列:
use arrow_schema::{DataType, Field, Schema};
// Create the columns as Arrow arrays
let ids = Int32Array::from(vec![1, 2, 3]);
let names = StringArray::from(vec![Some("alice"), None, Some("carol")]);
// Create the schema
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int32, false), // false means non-nullable
Field::new("name", DataType::Utf8, true), // true means nullable
]));
// Assemble the columns
let cols: Vec<ArrayRef> = vec![
Arc::new(ids),
Arc::new(names)
];
// Finally, create the RecordBatch
RecordBatch::try_new(schema, cols).expect("Failed to create RecordBatch");ArrayRef 与 RecordBatch 的使用
大多数 DataFusion API 都以 ArrayRef 和 RecordBatch 的形式提供。要操作底层数据,通常需要将 ArrayRef 向下转型为其具体类型(例如 Int32Array)。
为此,可以使用 as_any().downcast_ref::<T>() 方法,或者使用 AsArray trait 中的 as_::<T>() 辅助方法。
// First check the data type of the array
match arr.data_type() {
&DataType::Int32 => {
// Downcast to Int32Array
let int_array = arr.as_primitive::<Int32Type>();
// Now you can access Int32Array methods
for i in 0..int_array.len() {
println!("Value at index {}: {}", i, int_array.value(i));
}
}
_ => {
println ! ("Array is not of type Int32");
}
}以下两种向下转型方法是等价的:
// Downcast to Int32Array using as_any
let int_array1 = arr.as_any().downcast_ref::<Int32Array>().unwrap();
// This is the same as using the as_::<T>() helper
let int_array2 = arr.as_primitive::<Int32Type>();
assert_eq!(int_array1, int_array2);常见陷阱
在使用 Arrow 和 RecordBatch 时,请注意以下常见问题:
- Schema 一致性:流中的所有批次必须共享完全相同的
Schema。例如,你不能让某一批次中的某一列是Int32,而下一批次中同一列是Int64,即使这些值完全放得下 - 不可变性:数组是不可变的——要“修改”数据,你必须构建新的数组或新的 RecordBatch。例如,要更改数组中的某个值,你需要创建一个包含更新值的新数组
- 逐行处理:尽可能避免逐元素遍历 Array,改用 Arrow 内置的计算内核
- 类型不匹配:来自不同文件的混合输入类型可能需要显式转换。例如,CSV 文件中的字符串列
"123"不会自动与 Parquet 文件中的整数列123进行连接——你需要将其中一方转换以匹配另一方。在适当的地方使用 Arrow 的cast内核 - 批次大小假设:不要假设某个特定的批次大小;应始终迭代直到流结束。一个文件可能产生 8192 行的批次,而另一个文件可能产生 1024 行的批次
延伸阅读
Arrow 文档:
- Arrow 格式简介 - 了解 Arrow 规范及其为何能实现零拷贝数据共享
- Arrow 列式格式 - 深入了解内存布局以进行性能优化
- Arrow Rust 文档 - Rust 实现的完整 API 参考
关键 API 参考:
- RecordBatch - 列式数据的基础数据结构(一个表切片)
- ArrayRef - 表示一个引用计数的 Arrow 数组(单个列)
- DataType - 所有受支持的 Arrow 数据类型的枚举(例如 Int32、Utf8)
- Schema - 描述 RecordBatch 的结构(列名和类型)
评论
登录后参与评论
KnowForge