2亿行大型Parquet文件简单等值查询速度过慢的优化方案咨询
将目录下多个CSV文件合并转换为单个Parquet文件,总数据量超2.05亿行,转换代码如下:
from pyarrow import csv, parquet import pyarrow as pa import pyarrow.dataset as ds csvDir = 'my_csv_dir' fields = [ ('workId',pa.string()), ('authorId',pa.string()) ] schema = pa.schema(fields) dataset = ds.dataset(csvDir,format="csv",schema=schema) table = dataset.to_table() print(table.num_rows,'rows') parquet.write_table(table, 'my_parquet.parquet')
通过以下代码查看Parquet文件元信息:
import pyarrow.parquet as pq parquet_file = pq.ParquetFile('my_parquet.parquet') print(parquet_file.metadata) print(parquet_file.schema)
元信息输出:
<pyarrow._parquet.FileMetaData object at 0x7f79d36cb8b0> created_by: parquet-cpp-arrow version 8.0.0 num_columns: 2 num_rows: 205841002 num_row_groups: 4 format_version: 1.0 serialized_size: 1429 <pyarrow._parquet.ParquetSchema object at 0x7f795a15e500> required group field_id=-1 schema { optional binary field_id=-1 workId (String); optional binary field_id=-1 authorId (String); }
当前针对该文件的简单查询运行速度极慢,两种查询方式的表现如下:
DuckDB查询
查询代码:
import duckdb print(duckdb.query("select count(*) from'my_parquet.parquet'").fetchall()) print(duckdb.query("select workId,authorId from'my_parquet.parquet' where workId = 'W2137422493'").fetchall())
运行耗时:
time python query.py [(205841002,)] [('W2137422493', 'A1461130442 A2023515231')] real 0m10.690s user 0m30.743s sys 0m1.277s
PyArrow Dataset过滤查询
查询代码:
import pyarrow.dataset as ds dataset = ds.dataset('my_parquet.parquet') print(dataset.files) print(dataset.schema.to_string()) print(dataset.to_table(filter=ds.field('workId') == 'W2137422493').to_pandas())
运行耗时:
['my_parquet.parquet'] workId: string authorId: string workId authorId 0 W2137422493 A1461130442 A2023515231 real 1m23.121s user 1m23.197s sys 0m48.689s
已做尝试和现存疑问:
- 未找到为
workId字段构建索引提速的方法 - 已知Parquet有列索引特性,但了解到该特性多用于范围查询,不符合当前等值查询场景
- 曾尝试移除
workId前缀"W"将字段转为整数类型,同时对全量数据按workId排序,认为Arrow按文件名读取时会保持数值顺序,但调整后查询速度无提升,仍不清楚如何触发Parquet生成列索引
单文件仅4个行组、未开启细粒度统计信息是慢的根本原因:
- 2亿行仅拆分4个行组,平均每个行组超5000万行,单组大小达GB级。Parquet谓词下推、索引过滤的最小粒度是行组/数据页,行组过大时,哪怕有过滤条件也无法跳过大量无关数据,基本等于全表扫描。
- 直接调用
write_table写单文件默认不会开启细粒度列索引、数据页统计信息,谓词下推完全无法生效。 - 之前排序后无效果,是因为4个超大行组覆盖了全量
workId的取值范围,每个行组的min/max区间都包含查询目标值,统计信息根本过滤不掉任何行组。 - 对Parquet列索引的认知存在偏差:列索引存储的是每个数据块的min/max统计值,等值查询同样可以用——只要目标值不在块的[min, max]区间内,就可以直接跳过整个块,不需要扫描块内数据。
重写Parquet文件,拆分合理大小的行组,开启统计信息和列索引
不要一次性把全量数据写成单个大文件,用pyarrow.dataset.write_dataset分块写入,控制单行组大小在100MB500MB区间(对应行数10万100万,根据单条记录大小调整),写之前按workId排序保证同值数据尽量集中,同时显式开启列统计和页索引:import pyarrow.dataset as ds import pyarrow as pa csvDir = 'my_csv_dir' # workId转整数的优化保留,整数比较比字符串快3~5倍 fields = [ ('workId',pa.int64()), ('authorId',pa.string()) ] schema = pa.schema(fields) dataset = ds.dataset(csvDir,format="csv",schema=schema) # 按workId排序,让同值数据集中在连续的行组 sorted_table = dataset.to_table().sort_by('workId') # 分块写入Parquet目录 ds.write_dataset( sorted_table, 'optimized_parquet', format='parquet', existing_data_behavior='overwrite_or_ignore', max_rows_per_row_group=200000, # 单组20万行,两列场景下单组大小约20~50MB,符合最优区间 use_threads=True, file_options=ds.ParquetFileFormat().make_write_options( write_statistics=True, # 开启行组/数据页的min/max统计,谓词下推的基础 column_index_truncate_length=4096, # 不截断列索引的统计值,避免字符串/长数值统计不准 data_page_version='2.0' # 启用2.0数据页格式,支持更细粒度的页级索引 ) )这一步完成后,等值查询可以通过统计信息直接跳过99%以上的无关行组,扫描数据量会降到原来的几十分之一。
查询优先用DuckDB,直接读取分块后的Parquet目录
DuckDB对Parquet列索引、谓词下推、并行扫描的优化远好于PyArrow原生Dataset,重写后的数据集用DuckDB做点查,耗时可以降到100毫秒级别。注意不要先把全量数据转成Pandas再过滤,所有过滤逻辑直接写在SQL里,让DuckDB在扫描阶段就完成过滤。高频查询场景直接导入DuckDB持久化存储
如果这类点查频率很高,可以直接把清洗后的数据导入DuckDB本地数据库,DuckDB的存储引擎会自动生成更细粒度的Zone Map索引,点查性能还能再提升一个量级,2亿行数据的等值查询基本可以做到10毫秒级返回。
内容的提问来源于stack exchange,提问作者Sol

