尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Apache Arrow PyArrow 表格文件格式 API 全解析:CSV、Feather、JSON、Parquet 与 ORC

Apache Arrow PyArrow 表格文件格式 API 全解析:CSV、Feather、JSON、Parquet 与 ORC Apache Arrow PyArrow 表格文件格式 API 全解析CSV、Feather、JSON、Parquet 与 ORC【免费下载链接】arrowApache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing项目地址: https://gitcode.com/gh_mirrors/arrow13/arrow本文围绕 Apache Arrow 项目 Python 绑定 PyArrow 的表格文件格式 API对应文档 docs/source/python/api/formats.rst展开系统梳理pyarrow.csv、pyarrow.feather、pyarrow.json、pyarrow.parquet、pyarrow.orc五个模块的读写函数、选项类与元数据类。读完本文你将掌握如何用统一、高性能的方式在 Arrow 列式内存格式与五种主流表格文件格式之间往返转换理解每个选项参数的默认值与底层影响并能在实际数据管道中直接复用文中的示例代码。模块总览一张表看尽五种格式该 API 文档按格式划分为五节每一节对应一个独立模块且模块之间通过共同的pyarrow.Table/RecordBatch数据类型打通——无论从哪种格式读入得到的都是统一的 Arrow 内存数据可无缝衔接计算、Dataset 扫描与写出模块核心函数关键选项/辅助类元数据/进阶类pyarrow.csvread_csv/open_csv/write_csvReadOptions/ParseOptions/ConvertOptions/WriteOptionsCSVStreamingReader/CSVWriter/ISO8601/InvalidRowpyarrow.featherread_feather/read_table/write_feathercompression/compression_level/chunksize/versionFeatherDataset源码提供未列入 API 列表pyarrow.jsonread_jsonReadOptions/ParseOptions—pyarrow.parquetread_table/read_pandas/read_schema/read_metadata/write_table/write_to_dataset/write_metadataParquetDataset/ParquetFile/ParquetWriterFileMetaData/RowGroupMetaData/ColumnChunkMetaData/Statistics/SortingColumn/ParquetSchema/ColumnSchema/ParquetLogicalTypepyarrow.parquet.encryption下的CryptoFactory/KmsClient/KmsConnectionConfig/EncryptionConfiguration/DecryptionConfigurationpyarrow.orcread_table/write_tableORCFile/ORCWriter文件级属性metadata/schema/nrows/nstripes/compression等在 Python 侧pyarrow.csv见 python/pyarrow/csv.py与pyarrow.json见 python/pyarrow/json.py均为薄封装实际实现位于 Cython 层 _csv.pyx、_json.pyx最终落到 C 实现如 cpp/src/arrow/csv/reader.ccParquet 实现则位于 python/pyarrow/parquet/core.py 与 _parquet.pyx。CSV 文件pyarrow.csvCSV 是文本格式中读写控制最精细的一类。pyarrow.csv采用「读取 → 解析 → 类型转换」三段式流水线分别由ReadOptions、ParseOptions、ConvertOptions控制三者默认值均由 C 侧Defaults()构造见 _csv.pyx。读取与流式读取read_csv一次性读入整个文件并返回pyarrow.Table。它支持路径或类文件对象且会自动解压.gz、.bz2等压缩扩展名见 read_csv 文档import io from pyarrow import csv s ( animals,n_legs,entry\n Flamingo,2,2022-03-01\n Horse,4,2022-03-02\n Brittle stars,5,2022-03-03\n Centipede,100,2022-03-04 ) source io.BytesIO(s.encode()) table csv.read_csv(source) # animals: string, n_legs: int64, entry: date32[day]注意read_csv的类型推断能力示例中的2/4/5/100被推断为int64ISO 格式日期被推断为date32。open_csv则返回CSVStreamingReader继承自RecordBatchReader适合超大文件按批处理。与read_csv不同流式读取始终单线程见 open_csv 文档。ReadOptions控制数据读取ReadOptions决定「读多少、从哪读」见 ReadOptions 定义use_threads默认True是否多线程加速读取影响块间并行度block_size每次从输入流处理的字节数决定多线程粒度以及单个 record batch / table chunk 的大小最小合法值为 1skip_rows默认0在列名如有与 CSV 数据之前跳过的行数skip_rows_after_names默认0读取列名后跳过的行数可大于单个 block 的行数空行也计入执行顺序为先skip_rows再读列名除非指定column_names最后应用skip_rows_after_namescolumn_names目标表的列名为空时回退到autogenerate_column_namesautogenerate_column_names默认False为True时自动生成f0、f1… 形式的列名否则从skip_rows之后的首行读取列名encoding默认utf8CSV 数据的字符编码无法按该编码解码的列仍可读为 Binary。例如忽略首行编号行并替换列名read_options csv.ReadOptions( column_names[animals, n_legs, entry], skip_rows1, ) csv.read_csv(io.BytesIO(s.encode()), read_optionsread_options)若改用autogenerate_column_namesTrue, skip_rows1则列名变为f0/f1/f2。ParseOptions控制语法解析ParseOptions决定「怎么切分字段」见 ParseOptions 定义delimiter默认,单字符分隔符可改为;、\t等quote_char默认引号字符传False表示禁止引号double_quote默认True引号值内的两个引号是否表示数据中的一个引号escape_char默认False转义字符传False表示不允许转义newlines_in_values默认False值内是否允许换行开启会显著降低多线程读取性能ignore_empty_lines默认True是否忽略空行为False时空行被视为含单个空值仅当 CSV 为单列时invalid_row_handler默认None为每一行「列数不匹配」的失败解析行调用接收一个InvalidRow参数返回skip或error。invalid_row_handler是处理脏数据注释行、缺列行的常用手段。文档示例用分号分隔并跳过#开头的注释行def skip_comment(row): if row.text.startswith(# ): return skip return error parse_options csv.ParseOptions( delimiter;, invalid_row_handlerskip_comment, ) csv.read_csv(source, parse_optionsparse_options)对应测试 test_csv.py#test_invalid_row_handler 验证了skip/error两种返回值的完整行为handler 返回skip时坏行被丢弃返回error时抛出ValueError且InvalidRow记录期望列数expected_columns、实际列数actual_columns、物理行号多线程并行读取时可能为None与行文本见 InvalidRow 定义。Cython 侧_handle_invalid_row仅接受error与skip其他返回值会抛出ValueError。ConvertOptions控制类型转换ConvertOptions决定「文本如何变成 Arrow 类型」见 ConvertOptions 定义check_utf8默认True是否校验字符串列的 UTF-8 合法性column_typespyarrow.Schema或dict显式指定列类型指定后该列跳过类型推断null_values表示空值的字符串序列默认值适用于大多数场景注意默认情况下字符串列不做空值检查需开启strings_can_be_nullTruetrue_values/false_values表示布尔真/假的字符串序列decimal_point默认.浮点与 decimal 数据的小数点字符strings_can_be_null默认False字符串/二进制列是否可为空为True时null_values中的字符串在字符串列中视为空quoted_strings_can_be_null默认True被引号包裹的值是否可为空为False时带引号的值永不被视为空include_columns要包含进 Table 的列名列表为空则包含全部列非空则按给定顺序只包含这些列include_missing_columns默认False为False时include_columns中未出现的列会报错为True时生成一个空值列类型由column_types决定否则为 null仅在include_columns非空时生效auto_dict_encode默认False自动对字符串/二进制列做字典编码每个 chunk 超过auto_dict_max_cardinality个不同值后回退为普通编码对column_types中显式指定的列不生效auto_dict_max_cardinalityauto_dict_encode的最大字典基数按 chunk 计timestamp_parsersstrptime()兼容格式字符串序列按顺序尝试推断/转换时间戳也可传入ISO8601()特殊值默认使用内置快速 ISO-8601 解析器。几个典型用法import pyarrow as pa # 1) 显式指定列类型禁用该列推断 convert_options csv.ConvertOptions( column_types{n_legs: pa.float64()}) # 2) 非 ISO 日期格式的 timestamp 解析 convert_options csv.ConvertOptions( timestamp_parsers[%m/%d/%Y, %m-%d-%Y]) # 3) 只读子集列 convert_options csv.ConvertOptions( include_columns[animals, n_legs]) # 4) 子集列包含缺失列时补空值列 convert_options csv.ConvertOptions( include_columns[animals, n_legs, location], include_missing_columnsTrue) # 5) 开启字典编码并限制基数上限超过则回退普通编码 convert_options csv.ConvertOptions( auto_dict_encodeTrue, auto_dict_max_cardinality2) # 6) 空字符串视为空值 convert_options csv.ConvertOptions( strings_can_be_nullTrue)写出 CSVwrite_csv与WriteOptionswrite_csv将 Table 写出为 CSVCSVWriter则提供增量写入见 _csv.pyx。WriteOptions参数见 WriteOptions 定义include_header默认True是否写入列名表头batch_size默认1024转换与写出 CSV 时一次处理的批大小delimiter默认,分隔符quoting_style默认needed引号策略可选值needed仅在必要时给值加引号all_valid所有合法值都加引号空值不加none任何值都不加引号值内含引号、分隔符或换行等特殊字符时会抛错。import pyarrow as pa from pyarrow import csv table pa.table({n_legs: [2, 4, 5, 100], animal: [Flamingo, Horse, Brittle stars, Centipede]}) csv.write_csv(table, animals.csv) # 默认含表头Feather 文件pyarrow.featherFeather 是 Arrow 官方推荐的轻量列式交换格式基于 Arrow IPC读写极快、无需 schema 推断适合进程间与持久化列式数据。模块实现见 python/pyarrow/feather.py。写入write_featherwrite_feather签名及默认值见 write_feather 定义dfpandas.DataFrame或pyarrow.Tabledest本地目标路径compression默认None取值{zstd, lz4, uncompressed}None时 V2 文件在可用lz4_frame时默认用 LZ4否则不压缩compression_level默认None所选压缩器的特定压缩级别chunksize默认NoneV2 文件内部 Arrow IPC 写入时 RecordBatch chunk 的最大大小默认约 64K 行version默认2Feather 文件版本V2 为当前版本V1 为受限的旧格式。实现细节V2 下自动选择lz4需要Codec.is_available(lz4_frame)为真见 feather.py#L178V1 不支持压缩与chunksize且 DataFrame 索引处理规则不同V1 不保留非 RangeIndex 索引V2 通过preserve_indexNone自动处理。若写入失败模块会清理残留的目标文件。import pandas as pd import pyarrow.feather as feather df pd.DataFrame({n_legs: [2, 4, 5, 100], animal: [Flamingo, Horse, Brittle stars, Centipede]}) feather.write_feather(df, animals.feather) feather.write_feather(df, animals_zstd.feather, compressionzstd)读取read_feather与read_tableread_feather(source, columnsNone, use_threadsTrue, memory_mapFalse, **kwargs)读取为pandas.DataFrame**kwargs透传给pyarrow.Table.to_pandas见 read_feather 定义read_table(source, columnsNone, memory_mapFalse, use_threadsTrue)读取为pyarrow.Table见 read_table 定义。两者共同支持columns只读指定列可传列索引整数或列名字符串V2 读取时会按精确顺序/选择重新投影列memory_map在 source 为路径时启用内存映射读取use_threads控制多线程读取与转 pandas 的并行度。源码中read_table会依据列参数类型全为int→ 按索引读全为str→ 按名字读走不同读取路径。JSON 文件pyarrow.jsonpyarrow.json目前支持**行分隔 JSONNDJSON/JSON Lines**格式的读取API 面简洁ReadOptions、ParseOptions、read_json见 python/pyarrow/_json.pyx。ReadOptionsuse_threads默认True多线程加速读取block_size每次处理的字节数决定多线程粒度与 Table 中单个 chunk 的大小。ParseOptionsexplicit_schema默认None显式提供 schema此时不做类型推断并忽略其他推断字段newlines_in_values默认False对象是否可跨多行如 pretty-print 格式为False时输入必须以空行结束unexpected_field_behavior默认inferexplicit_schema之外字段的处理策略三者互斥ignore忽略意外字段error遇到意外字段直接报错infer对意外字段做类型推断并纳入输出。unexpected_field_behavior的 Cython 层实现了Ignore / Error / InferType三种枚举映射非法值会抛出ValueError见 _json.pyx#L205。read_jsonfrom pyarrow import json table json.read_json(data.jsonl)配合显式 schema 控制意外字段行为的用法与测试 test_json.py#test_explicit_schema_with_unexpected_behaviour 完全对应import pyarrow as pa from pyarrow import json schema pa.schema([(foo, pa.binary())]) # 默认 inferfoo 按 schema 读num 被推断为 int64 table json.read_json(data.jsonl, parse_optionsjson.ParseOptions(explicit_schemaschema)) # ignore只保留 foo 列 table json.read_json(data.jsonl, parse_optionsjson.ParseOptions( explicit_schemaschema, unexpected_field_behaviorignore))Parquet 文件pyarrow.parquetParquet 是文档中 API 面最庞大的格式包含高层函数、三个核心类、元数据类族与模块级加密子模块。实现位于 python/pyarrow/parquet/core.py高层封装与 python/pyarrow/_parquet.pyxCython 绑定。高层读写函数read_table(source, *, columnsNone, use_threadsTrue, schemaNone, use_pandas_metadataFalse, read_dictionaryNone, memory_mapFalse, buffer_size0, partitioninghive, filesystemNone, filtersNone, pre_bufferTrue, coerce_int96_timestamp_unitNone, decryption_propertiesNone, ...)见 read_table 定义columns只读的列名子集use_pandas_metadata为True且文件含 pandas schema 元数据时确保索引列也被加载read_dictionary直接以DictionaryArray读取的列名列表memory_map源为路径时启用内存映射buffer_size为正时对单个列块反序列化做读缓冲partitioning默认hive用于多文件分区目录发现filters过滤谓词pyarrow.compute.Expression或嵌套元组列表分区键可用于跳过整个文件pre_buffer默认True在高延迟文件系统如 S3上合并并发发起文件读取使用后台 I/O 线程池coerce_int96_timestamp_unit将 INT96 存储的时间戳强制转换为指定分辨率如msNone等价于ns。从源码看read_table优先走ParquetDataset基于pyarrow.dataset仅当 dataset 模块不可用时回退到单文件ParquetFile且此时filters、partitioning、schema参数不可用见 core.py#L1810-L1841。import pyarrow as pa import pyarrow.parquet as pq table pa.table({n_legs: [2, 2, 4, 4, 5, 100], animal: [Flamingo, Parrot, Dog, Horse, Brittle stars, Centipede]}) pq.write_table(table, example.parquet) t pq.read_table(example.parquet) # 读回 pyarrow.Table df pq.read_pandas(example.parquet) # 连索引一并读出其余高层函数read_pandas(source, columnsNone, **kwargs)等价于read_table(..., use_pandas_metadataTrue)将 DataFrame 索引值作为列读回见 read_pandasread_schema(where, memory_mapFalse, decryption_propertiesNone, filesystemNone)读取文件的有效 Arrow schema见 core.py#L2305read_metadata(where, memory_mapFalse, decryption_propertiesNone, filesystemNone)读取单个 Parquet 文件 footer 中的FileMetaData见 read_metadatawrite_table(table, where, row_group_sizeNone, version2.6, use_dictionaryTrue, compressionsnappy, write_statisticsTrue, use_deprecated_int96_timestampsNone, coerce_timestampsNone, allow_truncated_timestampsFalse, data_page_sizeNone, flavorNone, compression_levelNone, use_byte_stream_splitFalse, column_encodingNone, data_page_version1.0, use_compliant_nested_typeTrue, encryption_propertiesNone, write_batch_sizeNone, dictionary_pagesize_limitNone, store_schemaTrue, write_page_indexFalse, write_page_checksumFalse, sorting_columnsNone, store_decimal_as_integerFalse)见 write_table核心参数包括row_group_size每个行组的最大行数、version默认2.6格式版本、use_dictionary字典编码开关、compression默认snappy、write_statistics写入统计信息、data_page_version默认1.0与sorting_columns声明排序列等write_to_dataset(table, root_path, partition_colsNone, ...)按分区列将数据分布写入目录数据集见 core.py#L1993write_metadata(schema, where, metadata_collectorNone, filesystemNone)仅写入元数据的 Parquet 文件与write_to_dataset配合生成_common_metadata与_metadata侧车文件见 write_metadata 示例metadata_collector [] pq.write_to_dataset(table, dataset_metadata, metadata_collectormetadata_collector) pq.write_metadata(table.schema, dataset_metadata/_common_metadata) pq.write_metadata(table.schema, dataset_metadata/_metadata, metadata_collectormetadata_collector)三个核心类ParquetFile(source, metadataNone, common_metadataNone, read_dictionaryNone, memory_mapFalse, buffer_size0, pre_bufferFalse, coerce_int96_timestamp_unitNone, decryption_propertiesNone, thrift_string_size_limitNone, thrift_container_size_limitNone, filesystemNone, page_checksum_verificationFalse)单个 Parquet 文件的读取接口支持复用已有FileMetaData、read()/read_row_group()等见 ParquetFileParquetWriter(where, schema, ...)增量构建 Parquet 文件的写入类参数与write_table对齐支持metadata_collector选项收集元数据、flavor兼容模式如spark时自动使用 INT96 时间戳且实现了上下文管理器协议见 ParquetWriterwith pq.ParquetWriter(example.parquet, table.schema, compressionzstd) as writer: writer.write_table(table, row_group_size2048)ParquetDataset(path_or_paths, filesystemNone, schemaNone, *, filtersNone, read_dictionaryNone, memory_mapFalse, buffer_sizeNone, partitioninghive, ignore_prefixesNone, pre_bufferTrue, coerce_int96_timestamp_unitNone, decryption_propertiesNone, thrift_string_size_limitNone, thrift_container_size_limitNone, page_checksum_verificationFalse)封装一个由多个文件与子目录分区组成的完整 Parquet 数据集见 ParquetDataset。底层委托给pyarrow.datasetfilters在数据集发现阶段即可利用嵌套目录中的分区键跳过不含匹配行的文件ignore_prefixes默认[., _]发现阶段跳过以这些前缀开头的 basenamepre_buffer默认True面向高延迟文件系统。Parquet 元数据类族读取元数据是 Parquet 生态如查询优化、行组裁剪的重要能力相关类定义于 _parquet.pyxFileMetaDataL853单文件元数据属性含created_by、num_columns、num_rows、num_row_groups、format_version、serialized_size支持row_group(i)、to_dict()、append_row_groups()与write_metadata_file()RowGroupMetaDataL726行组元数据可定位column(chunk_index)与行组统计ColumnChunkMetaDataL315列块元数据含文件偏移、压缩大小、编码、路径等StatisticsL52单行组单列的统计信息属性包括has_min_max、min/max按逻辑类型返回 Python 等价对象如datetime.date、decimal.Decimal、min_raw/max_raw物理类型、null_count、distinct_count、num_values、physical_type、logical_type、converted_typelegacy并提供to_dict()与equals()SortingColumnL512描述文件声明排序的列ParquetSchemaL1077与ColumnSchemaL1164分别表示 Parquet 文件 schema 与单列 schema物理类型、逻辑类型、路径、重复类型等ParquetLogicalTypeL218逻辑类型描述。meta pq.read_metadata(example.parquet) print(meta.num_rows, meta.num_row_groups, meta.format_version) for i in range(meta.num_row_groups): rg meta.row_group(i) col rg.column(0) print(col.statistics) # has_min_max / min / max / null_count ...加密 Parquetpyarrow.parquet.encryptionAPI 文档单列了加密子模块支持 Parquet Modular Encryption相关类从 python/pyarrow/parquet/encryption.py 导出实现位于_parquet_encryptionCython 模块CryptoFactory创建加密/解密属性的工厂提供file_encryption_properties()与file_decryption_properties()KmsClientKMS密钥管理服务客户端抽象供CryptoFactory对接外部 KMSKmsConnectionConfigKMS 连接配置EncryptionConfiguration文件加密配置加密列、footer 密钥等DecryptionConfiguration文件解密配置。用法上加密读写在read_table/write_table/ParquetFile的decryption_properties参数中注入由CryptoFactory创建的解密属性。ORC 文件pyarrow.orcORC 模块提供单文件读写接口实现见 python/pyarrow/orc.py。其 API 文档仅列出四个符号ORCFile、ORCWriter、read_table、write_table。ORCFile单文件读取ORCFile(source)提供文件级属性与读取方法见 ORCFile 定义元数据属性metadataKeyValueMetadata、schemaArrow schema、nrows、nstripesstripe 数、file_version、software_version、compression、compression_size、writer、writer_version、row_index_stride、content_length、file_footer_length、file_length等读取方法read(columnsNone)读整文件为Tableread_stripe(n, columnsNone)读单个 stripe 为RecordBatch。columns支持嵌套字段前缀选择如a会选中a.b、a.c、a.d.e。ORCWriter单文件写出ORCWriter(where, *, ...)的全部参数及默认值见 _orc_writer_args_docsfile_version默认0.12ORC 文件版本可选0.11Hive 0.11 / ORC v0旧版与0.12Hive 0.12 / ORC v1新版batch_size默认1024ORC writer 一次写入的行数stripe_size默认64 * 1024 * 1024每个 ORC stripe 的大小字节compression默认uncompressed压缩编解码器合法值{UNCOMPRESSED, SNAPPY, ZLIB, LZ4, ZSTD}注意 LZO 目前不支持compression_block_size默认64 * 1024压缩块大小字节compression_strategy默认speed压缩策略{SPEED, COMPRESSION}权衡速度与压缩率row_index_stride默认10000行索引中每个条目覆盖的行数padding_tolerance默认0.0填充容差dictionary_key_size_threshold默认0.0字典键大小阈值0禁用字典编码1始终启用bloom_filter_columns默认None启用布隆过滤器的列bloom_filter_fpp默认0.05布隆过滤器误报率上限。import pyarrow.orc as orc # 写出 orc.write_table(table, example.orc, compressionZLIB, stripe_size64 * 1024 * 1024) # 读取 t orc.read_table(example.orc) # 或细粒度访问 with orc.ORCFile(example.orc) as f: # 支持上下文管理器 print(f.schema, f.nrows, f.nstripes) batch f.read_stripe(0)write_table(table, where, *, ...)与ORCWriter参数完全一致read_table(source, columnsNone, filesystemNone)支持本地路径或 URI 解析文件系统见 read_table。五格式横向对比与选型建议维度CSVFeatherJSONParquetORC数据类型文本Arrow IPC 列式文本NDJSON列式压缩列式压缩写入支持有write_csv/CSVWriter有write_feather无只读有write_table/ParquetWriter有write_table/ORCWriter类型推断ConvertOptions精细控制天然保留 schemaexplicit_schema/推断schema 内嵌schema 内嵌复杂能力无效行处理、字典编码、多线程压缩lz4/zstd意外字段策略分区、谓词下推、加密、统计信息stripe 级读取、布隆过滤典型场景与外部系统交换文本数据进程内/本地高速交换日志、API 输出的 NDJSON数仓、分析型查询Hive 生态列式存储选型要点基于本文 API 面推断可结合具体负载验证Feather面向「Arrow 生态内部快速往返」读写开销最小Parquet面向「分析型大规模存储」ParquetDatasetfilterspre_buffer组合适合分布式文件系统上的扫描ORC与 Hive 生态深度绑定ORCWriter的 stripe/布隆过滤参数适合需要行组级跳过读取的 Hive 工作负载CSV/JSON面向「异构系统文本交换」用ConvertOptions/ParseOptions解决脏数据与类型推断问题。延伸阅读其余 API 参考IPCs流式/文件格式见 docs/source/python/api/ipc.rst表格与数组见 tables.rst 与 arrays.rst数据集高层 API 见 dataset.rst各模块完整单元测试是理解行为边界的绝佳材料test_csv.py、test_json.py、test_feather.py、test_orc.pyParquet 测试位于python/pyarrow/tests/parquet/目录C 底层实现可追溯至 cpp/src/arrow/csv/reader.cc 等Cython 绑定层与高层封装的对应关系见 python/pyarrow/_csv.pyx、python/pyarrow/parquet/core.py。【免费下载链接】arrowApache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing项目地址: https://gitcode.com/gh_mirrors/arrow13/arrow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表