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)。
属性
一个表达式,对于该片段所查看的所有数据,其计算结果均为真。
返回此片段的物理模式(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)的行数。
- 参数:
- filter
Expression, defaultNone 扫描将仅返回与过滤器匹配的行。如果可能,谓词将被下推以利用分区信息或在数据源中找到的内部元数据,例如 Parquet 统计信息。否则,它会在产生记录批次之前过滤已加载的记录批次。
- batch_size
int, default 131_072 扫描的记录批次的最大行数。如果扫描的记录批次溢出内存,则可以调用此方法来减小其大小。
- batch_readahead
int, default 16 在一个文件中预读的批次数量。这可能不适用于所有文件格式。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。
- fragment_readahead
int, default 4 预读的文件数量。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。
- fragment_scan_options
FragmentScanOptions, defaultNone 特定于特定扫描和片段类型的选项,在同一数据集的不同扫描之间可能会有所不同。
- use_threadsbool, 默认
True 如果启用,将使用由可用 CPU 核心数确定的最大并行度。
- cache_metadatabool, default
True 如果启用,扫描时可能会缓存元数据以加快重复扫描。
- memory_pool
MemoryPool, 默认None 如果需要,用于内存分配。如果未指定,则使用默认池。
- filter
- 返回:
- count
int
- count
- 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_rows
int 要加载的行数。
- columns
listofstr, defaultNone 要投影的列。这可以是要包含的列名列表(顺序和重复项将保留),也可以是包含 {new_column_name: expression} 值的字典,用于更高级的投影。
列或表达式列表可以使用特殊字段 __batch_index(片段中批次的索引)、__fragment_index(数据集中片段的索引)、__last_in_fragment(批次是否是片段中的最后一个)和 __filename(源文件的名称或源片段的描述)。
列将被传递给数据集和相应的数据片段,以避免加载、复制和反序列化计算链下游不需要的列。默认情况下,所有可用列都会被投影。如果引用的任何列名不存在于数据集的模式中,则会引发异常。
- filter
Expression, defaultNone 扫描将仅返回与过滤器匹配的行。如果可能,谓词将被下推以利用分区信息或在数据源中找到的内部元数据,例如 Parquet 统计信息。否则,它会在产生记录批次之前过滤已加载的记录批次。
- batch_size
int, default 131_072 扫描的记录批次的最大行数。如果扫描的记录批次溢出内存,则可以调用此方法来减小其大小。
- batch_readahead
int, default 16 在一个文件中预读的批次数量。这可能不适用于所有文件格式。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。
- fragment_readahead
int, default 4 预读的文件数量。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。
- fragment_scan_options
FragmentScanOptions, defaultNone 特定于特定扫描和片段类型的选项,在同一数据集的不同扫描之间可能会有所不同。
- use_threadsbool, 默认
True 如果启用,将使用由可用 CPU 核心数确定的最大并行度。
- cache_metadatabool, default
True 如果启用,扫描时可能会缓存元数据以加快重复扫描。
- memory_pool
MemoryPool, 默认None 如果需要,用于内存分配。如果未指定,则使用默认池。
- num_rows
- 返回:
- 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),它公开了进一步的操作(例如,将所有数据加载为表,计算行数)。
- 参数:
- schema
Schema 扫描时使用的模式。这用于将片段统一到其数据集的模式。如果未指定,将使用该片段的物理模式,每个片段的物理模式可能不同。
- columns
listofstr, defaultNone 要投影的列。这可以是要包含的列名列表(顺序和重复项将保留),也可以是包含 {new_column_name: expression} 值的字典,用于更高级的投影。
列或表达式列表可以使用特殊字段 __batch_index(片段中批次的索引)、__fragment_index(数据集中片段的索引)、__last_in_fragment(批次是否是片段中的最后一个)和 __filename(源文件的名称或源片段的描述)。
列将被传递给数据集和相应的数据片段,以避免加载、复制和反序列化计算链下游不需要的列。默认情况下,所有可用列都会被投影。如果引用的任何列名不存在于数据集的模式中,则会引发异常。
- filter
Expression, defaultNone 扫描将仅返回与过滤器匹配的行。如果可能,谓词将被下推以利用分区信息或在数据源中找到的内部元数据,例如 Parquet 统计信息。否则,它会在产生记录批次之前过滤已加载的记录批次。
- batch_size
int, default 131_072 扫描的记录批次的最大行数。如果扫描的记录批次溢出内存,则可以调用此方法来减小其大小。
- batch_readahead
int, default 16 在一个文件中预读的批次数量。这可能不适用于所有文件格式。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。
- fragment_readahead
int, default 4 预读的文件数量。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。
- fragment_scan_options
FragmentScanOptions, defaultNone 特定于特定扫描和片段类型的选项,在同一数据集的不同扫描之间可能会有所不同。
- use_threadsbool, 默认
True 如果启用,将使用由可用 CPU 核心数确定的最大并行度。
- cache_metadatabool, default
True 如果启用,扫描时可能会缓存元数据以加快重复扫描。
- memory_pool
MemoryPool, 默认None 如果需要,用于内存分配。如果未指定,则使用默认池。
- schema
- 返回:
- scanner
Scanner
- scanner
- 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)#
按索引选择数据行。
- 参数:
- indices
Array或array-like 数据集中要选择的行索引。
- columns
listofstr, defaultNone 要投影的列。这可以是要包含的列名列表(顺序和重复项将保留),也可以是包含 {new_column_name: expression} 值的字典,用于更高级的投影。
列或表达式列表可以使用特殊字段 __batch_index(片段中批次的索引)、__fragment_index(数据集中片段的索引)、__last_in_fragment(批次是否是片段中的最后一个)和 __filename(源文件的名称或源片段的描述)。
列将被传递给数据集和相应的数据片段,以避免加载、复制和反序列化计算链下游不需要的列。默认情况下,所有可用列都会被投影。如果引用的任何列名不存在于数据集的模式中,则会引发异常。
- filter
Expression, defaultNone 扫描将仅返回与过滤器匹配的行。如果可能,谓词将被下推以利用分区信息或在数据源中找到的内部元数据,例如 Parquet 统计信息。否则,它会在产生记录批次之前过滤已加载的记录批次。
- batch_size
int, default 131_072 扫描的记录批次的最大行数。如果扫描的记录批次溢出内存,则可以调用此方法来减小其大小。
- batch_readahead
int, default 16 在一个文件中预读的批次数量。这可能不适用于所有文件格式。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。
- fragment_readahead
int, default 4 预读的文件数量。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。
- fragment_scan_options
FragmentScanOptions, defaultNone 特定于特定扫描和片段类型的选项,在同一数据集的不同扫描之间可能会有所不同。
- use_threadsbool, 默认
True 如果启用,将使用由可用 CPU 核心数确定的最大并行度。
- cache_metadatabool, default
True 如果启用,扫描时可能会缓存元数据以加快重复扫描。
- memory_pool
MemoryPool, 默认None 如果需要,用于内存分配。如果未指定,则使用默认池。
- indices
- 返回:
- 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)读取。
- 参数:
- schema
Schema,可选 扫描时使用的具体模式。
- columns
listofstr, defaultNone 要投影的列。这可以是要包含的列名列表(顺序和重复项将保留),也可以是包含 {new_column_name: expression} 值的字典,用于更高级的投影。
列或表达式列表可以使用特殊字段 __batch_index(片段中批次的索引)、__fragment_index(数据集中片段的索引)、__last_in_fragment(批次是否是片段中的最后一个)和 __filename(源文件的名称或源片段的描述)。
列将被传递给数据集和相应的数据片段,以避免加载、复制和反序列化计算链下游不需要的列。默认情况下,所有可用列都会被投影。如果引用的任何列名不存在于数据集的模式中,则会引发异常。
- filter
Expression, defaultNone 扫描将仅返回与过滤器匹配的行。如果可能,谓词将被下推以利用分区信息或在数据源中找到的内部元数据,例如 Parquet 统计信息。否则,它会在产生记录批次之前过滤已加载的记录批次。
- batch_size
int, default 131_072 扫描的记录批次的最大行数。如果扫描的记录批次溢出内存,则可以调用此方法来减小其大小。
- batch_readahead
int, default 16 在一个文件中预读的批次数量。这可能不适用于所有文件格式。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。
- fragment_readahead
int, default 4 预读的文件数量。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。
- fragment_scan_options
FragmentScanOptions, defaultNone 特定于特定扫描和片段类型的选项,在同一数据集的不同扫描之间可能会有所不同。
- use_threadsbool, 默认
True 如果启用,将使用由可用 CPU 核心数确定的最大并行度。
- cache_metadatabool, default
True 如果启用,扫描时可能会缓存元数据以加快重复扫描。
- memory_pool
MemoryPool, 默认None 如果需要,用于内存分配。如果未指定,则使用默认池。
- schema
- 返回:
- record_batchesiterator of
RecordBatch
- record_batchesiterator of
- 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)。
谨慎使用此便捷工具。它会在创建表之前在内存中串行具体化扫描结果。
- 参数:
- schema
Schema,可选 扫描时使用的具体模式。
- columns
listofstr, defaultNone 要投影的列。这可以是要包含的列名列表(顺序和重复项将保留),也可以是包含 {new_column_name: expression} 值的字典,用于更高级的投影。
列或表达式列表可以使用特殊字段 __batch_index(片段中批次的索引)、__fragment_index(数据集中片段的索引)、__last_in_fragment(批次是否是片段中的最后一个)和 __filename(源文件的名称或源片段的描述)。
列将被传递给数据集和相应的数据片段,以避免加载、复制和反序列化计算链下游不需要的列。默认情况下,所有可用列都会被投影。如果引用的任何列名不存在于数据集的模式中,则会引发异常。
- filter
Expression, defaultNone 扫描将仅返回与过滤器匹配的行。如果可能,谓词将被下推以利用分区信息或在数据源中找到的内部元数据,例如 Parquet 统计信息。否则,它会在产生记录批次之前过滤已加载的记录批次。
- batch_size
int, default 131_072 扫描的记录批次的最大行数。如果扫描的记录批次溢出内存,则可以调用此方法来减小其大小。
- batch_readahead
int, default 16 在一个文件中预读的批次数量。这可能不适用于所有文件格式。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。
- fragment_readahead
int, default 4 预读的文件数量。增加此数字会增加 RAM 使用量,但也可以提高 IO 利用率。
- fragment_scan_options
FragmentScanOptions, defaultNone 特定于特定扫描和片段类型的选项,在同一数据集的不同扫描之间可能会有所不同。
- use_threadsbool, 默认
True 如果启用,将使用由可用 CPU 核心数确定的最大并行度。
- cache_metadatabool, default
True 如果启用,扫描时可能会缓存元数据以加快重复扫描。
- memory_pool
MemoryPool, 默认None 如果需要,用于内存分配。如果未指定,则使用默认池。
- schema
- 返回:
- table
Table
- table