使用PyArrow Dataset API获取分区Parquet文件公共Schema异常
问题:PyArrow Dataset 推断Schema未返回公共列,而是首个文件的Schema
我正在处理按月分区的Parquet文件,文件间存在列增删差异,原本期望以下代码能返回数据集的公共Schema:
import pyarrow.dataset as ds import pyarrow.parquet as pq dataset = ds.dataset('data/Parquet') dataset.files # Infer a common schema for all files in the dataset schema = dataset.schema # Loop over each file in the dataset and count the inferred schema columns that exist in each file schema for file in dataset.files: parquet_file = pq.ParquetFile(file) file_schema = parquet_file.schema matching_cols = set(schema.names).intersection(set(file_schema.names)) non_matching_cols = set(schema.names) - matching_cols print("File:", file) print("Total matching columns:", len(matching_cols)) print("Total non-matching columns:", len(non_matching_cols)) print("")
但实际检查发现ds.dataset推断的Schema是分区中首个文件的完整副本,输出结果如下:
File: data/Parquet/File1.parquet Total matching columns: 14527 Total non-matching columns: 0 File: data/Parquet/File2.parquet Total matching columns: 11522 Total non-matching columns: 3005 File: data/Parquet/File3.parquet Total matching columns: 11637 Total non-matching columns: 2890
根据文档描述,schema应该是整个数据集的公共Schema,我原以为会返回所有文件的共有列。请问这是预期行为吗?
回答
这不是预期行为,PyArrow Dataset 默认的Schema推断逻辑会优先采样少量文件(比如第一个文件)来快速生成Schema,尤其是数据集文件较多时,不会默认扫描所有文件去计算真正的公共Schema。
要获取所有文件的共有列Schema,有两种可行方式:
方式1:创建Dataset时强制统一Schema
创建Dataset时指定unify_schemas=True,强制扫描所有文件并生成公共Schema:
import pyarrow.dataset as ds # 启用Schema统一,生成所有文件的公共Schema dataset = ds.dataset('data/Parquet', unify_schemas=True) schema = dataset.schema
方式2:手动计算公共Schema
如果需要更精细的控制,可以手动收集所有文件的Schema,再计算它们的交集:
import pyarrow.dataset as ds import pyarrow.parquet as pq dataset = ds.dataset('data/Parquet') # 收集所有文件的Arrow Schema all_schemas = [] for file in dataset.files: with pq.ParquetFile(file) as pf: all_schemas.append(pf.schema.to_arrow_schema()) # 计算所有Schema的交集(公共列) common_schema = all_schemas[0] for schema in all_schemas[1:]: common_schema = common_schema.intersect(schema) # 使用公共Schema重新加载Dataset dataset = ds.dataset('data/Parquet', schema=common_schema)
注意:如果数据集文件数量极大,扫描所有文件会增加初始化时间,但能确保得到真正的共有列Schema。
内容的提问来源于stack exchange,提问作者Stretch
相关产品推荐
相关产品推荐

