PyArrow Dataset对分区Parquet文件过滤无效问题解决
问题描述
使用PyArrow的pq.write_to_dataset将Pandas DataFrame保存为按pnr_group分区的Parquet文件后,直接调用pq.read_table读取单个文件可正常获取数据,但使用Dataset API进行过滤时返回空结果。
保存分区文件的代码
import pyarrow as pa import pyarrow.parquet as pq table = pa.Table.from_pandas(my_df) pq.write_to_dataset(table, root_path="data/bfl", partition_cols=['pnr_group'])
生成的目录结构
data/bfl/pnr_group=0/319a1fb5557a342c1b55356ce5123123-0.parquet
过滤代码(返回空DataFrame)
import pyarrow.dataset as ds import pyarrow as pa # 此代码能找到文件但返回空数据 dataset = ds.dataset( 'data/bfl', format='parquet', partitioning=ds.DirectoryPartitioning.discover(['pnr_group']) ) filter_expr = ds.field('pnr_group') == '0' filtered_dataset = dataset.filter(filter_expr) df = filtered_dataset.to_table().to_pandas() # 返回空DataFrame
已知信息:数据集Schema显示pnr_group为字符串类型,dataset.files可正确列出所有Parquet文件,但过滤后无数据返回。
原因分析与解决方法
核心问题是分区值的数据类型不匹配:
- 保存分区时,原DataFrame中的
pnr_group为数值类型(如int),因此write_to_dataset生成的目录后缀是=0(对应整数0)。 - 但使用
DirectoryPartitioning.discover(['pnr_group'])时,PyArrow默认将分区列推断为字符串类型,导致过滤条件中的字符串'0'与实际存储的整数0无法匹配,最终返回空结果。
解决方法
方法1:显式指定分区列的数据类型
创建Dataset时,明确指定分区列的类型与原数据一致(比如整数类型):
import pyarrow.dataset as ds import pyarrow as pa # 定义分区Schema,指定pnr_group为整数类型 partition_schema = pa.schema([('pnr_group', pa.int64())]) partitioning = ds.DirectoryPartitioning(partition_schema) dataset = ds.dataset( 'data/bfl', format='parquet', partitioning=partitioning ) # 使用整数0作为过滤条件 filter_expr = ds.field('pnr_group') == 0 df = dataset.filter(filter_expr).to_table().to_pandas()
方法2:保存前将分区列转为字符串类型
如果需要分区列保持字符串类型,在保存前先转换原DataFrame的对应列:
import pyarrow as pa import pyarrow.parquet as pq # 将pnr_group转为字符串类型 my_df['pnr_group'] = my_df['pnr_group'].astype(str) table = pa.Table.from_pandas(my_df) pq.write_to_dataset(table, root_path="data/bfl", partition_cols=['pnr_group'])
之后使用字符串'0'进行过滤即可正常匹配。
方法3:自动发现分区类型后匹配过滤
先让PyArrow自动识别分区的实际类型,再使用对应类型的条件过滤:
import pyarrow.dataset as ds # 自动发现分区的Schema partitioning = ds.partitioning.discover() dataset = ds.dataset('data/bfl', format='parquet', partitioning=partitioning) # 查看分区列的实际类型 print(dataset.schema) # 根据Schema显示的类型编写过滤条件,比如int64类型就用整数0 filter_expr = ds.field('pnr_group') == 0 df = dataset.filter(filter_expr).to_table().to_pandas()
验证步骤
- 检查原DataFrame中
pnr_group的类型:print(my_df['pnr_group'].dtype) - 查看Dataset的Schema:
print(dataset.schema),确保过滤条件的类型与Schema定义完全一致
内容的提问来源于stack exchange,提问作者FooBar
相关产品推荐
相关产品推荐

