Polars性能优化:partition_by与collect_all提速原因探究
示例设置
警告:创建占用5GB内存的DataFrame
import time import numpy as np import polars as pl rng = np.random.default_rng(1) nrows = 50_000_000 df = pl.DataFrame( dict( id=rng.integers(1, 50, nrows), id2=rng.integers(1, 500, nrows), v=rng.normal(0, 1, nrows), v1=rng.normal(0, 1, nrows), v2=rng.normal(0, 1, nrows), v3=rng.normal(0, 1, nrows), v4=rng.normal(0, 1, nrows), v5=rng.normal(0, 1, nrows), v6=rng.normal(0, 1, nrows), v7=rng.normal(0, 1, nrows), v8=rng.normal(0, 1, nrows), v9=rng.normal(0, 1, nrows), v10=rng.normal(0, 1, nrows), ) )
我有如下数据处理任务:
start = time.perf_counter() res = ( df.lazy() .with_columns( pl.col(f"v{i}") - pl.col(f"v{i}").mean().over("id", "id2") for i in range(1, 11) ) .group_by("id", "id2") .agg((pl.col(f"v{i}") * pl.col("v")).sum() for i in range(1, 11)) .collect() ) time.perf_counter() - start # 9.85
该任务在16核机器上耗时约10秒。但如果先将df按id分区,再执行相同计算,最后调用collect_all和concat,可获得近2倍提速:
start = time.perf_counter() res2 = pl.concat( pl.collect_all( dfi.lazy() .with_columns( pl.col(f"v{i}") - pl.col(f"v{i}").mean().over("id", "id2") for i in range(1, 11) ) .group_by("id", "id2") .agg((pl.col(f"v{i}") * pl.col("v")).sum() for i in range(1, 11)) for dfi in df.partition_by("id", maintain_order=False) ) ) time.perf_counter() - start # 5.60
若按id2分区,耗时更短,约4秒。且第二种方法(按id或id2分区)的CPU利用率更高。现提出技术疑问:
- 为何第二种方法速度更快且CPU利用率更高?
- 窗口/分组操作应能并行处理各窗口/分组、充分利用硬件资源,为何两种方法性能存在差异?
内容的提问来源于stack exchange,提问作者lebesgue
相关产品推荐
相关产品推荐

