Polars流引擎处理Parquet文件为何远慢于循环遍历实现?
Polars流处理比循环分块慢的原因分析
问题场景
我有一个Parquet文件目录,需要对所有文件应用函数并计算均值。原本以为Polars的LazyFrame与流处理功能会在此场景表现出色,但实际发现循环遍历分块处理的简易方法性能更优。
两段实现代码
循环分块实现(性能更优)
%%time unique_years = ( pl.scan_parquet(data_glob) .select(pl.col("time").dt.year().unique()) .collect() .to_series() .to_list() ) results = [] for year in tqdm(unique_years): results.append( pl.scan_parquet(data_glob) .filter(pl.col("time").dt.year()==year) .select(["time", "market", "close"]) .sort("time") .with_columns(pl.col("close").log().diff().over("market").alias("log_returns")) .group_by("time") .agg(pl.col("log_returns").mean()) .collect() )
执行耗时:CPU times: user 10min 10s, sys: 4min 44s, total: 14min 55s;Wall time: 2min 4s
流处理实现(效率极低)
%%time df = ( pl.scan_parquet(data_glob) .select(["time", "market", "close"]) .sort("time") .with_columns(pl.col("close").log().diff().over("market").alias("log_returns")) .group_by("time") .agg(pl.col("log_returns").mean()) .collect(streaming=True) )
核心原因分析
- 全局排序的开销差异:流处理版本的
sort("time")是全局排序,需要读取全量数据完成排序,IO和计算成本极高;而循环版本每次仅处理单一年份数据,排序规模极小,耗时大幅降低。 - 窗口操作的执行效率:
diff().over("market")需要在market分组内按time排序计算差值。循环版本中单年份的market分组数据量小,窗口操作可在内存高效完成;流处理版本则需跨分片维护状态,频繁的磁盘读写和状态同步拖慢了速度。 - 谓词下推的优化差异:循环版本的
filter会被Polars推到Parquet读取阶段,直接跳过不符合年份的行,实际读取的数据量远小于全量;流处理版本无提前过滤,需读取所有数据。 - 流处理的场景不匹配:Polars流处理更适合无法一次性放入内存的超大数据集,且操作链无全局排序、跨分片窗口这类高开销操作的场景。你的场景按年份拆分后的数据可完全放入内存,循环分块更贴合数据的自然处理逻辑。
优化建议
如果想保留流处理写法并提升性能,可尝试:
- 先按年份过滤再执行后续操作,将全局排序转化为分区内排序;
- 若Parquet文件本身按年份分区,直接指定分区路径扫描,避免全量扫描;
- 调整
sort操作位置,尽量在过滤后执行,缩小排序数据集规模。
内容的提问来源于stack exchange,提问作者J Griffiths
相关产品推荐
相关产品推荐

