手工构建Parquet分区路径引发ArrowTypeError类型不兼容问题求助
解决Parquet分区读取类型冲突与写入覆盖问题
一、解决分区目录读取的类型冲突问题
报错ArrowTypeError: Unable to merge: Field month has incompatible types: int32 vs dictionary<values=int32, indices=int32, ordered=0>是因为不同分区文件中month字段的编码类型不一致(部分为int32,部分为字典编码),可通过以下两种方式处理:
方法1:读取时强制指定字段类型
直接在pd.read_parquet中通过dtype参数强制统一month的类型为int32,让Arrow在合并时自动转换:
import pandas as pd df = pd.read_parquet( r'C:\Datasets\cn_data\dm\qmt\wqa_mfeatures\30m\year=2020', dtype={'month': 'int32'} )
方法2:使用PyArrow Dataset API统一类型
如果方法1无效,可借助PyArrow的Dataset接口手动指定字段类型后再读取:
import pyarrow as pa import pyarrow.dataset as ds import pandas as pd # 加载分区数据集 dataset = ds.dataset( r'C:\Datasets\cn_data\dm\qmt\wqa_mfeatures\30m\year=2020', format='parquet' ) # 构建新的schema,将month字段强制设为int32 new_schema = pa.schema( [(field.name, pa.int32() if field.name == 'month' else field.type) for field in dataset.schema] ) # 按新schema读取并转换为DataFrame table = dataset.to_table(schema=new_schema) df = table.to_pandas()
二、解决写入时的旧文件覆盖问题
由于pd.DataFrame.to_parquet配合partition_cols不会自动删除旧分区文件,重跑时会生成重复文件,可通过写入前手动清理目标分区文件的方式解决,同时不影响其他分区:
自定义覆盖写入函数
import os import pandas as pd def write_parquet_overwrite(df, base_path, partition_cols): # 提取当前数据的分区值(假设同分区内的year/month值一致) partition_values = {col: df[col].unique()[0] for col in partition_cols} # 拼接目标分区目录路径 partition_dir = os.path.join(base_path, *[f"{k}={v}" for k, v in partition_values.items()]) # 清理分区目录下的所有Parquet文件 if os.path.exists(partition_dir): for filename in os.listdir(partition_dir): if filename.endswith('.parquet'): os.remove(os.path.join(partition_dir, filename)) # 写入新数据 df.to_parquet(base_path, partition_cols=partition_cols)
使用示例
当需要重跑某批次数据时,直接调用该函数即可自动覆盖对应分区的旧文件:
# 假设processed_df是处理好的待写入数据 write_parquet_overwrite( df=processed_df, base_path=r'C:\Datasets\cn_data\dm\qmt\wqa_mfeatures\30m', partition_cols=['year', 'month'] )
内容的提问来源于stack exchange,提问作者Mithril
相关产品推荐
相关产品推荐

