news 2026/9/14 17:28:39

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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模块提供的工厂函数(datasetparquet_datasetpartitioningwrite_dataset等)、核心类(DatasetFileSystemDatasetScanner、各文件格式与分区方案等)以及辅助函数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 行)DatasetFragmentScanner、分区方案、CSV/Ipc/JSON 格式等类的 C++ 绑定
格式扩展绑定_dataset_parquet.pyx、_dataset_orc.pyxParquet 与 ORC 的读写选项与工厂

dataset.py 展示了可选依赖的加载策略:核心类从pyarrow._dataset导入,而OrcFileFormat与 Parquet 相关类分别从pyarrow._dataset_orcpyarrow._dataset_parquet可选导入;若对应扩展未编译,访问这些类时会通过__getattr__抛出明确的ImportError(如 "The pyarrow installation is not built with support for the ORC file format")。这一点决定了后文各格式类的使用前提。

文档参考页中列出的fieldscalar属于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 的一个字段,且所有字段都必须出现。例如 schema<year:int16, month:int8>,路径/2009/11解析为year == 2009 and month == 11;
  • HivePartitioning(flavor="hive"):Hive 风格的/key=value/嵌套目录。字段顺序无关,未识别的键被忽略。例如 schema<year:int16, month:int8, day:int8>,路径/year=2009/month=11/day=15是合法的,即使与 schema 字段顺序不一致;
  • FilenamePartitioning(flavor="filename"):分区值直接编码在文件名中,字段值之间用_分隔。例如2009_11_part-0.parquet解析为year == 2009 and month == 11

2.2 参数与返回类型

partitioning(schema=None, field_names=None, flavor=None, dictionaries=None)的语义在 dataset.py 的 docstring 中有完整说明:

  • schema:描述路径中分区的 Schema;若未提供而给了field_names/flavor,则类型从文件路径推断,此时返回的是PartitioningFactory而非确定的Partitioning;
  • field_names:字段名列表,仅对 Directory 方案有效,类型同样从路径推断;
  • flavor:缺省为目录分区,"hive"表示 Hive 分区,"filename"表示文件名分区;
  • dictionaries:若分区字段是字典类型,必须提供包含该列全部可能取值的数组,否则解析会报错;也可以传字符串"infer"让 Arrow 自行发现字典值(此时返回PartitioningFactory)。

注意 dataset.py 中的互斥约束:Directory 方案下不能同时给schemafield_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())) ]), dictionaries="infer") # Hive 方案:路径形如 "/year=2009/month=11" part = ds.partitioning( pa.schema([("year", pa.int16()), ("month", pa.int8())]), flavor="hive") # Hive 方案由目录结构自动发现(类型一并推断) part = ds.partitioning(flavor="hive")

对应的底层类定义在 _dataset.pyx 中:Partitioning(L2548)→KeyValuePartitioning(L2692)→ 三个具体子类DirectoryPartitioning(L2742)、HivePartitioning(L2870)、FilenamePartitioning(L3021);PartitioningFactory(L2642)则封装了"先扫描目录再确定类型/字典"的延迟发现流程。

三、工厂函数:dataset() 打开数据源

dataset()是整个模块的入口,位于 dataset.py。参考页中它与parquet_datasetwrite_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", format="parquet") dataset.to_table() # 打开单个文件并显式指定 schema(等价于投影列子集) myschema = pa.schema([('n_legs', pa.int64()), ('animal', pa.string())]) dataset = ds.dataset("file.parquet", schema=myschema, format="parquet") dataset.to_table() # 打开分区目录 ds.write_dataset(table, "partitioned_dataset", format="parquet", partitioning=['year']) dataset = ds.dataset("partitioned_dataset", format="parquet") # S3 桶中的目录(需要相应的凭证配置) ds.dataset("s3://mybucket/nyc-taxi/", format="parquet") # 从相对路径文件列表打开 dataset = ds.dataset([ "partitioned_dataset/2019/part-0.parquet", "partitioned_dataset/2020/part-0.parquet", "partitioned_dataset/2021/part-0.parquet", ], format='parquet') # 文件列表 + filesystem URI 前缀 paths = ['part0/data.parquet', 'part1/data.parquet', 'part3/data.parquet'] ds.dataset(paths, filesystem='s3://bucket/nested/directory', format='parquet') # 嵌套 UnionDataset:组合任意其他数据集 ds.dataset([ ds.dataset("s3://old-taxi-data", format="parquet"), ds.dataset("local/path/to/data", format="ipc") ])

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 对象列表";
  • 最终把partitioningpartition_base_direxclude_invalid_filesselector_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 一节中ParquetFileFormatParquetReadOptions等类被组合使用的地方。

五、工厂函数: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选择方案类型(缺省为目录分区);
  • schemafilesystemfile_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_partitionsmax_open_files为 1024、max_rows_per_file为 0、max_rows_per_group1 << 20min_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", format="parquet", partitioning=["year"], file_visitor=file_visitor, existing_data_behavior="overwrite_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_groupsfragment_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 级扫描选项,派生出CsvFragmentScanOptionsParquetFragmentScanOptions等)。

七、核心类:FileFormat 与各格式扩展

FileFormat基类定义在 _dataset.pyx,提供make_write_options()等通用接口。各格式绑定类与参考页的对应关系:

参考页类名源码位置说明
IpcFileFormat_dataset.pyxFeather v2 由其子类FeatherFileFormat(L2184)复用
CsvFileFormat/CsvFragmentScanOptions_dataset.pyx、L2310CSV 格式与其扫描选项
JsonFileFormat_dataset.pyx行式 JSON
ParquetFileFormat/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_datasetfile_options参数。

Parquet 还有加密支持:ParquetEncryptionConfigParquetDecryptionConfig从 _dataset_parquet_encryption.pyx 可选导入(见 dataset.py),供ParquetFileWriteOptions/扫描选项使用,前提是 pyarrow 编译时启用了相应支持。

八、核心类:Scanner 与扫描选项

Scanner是"绑定上下文与选项的物化扫描操作",定义在 _dataset.pyx。它把扫描任务、数据 fragment 与数据源粘合在一起,提供两个静态构造入口:

  • Scanner.from_dataset(dataset, *, columns=None, filter=None, ...):对整数据集扫描;
  • Scanner.from_fragment(fragment, *, schema=None, columns=None, filter=None, ...):对单个 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)数,同理以内存换 IO
fragment_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_keysget_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_datasetfile_visitor的元数据流(构建_metadata文件)也与 pyarrow.parquet 模块的write_metadata配套使用。

十、快速参考:从文档到源码的索引

文档条目源码入口关键事实
datasetdataset.py五类 source 分派;format 支持 parquet/ipc/feather/csv/json/orc
parquet_datasetdataset.py基于_metadata文件;底层ParquetDatasetFactory
partitioningdataset.py目录/Hive/文件名三方案;"infer" 与 factory 语义
write_datasetdataset.py默认 max_partitions=max_open_files=1024、max_rows_per_group=1<<20
Dataset及子类_dataset.pyxInMemory/Union/FileSystem 三实现
Scanner_dataset.pyxbatch_size=131072、readahead=16/4、谓词下推
Parquet*系列_dataset_parquet.pyx可选编译扩展,含加密配置
OrcFileFormat_dataset_orc.pyx不可用时抛 ImportError
get_partition_keysdataset.py分区键反解,保留旧别名

综上,pyarrow.dataset参考页所列的每个工厂函数、类与辅助函数,在当前仓库中都能在 dataset.py 与两个 Cython 绑定文件之间找到一一对应的实现与参数默认值;实际使用时应特别注意 Dataset API 的不稳定声明、Parquet/ORC 扩展的可选编译前提,以及UnionDatasetpreserve_ordermax_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/14 17:28:33

Less.js 备忘清单:CSS 预处理器核心语法与内置函数实战速查

Less.js 备忘清单&#xff1a;CSS 预处理器核心语法与内置函数实战速查 【免费下载链接】reference 面向开发者的技术速查清单&#xff08;Cheat Sheets&#xff09;集合&#xff0c;整理常见技术、工具与开发流程&#xff0c;帮助快速查阅关键信息&#xff0c;提高开发效率。 …

作者头像 李华
网站建设 2026/9/14 17:24:03

经销商数字化转型:Dify低代码平台实战解析

1. 经销商数字化转型的必然选择最近两年走访了上百家区域经销商&#xff0c;发现一个共性痛点&#xff1a;传统经营模式越来越难应对市场变化。上个月在山东聊城遇到一位做快消品批发的张总&#xff0c;他给我算了一笔账&#xff1a;人工统计订单的差错率高达8%&#xff0c;库存…

作者头像 李华
网站建设 2026/9/14 17:22:22

COMSOL频域感应加热模型构建与优化指南

1. COMSOL频域感应加热模型构建指南 感应加热技术在现代工业中应用广泛&#xff0c;从金属热处理到半导体加工都离不开这项技术。作为一名长期使用COMSOL进行电磁热耦合仿真的工程师&#xff0c;我将分享如何建立一个完整的底部电磁波频域感应加热模型&#xff0c;用于分析被加…

作者头像 李华
网站建设 2026/9/14 17:20:48

一文搞懂CORS与WebSocket跨域:原理、排查与配置实战

刚处理完一个线上事故&#xff0c;前端同事盯着浏览器控制台里那行经典的红色报错——has been blocked by CORS policy: No Access-Control-Allow-Origin header is present——一脸无辜地看着我&#xff1a;“后端不是都配了跨域吗&#xff1f;怎么 WebSocket 还是连不上&…

作者头像 李华