跳过内容

本文介绍了 arrow 提供的各种数据对象类型,并记录了这些对象的结构。

arrow 包提供了多种用于表示数据的对象类。RecordBatchTableDataset 对象是用于存储表格数据的二维矩形数据结构。对于列式的一维数据,提供了 ArrayChunkedArray 类。最后,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.parquetpart-1.parquet 等)。同样,没有特别的理由必须这样命名文件 part-0.parquet:如果我们愿意,完全可以把这些文件命名为 subset-a.parquetsubset-b.parquetsubset-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()filterprojection 参数,使用低级 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"
##   ]
## ]

附加说明

  • 前面讨论中忽略的一个区别是 FileSystemDatasetInMemoryDataset 对象之间的区别。在通常情况下,构成 Dataset 的数据存储在磁盘的文件中。毕竟,这是 Dataset 相对于 Table 的主要优势。然而,在某些情况下,利用已经存储在内存中的数据制作 Dataset 可能会很有用。在这种情况下,创建的对象类型将是 InMemoryDataset

  • 前面的讨论假设 Dataset 中存储的所有文件都具有相同的 Schema。在通常情况下这是正确的,因为每个文件在概念上都是单个矩形表的一个子集。但这不是严格要求的。

有关这些主题的更多信息,请参阅 help("Dataset", package = "arrow")

进一步阅读

  • 要了解有关 Array 内部结构的更多信息,请参阅关于数据对象布局的文章。
  • 要了解有关 Arrow 使用的不同数据类型的更多信息,请参阅关于数据类型的文章。
  • 要了解有关 Arrow 对象是如何实现的,请参阅 Arrow 规范页面。