pyarrow.dataset.Fragment#

class pyarrow.dataset.Fragment#

基类:_Weakrefable

来自数据集(Dataset)的数据片段(Fragment)。

__init__(*args, **kwargs)#

方法

__init__(*args, **kwargs)

count_rows(self, Expression filter=None, ...)

计算匹配扫描器过滤器(scanner filter)的行数。

head(self, int num_rows[, columns])

加载片段的前 N 行。

scanner(self, Schema schema=None[, columns])

针对该片段构建扫描操作。

take(self, indices[, columns])

按索引选择数据行。

to_batches(self, Schema schema=None[, columns])

将片段作为已具体化的记录批次(record batches)读取。

to_table(self, Schema schema=None[, columns])

将此片段转换为表(Table)。

属性

partition_expression

一个表达式,对于该片段所查看的所有数据,其计算结果均为真。

physical_schema

返回此片段的物理模式(schema)。

count_rows(self, Expression filter=None, int batch_size=_DEFAULT_BATCH_SIZE, int batch_readahead=_DEFAULT_BATCH_READAHEAD, int fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD, FragmentScanOptions fragment_scan_options=None, bool use_threads=True, bool cache_metadata=True, MemoryPool memory_pool=None)#

计算匹配扫描器过滤器(scanner filter)的行数。

参数:
filterExpression, default None

扫描将仅返回与过滤器匹配的行。如果可能,谓词将被下推以利用分区信息或在数据源中找到的内部元数据,例如 Parquet 统计信息。否则,它会在产生记录批次之前过滤已加载的记录批次。

batch_sizeint, default 131_072

扫描的记录批次的最大行数。如果扫描的记录批次溢出内存,则可以调用此方法来减小其大小。

batch_readaheadint, default 16

在一个文件中预读的批次数量。这可能不适用于所有文件格式。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。

fragment_readaheadint, default 4

预读的文件数量。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。

fragment_scan_optionsFragmentScanOptions, default None

特定于特定扫描和片段类型的选项,在同一数据集的不同扫描之间可能会有所不同。

use_threadsbool, 默认 True

如果启用,将使用由可用 CPU 核心数确定的最大并行度。

cache_metadatabool, default True

如果启用,扫描时可能会缓存元数据以加快重复扫描。

memory_poolMemoryPool, 默认 None

如果需要,用于内存分配。如果未指定,则使用默认池。

返回:
countint
head(self, int num_rows, columns=None, Expression filter=None, int batch_size=_DEFAULT_BATCH_SIZE, int batch_readahead=_DEFAULT_BATCH_READAHEAD, int fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD, FragmentScanOptions fragment_scan_options=None, bool use_threads=True, bool cache_metadata=True, MemoryPool memory_pool=None)#

加载片段的前 N 行。

参数:
num_rowsint

要加载的行数。

columnslist of str, default None

要投影的列。这可以是要包含的列名列表(顺序和重复项将保留),也可以是包含 {new_column_name: expression} 值的字典,用于更高级的投影。

列或表达式列表可以使用特殊字段 __batch_index(片段中批次的索引)、__fragment_index(数据集中片段的索引)、__last_in_fragment(批次是否是片段中的最后一个)和 __filename(源文件的名称或源片段的描述)。

列将被传递给数据集和相应的数据片段,以避免加载、复制和反序列化计算链下游不需要的列。默认情况下,所有可用列都会被投影。如果引用的任何列名不存在于数据集的模式中,则会引发异常。

filterExpression, default None

扫描将仅返回与过滤器匹配的行。如果可能,谓词将被下推以利用分区信息或在数据源中找到的内部元数据,例如 Parquet 统计信息。否则,它会在产生记录批次之前过滤已加载的记录批次。

batch_sizeint, default 131_072

扫描的记录批次的最大行数。如果扫描的记录批次溢出内存,则可以调用此方法来减小其大小。

batch_readaheadint, default 16

在一个文件中预读的批次数量。这可能不适用于所有文件格式。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。

fragment_readaheadint, default 4

预读的文件数量。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。

fragment_scan_optionsFragmentScanOptions, default None

特定于特定扫描和片段类型的选项,在同一数据集的不同扫描之间可能会有所不同。

use_threadsbool, 默认 True

如果启用,将使用由可用 CPU 核心数确定的最大并行度。

cache_metadatabool, default True

如果启用,扫描时可能会缓存元数据以加快重复扫描。

memory_poolMemoryPool, 默认 None

如果需要,用于内存分配。如果未指定,则使用默认池。

返回:
partition_expression#

一个表达式,对于该片段所查看的所有数据,其计算结果均为真。

physical_schema#

返回此片段的物理模式。此模式可能与数据集读取模式不同。

scanner(self, Schema schema=None, columns=None, Expression filter=None, int batch_size=_DEFAULT_BATCH_SIZE, int batch_readahead=_DEFAULT_BATCH_READAHEAD, int fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD, FragmentScanOptions fragment_scan_options=None, bool use_threads=True, bool cache_metadata=True, MemoryPool memory_pool=None)#

针对该片段构建扫描操作。

数据不会立即加载。相反,这将生成一个扫描器(Scanner),它公开了进一步的操作(例如,将所有数据加载为表,计算行数)。

参数:
schemaSchema

扫描时使用的模式。这用于将片段统一到其数据集的模式。如果未指定,将使用该片段的物理模式,每个片段的物理模式可能不同。

columnslist of str, default None

要投影的列。这可以是要包含的列名列表(顺序和重复项将保留),也可以是包含 {new_column_name: expression} 值的字典,用于更高级的投影。

列或表达式列表可以使用特殊字段 __batch_index(片段中批次的索引)、__fragment_index(数据集中片段的索引)、__last_in_fragment(批次是否是片段中的最后一个)和 __filename(源文件的名称或源片段的描述)。

列将被传递给数据集和相应的数据片段,以避免加载、复制和反序列化计算链下游不需要的列。默认情况下,所有可用列都会被投影。如果引用的任何列名不存在于数据集的模式中,则会引发异常。

filterExpression, default None

扫描将仅返回与过滤器匹配的行。如果可能,谓词将被下推以利用分区信息或在数据源中找到的内部元数据,例如 Parquet 统计信息。否则,它会在产生记录批次之前过滤已加载的记录批次。

batch_sizeint, default 131_072

扫描的记录批次的最大行数。如果扫描的记录批次溢出内存,则可以调用此方法来减小其大小。

batch_readaheadint, default 16

在一个文件中预读的批次数量。这可能不适用于所有文件格式。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。

fragment_readaheadint, default 4

预读的文件数量。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。

fragment_scan_optionsFragmentScanOptions, default None

特定于特定扫描和片段类型的选项,在同一数据集的不同扫描之间可能会有所不同。

use_threadsbool, 默认 True

如果启用,将使用由可用 CPU 核心数确定的最大并行度。

cache_metadatabool, default True

如果启用,扫描时可能会缓存元数据以加快重复扫描。

memory_poolMemoryPool, 默认 None

如果需要,用于内存分配。如果未指定,则使用默认池。

返回:
scannerScanner
take(self, indices, columns=None, Expression filter=None, int batch_size=_DEFAULT_BATCH_SIZE, int batch_readahead=_DEFAULT_BATCH_READAHEAD, int fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD, FragmentScanOptions fragment_scan_options=None, bool use_threads=True, bool cache_metadata=True, MemoryPool memory_pool=None)#

按索引选择数据行。

参数:
indicesArrayarray-like

数据集中要选择的行索引。

columnslist of str, default None

要投影的列。这可以是要包含的列名列表(顺序和重复项将保留),也可以是包含 {new_column_name: expression} 值的字典,用于更高级的投影。

列或表达式列表可以使用特殊字段 __batch_index(片段中批次的索引)、__fragment_index(数据集中片段的索引)、__last_in_fragment(批次是否是片段中的最后一个)和 __filename(源文件的名称或源片段的描述)。

列将被传递给数据集和相应的数据片段,以避免加载、复制和反序列化计算链下游不需要的列。默认情况下,所有可用列都会被投影。如果引用的任何列名不存在于数据集的模式中,则会引发异常。

filterExpression, default None

扫描将仅返回与过滤器匹配的行。如果可能,谓词将被下推以利用分区信息或在数据源中找到的内部元数据,例如 Parquet 统计信息。否则,它会在产生记录批次之前过滤已加载的记录批次。

batch_sizeint, default 131_072

扫描的记录批次的最大行数。如果扫描的记录批次溢出内存,则可以调用此方法来减小其大小。

batch_readaheadint, default 16

在一个文件中预读的批次数量。这可能不适用于所有文件格式。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。

fragment_readaheadint, default 4

预读的文件数量。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。

fragment_scan_optionsFragmentScanOptions, default None

特定于特定扫描和片段类型的选项,在同一数据集的不同扫描之间可能会有所不同。

use_threadsbool, 默认 True

如果启用,将使用由可用 CPU 核心数确定的最大并行度。

cache_metadatabool, default True

如果启用,扫描时可能会缓存元数据以加快重复扫描。

memory_poolMemoryPool, 默认 None

如果需要,用于内存分配。如果未指定,则使用默认池。

返回:
to_batches(self, Schema schema=None, columns=None, Expression filter=None, int batch_size=_DEFAULT_BATCH_SIZE, int batch_readahead=_DEFAULT_BATCH_READAHEAD, int fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD, FragmentScanOptions fragment_scan_options=None, bool use_threads=True, bool cache_metadata=True, MemoryPool memory_pool=None)#

将片段作为已具体化的记录批次(record batches)读取。

参数:
schemaSchema,可选

扫描时使用的具体模式。

columnslist of str, default None

要投影的列。这可以是要包含的列名列表(顺序和重复项将保留),也可以是包含 {new_column_name: expression} 值的字典,用于更高级的投影。

列或表达式列表可以使用特殊字段 __batch_index(片段中批次的索引)、__fragment_index(数据集中片段的索引)、__last_in_fragment(批次是否是片段中的最后一个)和 __filename(源文件的名称或源片段的描述)。

列将被传递给数据集和相应的数据片段,以避免加载、复制和反序列化计算链下游不需要的列。默认情况下,所有可用列都会被投影。如果引用的任何列名不存在于数据集的模式中,则会引发异常。

filterExpression, default None

扫描将仅返回与过滤器匹配的行。如果可能,谓词将被下推以利用分区信息或在数据源中找到的内部元数据,例如 Parquet 统计信息。否则,它会在产生记录批次之前过滤已加载的记录批次。

batch_sizeint, default 131_072

扫描的记录批次的最大行数。如果扫描的记录批次溢出内存,则可以调用此方法来减小其大小。

batch_readaheadint, default 16

在一个文件中预读的批次数量。这可能不适用于所有文件格式。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。

fragment_readaheadint, default 4

预读的文件数量。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。

fragment_scan_optionsFragmentScanOptions, default None

特定于特定扫描和片段类型的选项,在同一数据集的不同扫描之间可能会有所不同。

use_threadsbool, 默认 True

如果启用,将使用由可用 CPU 核心数确定的最大并行度。

cache_metadatabool, default True

如果启用,扫描时可能会缓存元数据以加快重复扫描。

memory_poolMemoryPool, 默认 None

如果需要,用于内存分配。如果未指定,则使用默认池。

返回:
record_batchesiterator of RecordBatch
to_table(self, Schema schema=None, columns=None, Expression filter=None, int batch_size=_DEFAULT_BATCH_SIZE, int batch_readahead=_DEFAULT_BATCH_READAHEAD, int fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD, FragmentScanOptions fragment_scan_options=None, bool use_threads=True, bool cache_metadata=True, MemoryPool memory_pool=None)#

将此片段转换为表(Table)。

谨慎使用此便捷工具。它会在创建表之前在内存中串行具体化扫描结果。

参数:
schemaSchema,可选

扫描时使用的具体模式。

columnslist of str, default None

要投影的列。这可以是要包含的列名列表(顺序和重复项将保留),也可以是包含 {new_column_name: expression} 值的字典,用于更高级的投影。

列或表达式列表可以使用特殊字段 __batch_index(片段中批次的索引)、__fragment_index(数据集中片段的索引)、__last_in_fragment(批次是否是片段中的最后一个)和 __filename(源文件的名称或源片段的描述)。

列将被传递给数据集和相应的数据片段,以避免加载、复制和反序列化计算链下游不需要的列。默认情况下,所有可用列都会被投影。如果引用的任何列名不存在于数据集的模式中,则会引发异常。

filterExpression, default None

扫描将仅返回与过滤器匹配的行。如果可能,谓词将被下推以利用分区信息或在数据源中找到的内部元数据,例如 Parquet 统计信息。否则,它会在产生记录批次之前过滤已加载的记录批次。

batch_sizeint, default 131_072

扫描的记录批次的最大行数。如果扫描的记录批次溢出内存,则可以调用此方法来减小其大小。

batch_readaheadint, default 16

在一个文件中预读的批次数量。这可能不适用于所有文件格式。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。

fragment_readaheadint, default 4

预读的文件数量。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。

fragment_scan_optionsFragmentScanOptions, default None

特定于特定扫描和片段类型的选项,在同一数据集的不同扫描之间可能会有所不同。

use_threadsbool, 默认 True

如果启用,将使用由可用 CPU 核心数确定的最大并行度。

cache_metadatabool, default True

如果启用,扫描时可能会缓存元数据以加快重复扫描。

memory_poolMemoryPool, 默认 None

如果需要,用于内存分配。如果未指定,则使用默认池。

返回:
tableTable