本文介绍了 arrow 提供的各种数据对象类型,并记录了这些对象的结构。
arrow 包提供了多种用于表示数据的对象类。RecordBatch、Table 和 Dataset 对象是用于存储表格数据的二维矩形数据结构。对于列式的一维数据,提供了 Array 和 ChunkedArray 类。最后,Scalar 对象表示单个值。下表总结了这些对象,并展示了如何使用 R6 类对象创建新实例,以及以更传统的 R 风格提供相同功能的便捷函数。
| 维度 | 类 | 如何创建实例 | 便捷函数 |
|---|---|---|---|
| 0 | Scalar |
Scalar$create(value, type) |
|
| 1 | 数组 |
Array$create(vector, type) |
as_arrow_array(x) |
| 1 | ChunkedArray |
ChunkedArray$create(..., type) |
chunked_array(..., type) |
| 2 | RecordBatch |
RecordBatch$create(...) |
record_batch(...) |
| 2 | 表 |
Table$create(...) |
arrow_table(...) |
| 2 | 数据集 |
Dataset$create(sources, schema) |
open_dataset(sources, schema) |
在本文的后面,我们将更详细地研究这些内容。目前我们需要注意的是,每个对象类都对应底层 Arrow C++ 库中同名的类。
除了这些数据对象外,arrow 还定义了以下用于表示元数据的类:
Schema(模式)是Field(字段)对象的列表,用于描述表格数据对象的结构;其中:Field指定了字符串名称和DataType(数据类型);并且DataType是控制值如何表示的属性。
这些元数据对象在确保数据被正确表示方面发挥着重要作用,所有三种表格数据对象类型(Record Batch、Table 和 Dataset)都包含用于表示元数据的显式 Schema 对象。要了解有关这些元数据类的更多信息,请参阅元数据文章。
标量(Scalars)
Scalar 对象只是一个可以是任何类型的单个值。它可以是整数、字符串、时间戳或 Arrow 支持的任何不同 DataType 对象。大多数 arrow R 包的用户不太可能直接创建 Scalar,但如果有需要,可以通过调用 Scalar$create() 方法来实现。
Scalar$create("hello")## Scalar
## hello
数组(Arrays)
Array 对象是 Scalar 值的有序集合。与 Scalar 一样,大多数用户不需要直接创建 Array,但如果需要,有一个 Array$create() 方法允许你创建新的 Array。
integer_array <- Array$create(c(1L, NA, 2L, 4L, 8L))
integer_array## Array
## <int32>
## [
## 1,
## null,
## 2,
## 4,
## 8
## ]
string_array <- Array$create(c("hello", "amazing", "and", "cruel", "world"))
string_array## Array
## <string>
## [
## "hello",
## "amazing",
## "and",
## "cruel",
## "world"
## ]
可以使用方括号对 Array 进行子集化,如下所示:
string_array[4:5]## Array
## <string>
## [
## "cruel",
## "world"
## ]
Array 是不可变对象:一旦创建,就不能修改或扩展。
分块数组(Chunked Arrays)
在实践中,大多数 arrow R 包的用户更可能使用分块数组(Chunked Arrays)而不是简单的 Array。在底层,分块数组是一个或多个 Array 的集合,可以像它们是单个 Array 一样进行索引。Arrow 提供此功能的原因在数据对象布局文章中有所描述,但目前只需注意到分块数组在常规数据分析中表现得像 Array 即可。
为了说明这一点,让我们使用 chunked_array() 函数。
chunked_string_array <- chunked_array(
string_array,
c("I", "love", "you")
)chunked_array() 函数只是 ChunkedArray$create() 所提供功能的包装器。让我们打印这个对象:
chunked_string_array## ChunkedArray
## <string>
## [
## [
## "hello",
## "amazing",
## "and",
## "cruel",
## "world"
## ],
## [
## "I",
## "love",
## "you"
## ]
## ]
此输出中的双括号旨在突出显示分块数组是一个或多个 Array 的包装器这一事实。尽管由多个不同的 Array 组成,但分块数组可以像它们被首尾相连地放置在一个“类向量”对象中一样进行索引。这说明如下:

我们可以使用 chunked_string_array 来演示这一点。
chunked_string_array[4:7]## ChunkedArray
## <string>
## [
## [
## "cruel",
## "world"
## ],
## [
## "I",
## "love"
## ]
## ]
需要注意的重要一点是,“分块”在语义上没有意义。它仅仅是一个实现细节:用户永远不应该将块视为一个有意义的单位。例如,将数据写入磁盘通常会导致数据被组织成不同的块。同样,包含分配给不同块的相同值的两个分块数组被认为是等价的。为了说明这一点,我们可以创建一个包含与 chunked_string_array[4:7] 相同的四个值的单块分块数组,而不是将其拆分为两个块:
cruel_world <- chunked_array(c("cruel", "world", "I", "love"))
cruel_world## ChunkedArray
## <string>
## [
## [
## "cruel",
## "world",
## "I",
## "love"
## ]
## ]
使用 == 测试相等性会产生逐元素的比较,结果是一个新的包含四个(布尔类型)true 值的分块数组。
cruel_world == chunked_string_array[4:7]## ChunkedArray
## <bool>
## [
## [
## true,
## true,
## true,
## true
## ]
## ]
简而言之,其目的是让用户像与普通的一维数据结构交互一样与分块数组交互,而无需过多考虑底层的分块排列。
分块数组在特定意义上是可变的:可以向分块数组添加 Array 或从中移除 Array。
记录批次(Record Batches)
Record Batch 是由命名 Array 和随附的 Schema(指定与每个 Array 关联的名称和数据类型)组成的表格数据结构。Record Batch 是 Arrow 中进行数据交换的基本单位,但通常不用于数据分析。在分析场景中,Table 和 Dataset 通常更方便。
这些 Array 可以是不同的类型,但长度必须相同。每个 Array 被称为 Record Batch 的“字段”或“列”之一。你可以使用 record_batch() 函数或 RecordBatch$create() 方法创建 Record Batch。这些函数很灵活,可以接受多种格式的输入:你可以传递数据框、一个或多个命名向量、输入流,甚至是包含适当二进制数据的原始向量。例如:
rb <- record_batch(
strs = string_array,
ints = integer_array,
dbls = c(1.1, 3.2, 0.2, NA, 11)
)
rb## RecordBatch
## 5 rows x 3 columns
## $strs <string>
## $ints <int32>
## $dbls <double>
这是一个包含 5 行 3 列的 Record Batch,其概念结构如下所示:

arrow 包为 Record Batch 对象提供了 $ 方法,用于按名称提取单个列。
rb$strs## Array
## <string>
## [
## "hello",
## "amazing",
## "and",
## "cruel",
## "world"
## ]
你可以使用双括号 [[ 按位置引用列。rb$ints 数组是我们 Record Batch 中的第二列,因此我们可以用此方法将其提取出来:
rb[[2]]## Array
## <int32>
## [
## 1,
## null,
## 2,
## 4,
## 8
## ]
还有一个 [ 方法允许你像对数据框那样提取记录批次的子集。命令 rb[1:3, 1:2] 提取前三行和前两列。
rb[1:3, 1:2]## RecordBatch
## 3 rows x 2 columns
## $strs <string>
## $ints <int32>
Record Batch 不能被拼接:因为它们由 Array 组成,而 Array 是不可变对象,一旦创建,就不能向 Record Batch 添加新行。
表(Tables)
Table 由命名分块数组组成,就像 Record Batch 由命名 Array 组成一样。与 Record Batch 一样,Table 包含一个显式的 Schema,用于指定每个分块数组的名称和数据类型。
你可以像对 Record Batch 一样使用 $、[[ 和 [ 对 Table 进行子集化。与 Record Batch 不同,Table 可以拼接(因为它们由分块数组组成)。假设第二个 Record Batch 到达了:
new_rb <- record_batch(
strs = c("I", "love", "you"),
ints = c(5L, 0L, 0L),
dbls = c(7.1, -0.1, 2)
)如果不创建全新的内存对象,就不可能创建一个将 new_rb 中的数据追加到 rb 中数据的 Record Batch。然而,使用 Table,我们可以这样做:
df <- arrow_table(rb)
new_df <- arrow_table(new_rb)现在我们有了作为 Table 表示的数据集的两个片段。Table 和 Record Batch 之间的区别在于列都表示为分块数组。原始 Record Batch 中的每个 Array 都是 Table 中相应分块数组中的一个块。
rb$strs## Array
## <string>
## [
## "hello",
## "amazing",
## "and",
## "cruel",
## "world"
## ]
df$strs## ChunkedArray
## <string>
## [
## [
## "hello",
## "amazing",
## "and",
## "cruel",
## "world"
## ]
## ]
这是相同的基础数据——实际上两者引用的是相同的不可变 Array——只是被一个新的、灵活的分块数组包装器所包围。然而,正是这个包装器允许我们拼接 Table。
concat_tables(df, new_df)## Table
## 8 rows x 3 columns
## $strs <string>
## $ints <int32>
## $dbls <double>
生成的对象结构示意如下:

注意,新 Table 中的分块数组保留了这种分块结构,因为原始的 Array 都没有被移动。
df_both <- concat_tables(df, new_df)
df_both$strs## ChunkedArray
## <string>
## [
## [
## "hello",
## "amazing",
## "and",
## "cruel",
## "world"
## ],
## [
## "I",
## "love",
## "you"
## ]
## ]
数据集(Datasets)
像 Record Batch 和 Table 对象一样,Dataset 用于表示表格数据。在抽象层面上,Dataset 可以看作是一个由行和列组成的对象,并且就像 Record Batch 和 Table 一样,它包含一个显式的 Schema,用于指定与每列关联的名称和数据类型。
然而,Table 和 Record Batch 是明确表示在内存中的数据,而 Dataset 则不是。相反,Dataset 是一种抽象,指的是存储在磁盘上一个或多个文件中的数据。存储在数据文件中的值作为批处理过程加载到内存中。加载仅在需要时进行,并且仅在针对数据执行查询时进行。在这方面,Arrow Dataset 是一种与 Arrow Table 非常不同的对象,但用于分析它们的 dplyr 命令基本相同。在本节中,我们将讨论 Dataset 的结构。如果你想了解更多关于分析 Dataset 的实际细节,请参阅关于分析多文件数据集的文章。
磁盘上的数据文件
简化到最简单的形式,Dataset 的磁盘结构仅仅是一个数据文件集合,每个文件存储数据的一个子集。这些子集有时被称为“片段”(fragments),分区过程有时被称为“分片”(sharding)。按照惯例,这些文件被组织成称为 Hive 风格分区的文件夹结构:详情请参阅 hive_partition()。
为了说明这是如何工作的,让我们手动将一个多文件数据集写入磁盘,而不使用任何 Arrow Dataset 功能来完成工作。我们将从三个小数据框开始,每个数据框包含我们想要存储的数据的一个子集:
df_a <- data.frame(id = 1:5, value = rnorm(5), subset = "a")
df_b <- data.frame(id = 6:10, value = rnorm(5), subset = "b")
df_c <- data.frame(id = 11:15, value = rnorm(5), subset = "c")我们的意图是每个数据框都应存储在单独的数据文件中。正如你所看到的,这是一个结构非常清晰的分区:所有 subset = "a" 的数据属于一个文件,所有 subset = "b" 的数据属于另一个文件,所有 subset = "c" 的数据属于第三个文件。
第一步是定义并创建一个将容纳所有文件的文件夹:
ds_dir <- "mini-dataset"
dir.create(ds_dir)下一步是手动创建 Hive 风格的文件夹结构:
ds_dir_a <- file.path(ds_dir, "subset=a")
ds_dir_b <- file.path(ds_dir, "subset=b")
ds_dir_c <- file.path(ds_dir, "subset=c")
dir.create(ds_dir_a)
dir.create(ds_dir_b)
dir.create(ds_dir_c)注意,我们以“键=值”格式命名了每个文件夹,准确描述了将写入该文件夹的数据子集。这种命名结构是 Hive 风格分区的精髓。
现在我们有了文件夹,我们将使用 write_parquet() 为这三个子集中的每一个创建一个单独的 parquet 文件:
write_parquet(df_a, file.path(ds_dir_a, "part-0.parquet"))
write_parquet(df_b, file.path(ds_dir_b, "part-0.parquet"))
write_parquet(df_c, file.path(ds_dir_c, "part-0.parquet"))如果我们愿意,还可以进一步细分数据集。如果需要,一个文件夹可以包含多个文件(part-0.parquet、part-1.parquet 等)。同样,没有特别的理由必须这样命名文件 part-0.parquet:如果我们愿意,完全可以把这些文件命名为 subset-a.parquet、subset-b.parquet 和 subset-c.parquet。我们可以编写其他文件格式(如果我们想要的话),并且不一定要使用 Hive 风格的文件夹。你可以通过阅读 open_dataset() 的帮助文档来了解有关支持格式的更多信息,并通过 help("Dataset", package = "arrow") 了解如何进行细粒度控制。
总之,我们创建了一个使用 Hive 风格分区的磁盘 parquet Dataset。我们的 Dataset 由这些文件定义:
list.files(ds_dir, recursive = TRUE)## [1] "subset=a/part-0.parquet" "subset=b/part-0.parquet"
## [3] "subset=c/part-0.parquet"
为了验证一切是否有效,让我们用 open_dataset() 打开数据,并调用 glimpse() 来检查其内容:
ds <- open_dataset(ds_dir)
glimpse(ds)## FileSystemDataset with 3 Parquet files
## 15 rows x 3 columns
## $ id <int32> 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15
## $ value <double> -1.400043517, 0.255317055, -2.437263611, -0.005571287, 0.62155~
## $ subset <string> "a", "a", "a", "a", "a", "b", "b", "b", "b", "b", "c", "c", "c~
## Call `print()` for full schema details
正如你所看到的,ds Dataset 对象聚合了这三个单独的数据文件。事实上,在这种特定情况下,Dataset 非常小,以至于所有三个文件中的值都会出现在 glimpse() 的输出中。
需要注意的是,在日常数据分析工作中,你不需要以这种方式手动编写数据文件。上面的例子完全是为了说明目的。完全相同的数据集可以通过以下命令创建:
ds |>
group_by(subset) |>
write_dataset("mini-dataset")事实上,即使 ds 恰好指向一个大于内存的数据源,这个命令也应该能工作,因为 Dataset 功能的编写是为了确保在这样的流水线过程中,数据是分段加载的,以避免耗尽内存。
Dataset 对象
在上一节中,我们检查了 Dataset 的磁盘结构。现在我们转向 Dataset 对象本身的内存结构(即前面例子中的 ds)。当创建 Dataset 对象时,arrow 会在数据集文件夹中搜索合适的文件,但不会加载这些文件的内容。这些文件的路径存储在活动绑定 ds$files 中:
ds$files## [1] "/build/r/vignettes/mini-dataset/subset=a/part-0.parquet"
## [2] "/build/r/vignettes/mini-dataset/subset=b/part-0.parquet"
## [3] "/build/r/vignettes/mini-dataset/subset=c/part-0.parquet"
当调用 open_dataset() 时发生的另一件事是,构建了 Dataset 的显式 Schema 并将其存储为 ds$schema:
ds$schema## Schema
## id: int32
## value: double
## subset: string
##
## See $metadata for additional Schema metadata
默认情况下,此 Schema 仅通过检查第一个文件来推断,尽管可以在检查所有文件后构建统一的模式。为此,在调用 open_dataset() 时设置 unify_schemas = TRUE。也可以使用 open_dataset() 的 schema 参数来显式指定 Schema(详情请参阅 schema() 函数)。
读取数据的操作由 Scanner 对象执行。在使用 dplyr 接口分析 Dataset 时,你从不需要手动构建 Scanner,但为了说明目的,我们在这里这样做:
scan <- Scanner$create(dataset = ds)调用 ToTable() 方法会将 Dataset(磁盘上的)物化为 Table(内存中的):
scan$ToTable()## Table
## 15 rows x 3 columns
## $id <int32>
## $value <double>
## $subset <string>
##
## See $metadata for additional Schema metadata
此扫描过程默认是多线程的,但如有必要,可以通过在调用 Scanner$create() 时设置 use_threads = FALSE 来禁用多线程。
查询数据集
当针对 Dataset 执行查询时,会启动一个新的扫描并将结果提取回 R。例如,考虑以下 dplyr 表达式:
ds |>
filter(value > 0) |>
mutate(new_value = round(100 * value)) |>
select(id, subset, new_value) |>
collect()## # A tibble: 6 x 3
## id subset new_value
## <int> <chr> <dbl>
## 1 2 a 26
## 2 5 a 62
## 3 6 b 115
## 4 12 c 63
## 5 13 c 207
## 6 15 c 51
我们可以通过指定 Scanner$create() 的 filter 和 projection 参数,使用低级 Dataset 接口来复制这一点。要使用这些参数,你需要了解一点 Arrow 表达式(Arrow Expressions),你可以通过阅读 help("Expression", package = "arrow") 中的帮助文档找到帮助。
下面定义的扫描器模拟了上面显示的 dplyr 流水线:
scan <- Scanner$create(
dataset = ds,
filter = Expression$field_ref("value") > 0,
projection = list(
id = Expression$field_ref("id"),
subset = Expression$field_ref("subset"),
new_value = Expression$create("round", 100 * Expression$field_ref("value"))
)
)如果我们调用 as.data.frame(scan$ToTable()),它将产生与 dplyr 版本相同的结果,尽管行的顺序可能不同。
为了更好地理解查询执行时发生了什么,我们将在这里调用 scan$ScanBatches()。与 ToTable() 方法非常相似,ScanBatches() 方法针对每个文件分别执行查询,但它返回一个 Record Batch 列表,每个文件一个。此外,我们将单独把这些 Record Batch 转换为数据框:
lapply(scan$ScanBatches(), as.data.frame)## [[1]]
## id subset new_value
## 1 2 a 26
## 2 5 a 62
##
## [[2]]
## id subset new_value
## 1 6 b 115
##
## [[3]]
## id subset new_value
## 1 12 c 63
## 2 13 c 207
## 3 15 c 51
如果我们回到我们之前进行的 dplyr 查询,并使用 compute() 返回一个 Table,而不是使用 collect() 返回一个数据框,我们可以看到这个过程在工作中的证据。Table 对象是通过连接针对三个数据文件执行查询时产生的三个 Record Batch 而创建的,因此,定义 Table 一列的分块数组反映了数据文件中存在的分区结构。
tbl <- ds |>
filter(value > 0) |>
mutate(new_value = round(100 * value)) |>
select(id, subset, new_value) |>
compute()
tbl$subset## ChunkedArray
## <string>
## [
## [
## "a",
## "a"
## ],
## [
## "b"
## ],
## [
## "c",
## "c",
## "c"
## ]
## ]
附加说明
前面讨论中忽略的一个区别是
FileSystemDataset和InMemoryDataset对象之间的区别。在通常情况下,构成 Dataset 的数据存储在磁盘的文件中。毕竟,这是 Dataset 相对于 Table 的主要优势。然而,在某些情况下,利用已经存储在内存中的数据制作 Dataset 可能会很有用。在这种情况下,创建的对象类型将是InMemoryDataset。前面的讨论假设 Dataset 中存储的所有文件都具有相同的 Schema。在通常情况下这是正确的,因为每个文件在概念上都是单个矩形表的一个子集。但这不是严格要求的。
有关这些主题的更多信息,请参阅 help("Dataset", package = "arrow")。