Polars如何高效获取满足条件行的±N行窗口数据?
如何用Polars获取满足条件行的±N行窗口数据
现有如下Polars DataFrame:
df = pl.DataFrame({ "letters": ["A", "B", "C", "D", "E", "F", "G", "H"], "values": ["aa", "bb", "cc", "dd", "ee", "ff", "gg", "hh"] }) print(df)
输出结果:
shape: (8, 2) ┌─────────┬────────┐ │ letters ┆ values │ │ --- ┆ --- │ │ str ┆ str │ ╞═════════╪════════╡ │ A ┆ aa │ │ B ┆ bb │ │ C ┆ cc │ │ D ┆ dd │ │ E ┆ ee │ │ F ┆ ff │ │ G ┆ gg │ │ H ┆ hh │ └─────────┴────────┘
需求:获取满足指定条件的行周围±N行的窗口数据。例如条件为pl.col("letters").contains("D|F"),N=2时,预期输出为:
┌─────────┬────────────────────────────────┐ │ letters ┆ output │ │ --- ┆ --- │ │ str ┆ list[str] │ ╞═════════╪════════════════════════════════╡ │ D ┆ ["bb", "cc", "dd", "ee", "ff"] │ │ F ┆ ["dd", "ee", "ff", "gg", "hh"] │ └─────────┴────────────────────────────────┘
注意:场景中窗口存在重叠,实际N值较大(约10-20),且数据集规模大,需要高效低内存的实现方案。
补充说明:以下是可得到正确结果的DuckDB SQL查询语句,需转换为Polars实现:
df_table = df.to_arrow() con = duckdb.connect() query = """ SELECT letters, list(values) OVER ( ROWS BETWEEN 2 PRECEDING AND 2 FOLLOWING ) as combined FROM df_table QUALIFY letters in ('D', 'F') """ print(pl.from_arrow(con.execute(query).arrow()))
输出结果:
shape: (2, 2) ┌─────────┬────────────────────────┐ │ letters ┆ combined │ │ --- ┆ --- │ │ str ┆ list[str] │ ╞═════════╪════════════════════════╡ │ D ┆ ["bb", "cc", ... "ff"] │ │ F ┆ ["dd", "ee", ... "hh"] │ └─────────┴────────────────────────┘
各方案基准测试
测试环境:Amazon ml.c5.xlarge机器的Jupyter Notebook,通过htop监控CPU与内存使用,数据集包含1200万+行。测试了两种Polars方案(即时/延迟API)、Python循环方案及DuckDB方案,结果如下:
汇总表
Polars的@jqurious方案基于shift()的无拷贝实现,性能稳健且内存占用合理;优化后的Python循环方案性能相当;DuckDB在速度与内存占用上表现较差。Polars与DuckDB均仅使用单核。
| 方法 | CPU使用 | 内存使用 | 耗时 |
|---|---|---|---|
| ΩΠΟΚΕΚΡΥΜΜΕΝΟΣ方案 | 单核 | 内存暴增 | - |
| jqurious方案 | 单核 | 2.53G 保持2.53G | 4.63 s |
| (优化)Python循环方案 | 单核 | 2.53G 升至2.58G | 4.91 s |
| DuckDB方案 | 单核 | 1.62G 升至6.13G | 38.6 s |
- CPU使用:显示操作期间是否占用多核
- 内存使用:显示操作前内存占用与操作期间最大内存占用
ΩΠΟΚΕΚΡΥΜΜΕΝΟΣ的方案
preceding = 2 following = 2 look_around = [pl.col("body").shift(-i) for i in range(-preceding, following + 1)] ( df .with_columns( pl.when(pl.col('body').str.contains(regex)) .then(pl.concat_list(look_around)) .alias('combined') ) .filter(pl.col('combined').is_not_null()) )
该方案在大数据集上导致内存暴增,无论使用即时还是延迟API均会使内核崩溃。
jqurious的方案
preceding = 2 following = 2 look_around = [ pl.col("body").shift(-i).alias(f"lag_{i}") for i in range(-preceding, following + 1) ] ( df .with_columns( look_around ) .filter(pl.col("body").str.contains(regex)) .select( pl.col("body"), pl.concat_list([f"lag_{i}" for i in range(-2, 3)]).alias("output") ) )
- 即时模式:
- CPU使用: 单核
- 内存使用: 2.53G → 2.53G
- 耗时: 4.63 s ± 6.6 ms/循环(7次运行均值±标准差,每次1循环)
- 延迟模式:
- CPU使用: 单核
- 内存使用: 2.53G → 2.53G
- 耗时: 4.63 s ± 3.85 ms/循环(7次运行均值±标准差,每次1循环)
(优化)Python循环方案
preceding = 2 following = 2 output = [] indices = df.with_row_index().select( pl.col("index").filter(pl.col("body").str.contains(regex)) )["index"] for idx, x in enumerate(indices): offset = max(0, x - preceding) length = preceding + following + 1 output.append(df["body"].slice(offset, length))
- CPU使用: 单核
- 内存使用: 2.53G → 2.58G
- 耗时: 4.91 s ± 24.5 ms/循环(7次运行均值±标准差,每次1循环)
DuckDB方案
注:先将df转换为Arrow.Table,使DuckDB可直接操作。不确定结果转回Arrow是否占用大量计算资源,对其是否公平存疑。
preceding = 2 following = 2 query = f""" SELECT body, list(body) OVER ( ROWS BETWEEN {preceding} PRECEDING AND {following} FOLLOWING ) as combined FROM df_table QUALIFY regexp_matches(body, '{regex}') """ result = con.execute(query).arrow()
首次运行计算崩溃,需直接读取Arrow Table而非通过Polars(节省约1GB内存)重试:
- 首次尝试:
- CPU: 单核
- 内存: 2.53G → 6.93G → 崩溃!
- 耗时: 无
- 二次尝试:
- CPU: 单核
- 内存: 1.62G → 6.13G
- 耗时: 38.6 s ± 311 ms/循环(7次运行均值±标准差,每次1循环)
内容的提问来源于stack exchange,提问作者anonymous1a
相关产品推荐
相关产品推荐

