用户指南

Arrow 简要介绍

师成师成· 更新于 2026-09-28· 阅读 19 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

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 > 10

DataFusion 流水线的结构如下:

┌─────────────┐    ┌──────────────┐    ┌────────────────┐    ┌──────────────────┐    ┌──────────┐
│ 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 文档:

关键 API 参考:

  • RecordBatch - 列式数据的基础数据结构(一个表切片)
  • ArrayRef - 表示一个引用计数的 Arrow 数组(单个列)
  • DataType - 所有受支持的 Arrow 数据类型的枚举(例如 Int32、Utf8)
  • Schema - 描述 RecordBatch 的结构(列名和类型)

评论

登录后参与评论

正在加载评论…