
如何读取和写入 Parquet 分区数据集并在 PyArrow 中按目录组织多文件数据【免费下载链接】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当你需要把一张表落盘为多个 Parquet 文件、按目录分区组织之后再从根目录把整个数据集读回一张表时PyArrow 提供了两条可直接使用的主路径pyarrow.parquet.write_to_dataset/pq.read_table这一对面向 Parquet 的函数以及pyarrow.dataset下称 dataset API里更通用的write_dataset/dataset。下面按“写入分区目录 → 读回 → 验证结果”的顺序给出完整做法以及元数据文件和分区规模方面的注意事项。环境前提按官方安装文档PyArrow 目前兼容 Python 3.11、3.12、3.13 和 3.14建议 64 位系统。准备安装带 Parquet 支持的 PyArrow通过 pip 或 conda 安装的 PyArrow 默认就捆绑了 Parquet 支持pip install pyarrow或者conda install -c conda-forge pyarrow安装后在 Python 中确认pyarrow.parquet可导入即可import pyarrow.parquet as pq如果你是从源码构建 PyArrow则必须在编译 C 库时使用-DARROW_PARQUETON并在构建 Python 包时启用 Parquet 扩展否则没有 Parquet 读写能力。用 write_to_dataset 把表写入分区目录先用一个普通的 Arrow Table 作为数据源以下数据直接来自官方文档示例import pyarrow as pa import pyarrow.parquet as pq table pa.table({one: [-1, None, 2.5], two: [foo, bar, baz], three: [True, False, True]})写入分区数据集pq.write_to_dataset(table, root_pathdataset_name, partition_cols[one, two])这条命令的行为在文档中有明确定义root_path指定数据保存的父目录partition_cols是要按列名划分的列列按给出的顺序分区每个分区内的文件划分由分区列的取值唯一性决定即按唯一值切分不传filesystem参数时默认使用本地文件系统。写入后目录结构类似文档示例按分区列不同会有多个层级dataset_name/ year2007/ month01/ 0.parq 1.parq ... month02/ 0.parq 1.parq ... year2008/ month01/ ...两点使用限制如果生成的表后续要交给 HIVE 使用分区列的取值必须与你所用 HIVE 版本允许的字符集兼容需要写到远端文件系统时只需追加filesystem参数单个表的写入在函数内部已经用with语句包装无需在外层再包# 可选分支Hadoop 远端文件系统示例 # host、port、user、ticket_cache_path 需替换为你实际的集群连接参数 from pyarrow.fs import HadoopFileSystem fs HadoopFileSystem(host, port, useruser, kerb_ticketticket_cache_path) pq.write_to_dataset(table, root_pathdataset_name, partition_cols[one, two], filesystemfs)可选生成 _metadata 与 _common_metadata 索引文件Spark 或 Dask 等框架可选地会读取分区数据集根目录下的_metadata和_common_metadata文件前者包含全部文件的行组元数据后者包含整个数据集的 schema。它们本身是只含元数据的 Parquet 文件。使用这两个文件可以让后续创建 Parquet Dataset 时直接复用已存的 schema 和文件路径而不必推断 schema 并遍历目录查找全部 Parquet 文件——对文件访问开销较大的文件系统收益明显。注意这不是 Parquet 标准而是这些框架在实践中形成的约定pq.write_to_dataset默认不会写这两个文件。如果需要用metadata_collector收集元数据后手动合并写入# 写入数据集并收集所有已写文件的元数据 metadata_collector [] root_path dataset_name_1 pq.write_to_dataset(table, root_path, metadata_collectormetadata_collector) # 写不含行组统计的 _common_metadata 文件 pq.write_metadata(table.schema, root_path /_common_metadata) # 写包含所有文件行组统计的 _metadata 文件 pq.write_metadata( table.schema, root_path /_metadata, metadata_collectormetadata_collector )如果你不用write_to_dataset而是用pq.write_table或ParquetWriter逐个文件写入metadata_collector关键字同样可用但此时需要你自己调用set_file_path设置行组元数据里的相对文件路径且所有文件的 schema 必须一致。写入完成后可以用下面命令读取索引文件做检查文档示例输出如下其中省略号部分按实际构建版本变化pq.read_metadata(_metadata) pyarrow._parquet.FileMetaData object at ... created_by: parquet-cpp-arrow version ... num_columns: 3 num_rows: 3 num_row_groups: 1 format_version: 2.6 serialized_size: ...用 Dataset API 写入目录并控制分区 schemapyarrow.dataset模块提供统一的多文件数据集接口支持 Parquet、Feather / Arrow IPC、CSV 和 ORCORC 目前只能读不能写。最基础的写入指定目录而不是单个文件名formatparquetimport pyarrow.dataset as ds import numpy as np table pa.table({a: range(10), b: np.random.randn(10), c: [1, 2] * 5}) ds.write_dataset(table, sample_dataset, formatparquet)这会在sample_dataset目录下生成单个文件part-0.parquet。文档明确提醒重复执行同一命令会直接覆盖已有的 part-0.parquet如果要在已有数据集上追加文件每次调用ds.write_dataset都需要指定新的basename_template以避免覆盖。要按目录组织分区数据用一个 partitioning 对象声明分区 schema 和风格示例来自官方文档part ds.partitioning( pa.schema([(c, pa.int16())]), flavorhive ) ds.write_dataset(table, partitioned_dataset, formatparquet, partitioningpart)执行后数据的一半在dataset_root/c1目录另一半在dataset_root/c2目录。分区列本身不再写入各个 Parquet 文件只体现在目录名里。写大文件时两个常用的控制参数max_rows_per_file限制单个文件的行数是控制文件大小的主要手段。不设置时每个输出目录默认只有一个文件数据量大时文件可能过大导致下游读取时内存不足max_open_files限制写入过程中保持打开的文件数默认 900。它只作用于分区写入行按分区值分发到对应文件设得太低会把数据切碎成许多小文件。如果你的进程同时还在扫描数据集打开文件总数可能超过系统上限——文档给出的例子是扫描 300 个文件同时写 900 个文件合计 1200 个在 Linux 上可能报 “Too Many Open Files” 错误默认 Linux 限制 1024。降低max_open_files或提高系统文件句柄上限二选一。读回分区数据集用 pyarrow.parquet 读取pq.ParquetDataset接受目录名或文件路径列表能发现并推断 Hive 风格等常见分区结构dataset pq.ParquetDataset(dataset_name/) table dataset.read()文档示例输出列one、two是分区列读回后表现为字典类型pyarrow.Table three: bool one: dictionaryvaluesstring, indicesint32, ordered0 two: dictionaryvaluesstring, indicesint32, ordered0 ---- three: [[true],[true],[false]] one: [ -- dictionary: [-1,2.5] -- indices: [0], ... two: [ -- dictionary: [foo,baz,bar] -- indices: [0], ...也可以跳过 Dataset 对象直接用便捷函数table pq.read_table(dataset_name)用 pyarrow.dataset 读取用 dataset API 创建数据集时显式声明partitioninghive它会递归遍历目录发现文件但不会开始读数据np.random.seed(0) table pa.table({a: range(10), b: np.random.randn(10), c: [1, 2] * 5, part: [a] * 5 [b] * 5}) pq.write_to_dataset(table, parquet_dataset_partitioned, partition_cols[part]) dataset ds.dataset(parquet_dataset_partitioned, formatparquet, partitioninghive) dataset.files文档示例中dataset.files返回[parquet_dataset_partitioned/parta/...-0.parquet, parquet_dataset_partitioned/partb/...-0.parquet]分区字段虽然不在 Parquet 文件里扫描时会被加回结果表dataset.to_table().to_pandas().head(3)文档示例输出a b c part 0 0 0.144044 1 a 1 1 1.454274 2 a 2 2 0.761038 1 a按分区键过滤时不匹配的文件根本不会被加载dataset.to_table(filterds.field(part) b).to_pandas()除 Hive 风格外dataset API 还支持两种分区风格按需替换上面的partitioning参数显式声明分区键 schemads.partitioning(pa.schema([(year, pa.int16()), (month, pa.int8()), (day, pa.int32())]), flavorhive)目录分区路径段直接是取值如/2019/11/15字段名不在路径中必须用ds.partitioning(field_names[year, month, day])指定。验证写入与读取结果文档给出的可核对点有三个检查目录里的文件dataset.files应列出全部数据文件且路径体现分区目录结构检查读回的列to_pandas()或table打印中应包含数据列和被加回的分区列pq.ParquetDataset路径下分区列是 dictionary 类型dataset API 路径下按声明的 schema 或从路径推断的类型出现检查元数据索引如果写了_metadatapq.read_metadata应返回包含正确num_columns、num_rows、num_row_groups的FileMetaData对象。注意事项与限制分区列的类型与顺序原始表中的分区列在保存/读取过程中会被转成 Arrow dictionary 类型对应 pandas 的 categorical且分区列的顺序不会在保存/读取过程中保留。从远端文件系统读入 pandas 时如需保持行顺序可能需要sort_index前提是写入时启用了preserve_index选项。读列子集时必须显式带上分区键当你用columns关键字只读部分列时分区键必须显式包含在columns里否则结果中不会出现这些列。单文件路径不推断分区列把单个文件路径传给pq.read_table或pq.ParquetDataset时即使路径里有 Hive 风格的段如dataset_name/year2017/data1.parquet也不会推断出分区列。要拿到分区列请传入父目录# 不含 year 列 pq.read_table(dataset_name/year2017/data1.parquet) # 含 year 分区列 pq.read_table(dataset_name/)分区规模要克制分区一方面增加文件数、一方面加深目录树。按日期分区一年数据至少有 365 个文件再叠加一个 1000 个唯一值的维度最多 365,000 个文件——过细的分区会产出以元数据为主的小文件目录发现时也需要相应次数的 list 调用365 次变 365,365 次。文档给出的经验边界是避免小于 20MB 或大于 2GB 的文件避免超过 10,000 个不同分区的布局。没有事务和 ACID 保证dataset API 对读写都不提供事务支持。并发读取没问题但并发写入或与读并发可能产生意外行为写进行程中被意外杀掉会留下不一致状态。文档建议的做法包括为每个写入方使用唯一的 basename 模板、用临时目录存放新文件、或者把文件清单单独存储而不依赖目录发现。相关文档分区数据集多文件说明Tabular Datasetsdataset API说明单文件 Parquet 读写安装 PyArrow【免费下载链接】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),仅供参考