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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 10:48:18