You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 01:20:13