ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Arrow Python Dataset API 详解:pyarrow.dataset 的工厂函数、核心类与读写实现

Apache Arrow Python Dataset API 详解:pyarrow.dataset 的工厂函数、核心类与读写实现 Apache Arrow Python Dataset API 详解:pyarrow.dataset 的工厂函数、核心类与读写实现【免费下载链接】arrowApache Arrow is the universal columnar format and multi-language toolbox for fast data interchange and in-memory analytics项目地址: https://gitcode.com/GitHub_Trending/arrow3/arrow本文以 Apache Arrow 官方文档 Dataset API 参考页 为主线,完整梳理pyarrow.dataset模块提供的工厂函数(dataset、parquet_dataset、partitioning、write_dataset等)、核心类(Dataset、FileSystemDataset、Scanner、各文件格式与分区方案等)以及辅助函数get_partition_keys,并结合源码 dataset.py、_dataset.pyx 和 _dataset_parquet.pyx 深入解析其参数取值、默认值与底层调用链,帮助读者掌握在 Python 中读写跨内存、多文件、分区化数据集的完整方案。需要说明的是,dataset.py 模块头部明确标注 Dataset is currently unstable. APIs subject to change without notice,即 Dataset API 仍属于不稳定接口,本文内容以当前仓库版本为准。一、Dataset API 的总体定位与模块结构pyarrow.dataset的目标是为可能大于内存的多文件表格数据提供统一接口,其官方文档(dataset.py 中dataset()的 docstring)概括了三大能力:统一的数据源接口:Parquet、Feather(Ipc)、CSV、JSON、ORC 等格式共用同一套打开与扫描 API;数据源发现:递归爬取目录、识别基于目录的分区数据集、做基础的 schema 归一化;优化读取:谓词下推(行过滤)、投影(列裁剪)、并行读取或细粒度任务管理。从源码结构看,该模块由三层组成:层次文件职责Python 高层封装dataset.py(共 1040 行)工厂函数dataset/parquet_dataset/partitioning/write_dataset的参数校验与分发Cython 核心绑定_dataset.pyx(共 4240 行)Dataset、Fragment、Scanner、分区方案、CSV/Ipc/JSON 格式等类的 C 绑定格式扩展绑定_dataset_parquet.pyx、_dataset_orc.pyxParquet 与 ORC 的读写选项与工厂dataset.py 展示了可选依赖的加载策略:核心类从pyarrow._dataset导入,而OrcFileFormat与 Parquet 相关类分别从pyarrow._dataset_orc、pyarrow._dataset_parquet可选导入;若对应扩展未编译,访问这些类时会通过__getattr__抛出明确的ImportError(如 The pyarrow installation is not built with support for the ORC file format)。这一点决定了后文各格式类的使用前提。文档参考页中列出的field与scalar属于pyarrow.compute模块,在 dataset.py 中仅为向后兼容而从pyarrow.compute重新导出:# keep Expression functionality exposed here for backwards compatibility from pyarrow.compute import Expression, scalar, field因此构造表达式过滤条件时,ds.field(...)、ds.scalar(...)与pa.field(...)、pa.scalar(...)是同一批对象。二、工厂函数:partitioning() 定义分区方案参考页 Factory functions 一节列出的第一个核心函数是partitioning,源码位于 dataset.py。它支持三种分区方案,返回值取决于参数组合:2.1 三种分区方案DirectoryPartitioning(默认):文件路径中每个段对应 schema 的一个字段,且所有字段都必须出现。例如 schemayear:int16, month:int8,路径/2009/11解析为year 2009 and month 11;HivePartitioning(flavorhive):Hive 风格的/keyvalue/嵌套目录。字段顺序无关,未识别的键被忽略。例如 schemayear:int16, month:int8, day:int8,路径/year2009/month11/day15是合法的,即使与 schema 字段顺序不一致;FilenamePartitioning(flavorfilename):分区值直接编码在文件名中,字段值之间用_分隔。例如2009_11_part-0.parquet解析为year 2009 and month 11。2.2 参数与返回类型partitioning(schemaNone, field_namesNone, flavorNone, dictionariesNone)的语义在 dataset.py 的 docstring 中有完整说明:schema:描述路径中分区的 Schema;若未提供而给了field_names/flavor,则类型从文件路径推断,此时返回的是PartitioningFactory而非确定的Partitioning;field_names:字段名列表,仅对 Directory 方案有效,类型同样从路径推断;flavor:缺省为目录分区,hive表示 Hive 分区,filename表示文件名分区;dictionaries:若分区字段是字典类型,必须提供包含该列全部可能取值的数组,否则解析会报错;也可以传字符串infer让 Arrow 自行发现字典值(此时返回PartitioningFactory)。注意 dataset.py 中的互斥约束:Directory 方案下不能同时给schema与field_names;Hive 方案下不允许field_names;不支持的 flavor 会抛出Unsupported flavor。文档中的示例可完整继承:import pyarrow as pa import pyarrow.dataset as ds # 显式 Schema:路径形如 /2009/June part ds.partitioning(pa.schema([(year, pa.int16()), (month, pa.string())])) # 仅给字段名,类型由路径推断(year 推断为 int32,month 推断为 string) part ds.partitioning(field_names[year, month]) # 字典编码分区:显式提供字典值 part ds.partitioning( pa.schema([ (year, pa.int16()), (month, pa.dictionary(pa.int8(), pa.string())) ]), dictionaries{ month: pa.array([January, February, March]), }) # 字典编码分区:让 Arrow 推断字典值 part ds.partitioning( pa.schema([ (year, pa.int16()), (month, pa.dictionary(pa.int8(), pa.string())) ]), dictionariesinfer) # Hive 方案:路径形如 /year2009/month11 part ds.partitioning( pa.schema([(year, pa.int16()), (month, pa.int8())]), flavorhive) # Hive 方案由目录结构自动发现(类型一并推断) part ds.partitioning(flavorhive)对应的底层类定义在 _dataset.pyx 中:Partitioning(L2548)→KeyValuePartitioning(L2692)→ 三个具体子类DirectoryPartitioning(L2742)、HivePartitioning(L2870)、FilenamePartitioning(L3021);PartitioningFactory(L2642)则封装了先扫描目录再确定类型/字典的延迟发现流程。三、工厂函数:dataset() 打开数据源dataset()是整个模块的入口,位于 dataset.py。参考页中它与parquet_dataset、write_dataset并列为工厂函数,docstring 完整覆盖了source的五种形态:source 类型行为单个文件路径从单文件打开FileSystemDataset目录路径递归发现,如指定分区方案则按分区解析文件路径列表从显式文件列表构造;所有文件必须位于filesystem参数指定的同一文件系统,且不允许以 URI 形式传路径Dataset 列表构造嵌套的UnionDataset,不允许再传其他关键字参数Table/RecordBatch(列表)、batches 可迭代对象、RecordBatchReader构造InMemoryDataset;可迭代对象或空列表必须同时给 schema,且可迭代/Reader 来源的数据集只能扫描一次3.1 关键参数说明schema:可选,显式提供后不再从源推断;format:字符串或FileFormat实例。_ensure_format() 中确认当前支持的字符串取值为parquet、ipc/arrow、feather、csv、json、orc,其中 Feather 仅支持 v2 文件;ORC 与 Parquet 分别要求对应的编译扩展;filesystem:可为FileSystem对象或 URI 字符串;URI 的 path 部分会作为目录前缀(等价于SubTreeFileSystem),Windows 下必须使用file:///C:...或file:/C:...形式;partitioning:可传Partitioning/PartitioningFactory对象、flavor 字符串(如hive)或字段名列表(等价于目录分区推断);partition_base_dir:应用分区时路径会先剥离该前缀;不匹配前缀的文件仍属于数据集,只是不带分区信息;exclude_invalid_files:默认 False;置 True 时会逐个文件串行做格式合法性检查(产生额外 IO),关闭则可能把不支持的文件留在数据集中、直到扫描时才报错;ignore_prefixes:发现过程忽略匹配这些前缀的文件(与路径 basename 匹配),默认[., _],即隐藏文件和下划线开头的文件(如_metadata)默认不参与发现。3.2 文档示例(完整继承)以下为 dataset.py 中 docstring 的官方示例,展示了从单文件到 S3、从显式 schema 到嵌套 UnionDataset 的典型用法:import pyarrow as pa import pyarrow.parquet as pq import pyarrow.dataset as ds table pa.table({year: [2020, 2022, 2021, 2022, 2019, 2021], n_legs: [2, 2, 4, 4, 5, 100], animal: [Flamingo, Parrot, Dog, Horse, Brittle stars, Centipede]}) pq.write_table(table, file.parquet) # 打开单个文件 dataset ds.dataset(file.parquet, formatparquet) dataset.to_table() # 打开单个文件并显式指定 schema(等价于投影列子集) myschema pa.schema([(n_legs, pa.int64()), (animal, pa.string())]) dataset ds.dataset(file.parquet, schemamyschema, formatparquet) dataset.to_table() # 打开分区目录 ds.write_dataset(table, partitioned_dataset, formatparquet, partitioning[year]) dataset ds.dataset(partitioned_dataset, formatparquet) # S3 桶中的目录(需要相应的凭证配置) ds.dataset(s3://mybucket/nyc-taxi/, formatparquet) # 从相对路径文件列表打开 dataset ds.dataset([ partitioned_dataset/2019/part-0.parquet, partitioned_dataset/2020/part-0.parquet, partitioned_dataset/2021/part-0.parquet, ], formatparquet) # 文件列表 filesystem URI 前缀 paths [part0/data.parquet, part1/data.parquet, part3/data.parquet] ds.dataset(paths, filesystems3://bucket/nested/directory, formatparquet) # 嵌套 UnionDataset:组合任意其他数据集 ds.dataset([ ds.dataset(s3://old-taxi-data, formatparquet), ds.dataset(local/path/to/data, formatipc) ])3.3 底层调用链从源码看,dataset()只做类型分派(dataset.py):路径或路径列表走 _filesystem_dataset(),Table/RecordBatch 列表走 _in_memory_dataset(),Dataset 列表走 _union_dataset()。其中:单路径经 _ensure_single_source() 解析:目录返回递归的FileSelector,单文件返回单元素列表,不存在则抛FileNotFoundError;路径列表经 _ensure_multiple_sources() 校验:本地文件系统下会逐一检查路径必须是真实文件,目录会抛IsADirectoryError并提示要构造嵌套或 union 数据集请传 Dataset 对象列表;最终把partitioning、partition_base_dir、exclude_invalid_files、selector_ignore_prefixes打包进FileSystemFactoryOptions,交给FileSystemDatasetFactory.finish(schema)完成 schema 推断与数据集构建。UnionDataset有一条值得注意的限制:在 _union_dataset() 中,子数据集若带过滤/投影(_scan_options非空)会直接抛错,官方建议先 union 再统一施加 filter;同时未指定schema时会自动pa.unify_schemas统一子数据集 schema。四、工厂函数:parquet_dataset() 从 _metadata 文件打开parquet_dataset()位于 dataset.py,专门从pyarrow.parquet.write_metadata生成的_metadata文件创建FileSystemDataset,免去逐文件扫描。参数:metadata_path:指向单个 Parquet 元数据文件的路径;schema:可选,提供后不从源推断;filesystem:同dataset()的语义,缺省为LocalFileSystem;format:必须是ParquetFileFormat实例(默认新建一个),否则抛ValueError;partitioning/partition_base_dir:语义与dataset()相同。内部实现用ParquetFactoryOptions携带partition_base_dir与分区方案,交由 ParquetDatasetFactory 的finish(schema)返回数据集。这条工厂链(ParquetFactoryOptions在 _dataset_parquet.pyx、ParquetDatasetFactory在 L1087)正是参考页 Classes 一节中ParquetFileFormat、ParquetReadOptions等类被组合使用的地方。五、工厂函数:write_dataset() 写入与分区写出write_dataset()位于 dataset.py,是把 Table/RecordBatch、Scanner 或 Dataset 写出为指定格式与分区结构的入口。参数较多,按官方 docstring 逐项说明:data:Dataset、Table/RecordBatch、RecordBatchReader、Table/RecordBatch 列表或 RecordBatch 可迭代对象(可迭代对象必须同时给schema);base_dir:写出根目录;basename_template:文件名模板,{i}会被自动递增的整数替换;缺省为part-{i}. format.default_extname;format:支持parquet、ipc/arrow/feather、csv;写入FileSystemDataset且未指定 format 时,沿用源数据集格式;写 Table/RecordBatch 时该参数必填;partitioning:分区对象或字段名列表,配合partitioning_flavor选择方案类型(缺省为目录分区);schema、filesystem、file_options(FileFormat.make_write_options()创建,格式相关的写选项);use_threads:默认 True,按 CPU 核数并行写文件,但可能打乱行序;preserve_order:默认 False,置 True 可保证多线程下仍保序,可能带来明显性能损耗;max_partitions:默认 1024,单个 batch 最多写入的分区数;max_open_files:默认 1024,限制同时打开的文件数,超出时关闭最近最少使用的文件;设置过低会把数据碎片化成大量小文件;max_rows_per_file:默认 0(不限制),大于 0 时限制单文件行数,否则每个输出目录一个文件(除非需要关闭文件以满足max_open_files);min_rows_per_group:默认 0,大于 0 时写出器先攒批,行数足够才写 row group;max_rows_per_group:默认 1024 × 1024,大于 0 时可把大批次拆成多个 row group;docstring 同时提示此时应一并设置min_rows_per_group,否则可能产生过小 row group;file_visitor:每创建一个文件就以WrittenFile实例回调,其path属性为文件路径,metadata属性为 Parquet 文件元数据(非 Parquet 格式为 None),可用于构建_metadata文件;existing_data_behavior:error(默认,目标已有数据即报错)|overwrite_or_ignore(同名文件覆盖、其余忽略,配合唯一basename_template可实现追加工作流)|delete_matching(首次遇到分区目录时整个删除,用于完整覆盖旧分区);create_dir:默认 True,置 False 时不创建目录,适用于不需要目录概念的文件系统。默认值并非只写在 docstring 里,代码在 dataset.py 中有对应落实:max_partitions与max_open_files为 1024、max_rows_per_file为 0、max_rows_per_group为1 20、min_rows_per_group为 0。函数尾部统一收敛到 Cython 的_filesystemdataset_write(导入自pyarrow._dataset,见 dataset.py);WrittenFile类定义在 _dataset.pyx。官方回调示例可直接照抄:visited_paths [] def file_visitor(written_file): visited_paths.append(written_file.path) ds.write_dataset(table, out_dir, formatparquet, partitioning[year], file_visitorfile_visitor, existing_data_behavioroverwrite_or_ignore)六、核心类:Dataset 体系与 Fragment参考页 Classes 一节的核心是数据集本体。在 _dataset.pyx 中,类继承关系为:Dataset(L162,抽象基类)→ 三个具体实现:InMemoryDataset(L996):由 Table/RecordBatch 或 batches 迭代构造,只存在于内存;UnionDataset(L1057):由多个子数据集组合,支持任意嵌套;FileSystemDataset(L1100):由文件系统 文件/选择器 FileFormat 发现选项构造,是最常用的形态;Fragment(L1436,数据的物理切分单元)→FileFragment(L1977)→ParquetFileFragment(_dataset_parquet.pyx),后者暴露 Parquet 特有的physical_row_groups、fragment_scan_options等元信息(配套的RowGroupInfo命名元组定义在 _dataset_parquet.pyx);TaggedRecordBatch(_dataset.pyx):带标签的 RecordBatch,是 fragment 级扫描的产出,可被 Scanner 复用。围绕发现与构建的还有:FileSystemFactoryOptions(L3241,承载 partitioning、partition_base_dir、exclude_invalid_files、selector_ignore_prefixes)、FileSystemDatasetFactory(L3362,finish(schema)返回数据集)、UnionDatasetFactory(L3437)、FragmentScanOptions(L2080,格式无关的 fragment 级扫描选项,派生出CsvFragmentScanOptions、ParquetFragmentScanOptions等)。七、核心类:FileFormat 与各格式扩展FileFormat基类定义在 _dataset.pyx,提供make_write_options()等通用接口。各格式绑定类与参考页的对应关系:参考页类名源码位置说明IpcFileFormat_dataset.pyxFeather v2 由其子类FeatherFileFormat(L2184)复用CsvFileFormat/CsvFragmentScanOptions_dataset.pyx、L2310CSV 格式与其扫描选项JsonFileFormat_dataset.pyx行式 JSONParquetFileFormat/ParquetReadOptions/ParquetFragmentScanOptions/ParquetFileFragment_dataset_parquet.pyx、L503、L730、L353Parquet 专用,ParquetReadOptions控制元数据读取与统计信息开关OrcFileFormat_dataset_orc.pyx可选编译扩展,不可用时ds.OrcFileFormat直接抛ImportError写选项侧的对应类有IpcFileWriteOptions(L2121)、CsvFileWriteOptions(L2388)、ParquetFileWriteOptions(_dataset_parquet.pyx),它们都是FileWriteOptions(L1260) 的子类,通过FileFormat.make_write_options()创建后传给write_dataset的file_options参数。Parquet 还有加密支持:ParquetEncryptionConfig与ParquetDecryptionConfig从 _dataset_parquet_encryption.pyx 可选导入(见 dataset.py),供ParquetFileWriteOptions/扫描选项使用,前提是 pyarrow 编译时启用了相应支持。八、核心类:Scanner 与扫描选项Scanner是绑定上下文与选项的物化扫描操作,定义在 _dataset.pyx。它把扫描任务、数据 fragment 与数据源粘合在一起,提供两个静态构造入口:Scanner.from_dataset(dataset, *, columnsNone, filterNone, ...):对整数据集扫描;Scanner.from_fragment(fragment, *, schemaNone, columnsNone, filterNone, ...):对单个 fragment 扫描。参数默认值在 _dataset.pyx 的签名与 docstring 中均有明确定义:参数默认值含义columnsNone列投影:列名列表(保持顺序与重复)或{新列名: Expression}字典;支持特殊列__batch_index、__fragment_index、__last_in_fragment、__filename;投影会被下推到数据源,避免加载/反序列化不需要的列filterNone谓词过滤 Expression;尽可能下推以利用分区信息与 Parquet 统计信息,否则在产出的 RecordBatch 上过滤batch_size131072扫描 RecordBatch 的最大行数batch_readahead16单个文件内预读 batch 数;并非所有格式都支持,增大可提高 IO 利用率但增加内存占用fragment_readahead4预读文件(fragment)数,同理以内存换 IOfragment_scan_optionsNone特定 fragment 类型的扫描选项,同一数据集不同扫描可不同use_threadsTrue按 CPU 核数取最大并行度cache_metadataTrue缓存元数据以加速重复扫描memory_poolNone内存池,缺省用默认池实现上,_make_scan_options(见 _dataset.pyx)将上述选项交给ScannerBuilder/ScanOptions,use_threads还会进一步经dataset._scanner_options()做线程数解析。Dataset.scanner()的常用快捷调用即基于同一套选项,配合前文write_dataset(use_threads...)的读侧对应关系一致。Scanner还支持Scanner.from_batches(...)直接从批次流构造(见 _dataset.pyx),write_dataset接收 Scanner 作为数据源时也走这条通道。九、辅助函数:get_partition_keys 与测试验证参考页 Helper functions 一节列出的get_partition_keys与get_partitioning相关能力用于从路径/值反解分区键,供上层框架(如 Spark 连接器)复用。get_partition_keys在 dataset.py 中从pyarrow._dataset导入,并保留了别名_get_partition_keys以维持向后兼容。该模块的行为有大量测试覆盖,可作为本文结论的验证入口:python/pyarrow/tests/test_dataset.py:覆盖dataset()的多种 source 形态、分区发现、Scanner 参数、write_dataset各默认值与existing_data_behavior等;python/pyarrow/tests/parquet/test_dataset.py:Parquet 工厂、ParquetReadOptions、row group 信息与谓词下推等格式特定行为;此外,parquet_dataset与file_visitor的元数据流(构建_metadata文件)也与 pyarrow.parquet 模块的write_metadata配套使用。十、快速参考:从文档到源码的索引文档条目源码入口关键事实datasetdataset.py五类 source 分派;format 支持 parquet/ipc/feather/csv/json/orcparquet_datasetdataset.py基于_metadata文件;底层ParquetDatasetFactorypartitioningdataset.py目录/Hive/文件名三方案;infer 与 factory 语义write_datasetdataset.py默认 max_partitionsmax_open_files1024、max_rows_per_group120Dataset及子类_dataset.pyxInMemory/Union/FileSystem 三实现Scanner_dataset.pyxbatch_size131072、readahead16/4、谓词下推Parquet*系列_dataset_parquet.pyx可选编译扩展,含加密配置OrcFileFormat_dataset_orc.pyx不可用时抛 ImportErrorget_partition_keysdataset.py分区键反解,保留旧别名综上,pyarrow.dataset参考页所列的每个工厂函数、类与辅助函数,在当前仓库中都能在 dataset.py 与两个 Cython 绑定文件之间找到一一对应的实现与参数默认值;实际使用时应特别注意 Dataset API 的不稳定声明、Parquet/ORC 扩展的可选编译前提,以及UnionDataset、preserve_order、max_open_files等对行为有实质影响的约束条件。【免费下载链接】arrowApache Arrow is the universal columnar format and multi-language toolbox for fast data interchange and in-memory analytics项目地址: https://gitcode.com/GitHub_Trending/arrow3/arrow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表