如何用Polars高效计算满足阈值条件的下一行行间距(大数据集)
问题
给定如下Polars DataFrame:
import polars as pl df = pl.DataFrame({ "Column A": [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], "Column B": [2, 3, 1, 4, 1, 7, 3, 2, 12, 0] })
需要创建新列Column C,记录当前行Column B的值到下一个满足「大于等于当前值150%」条件的Column B值所在行的行数差(行间距)。预期结果如下:
df_result = pl.DataFrame({ "Column A": [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], "Column B": [2, 3, 1, 4, 1, 7, 3, 2, 12, 0], "Column C": [1, 4, 1, 2, 1, 3, 2, 1, None, None] })
由于处理的是大型DataFrame,需用Polars高效实现该需求。
高效实现方案
针对大型数据集,必须规避逐行循环,利用Polars的向量化操作结合Numpy的广播特性,实现高性能计算。以下提供两种方案,分别适配不同规模的数据:
方案一:Polars原生API实现(中小规模数据)
代码简洁可读性强,适合十万级以内的数据:
import polars as pl df = pl.DataFrame({ "Column A": [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], "Column B": [2, 3, 1, 4, 1, 7, 3, 2, 12, 0] }) # 添加行索引用于计算间距 df = df.with_row_index("row_idx") # 为每行找到后续第一个满足条件的行索引 df = df.with_columns( pl.col("row_idx").map_elements( lambda idx: ( df.filter(pl.col("row_idx") > idx) .filter(pl.col("Column B") >= 1.5 * df.get_row(idx)["Column B"]) .select("row_idx") .min() ), return_dtype=pl.Int64 ).alias("target_idx") ) # 计算行间距,无匹配项设为None df = df.with_columns( (pl.col("target_idx") - pl.col("row_idx")).alias("Column C") ).drop("row_idx", "target_idx") print(df)
方案二:广播+Numpy优化(大规模数据)
利用Numpy的向量化广播,避免逐行Filter操作,百万级数据处理速度提升显著:
import polars as pl import numpy as np df = pl.DataFrame({ "Column A": [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], "Column B": [2, 3, 1, 4, 1, 7, 3, 2, 12, 0] }) # 提取列数据转为Numpy数组 b_values = df["Column B"].to_numpy() row_indices = np.arange(len(df)) # 广播生成条件矩阵:每行i与所有j>i的行比较是否满足B[j] >= 1.5*B[i] condition_mask = b_values[:, np.newaxis] * 1.5 <= b_values # 仅保留j > i的上三角区域 condition_mask = np.triu(condition_mask, k=1) # 找到每行第一个满足条件的列索引 first_match_idx = np.argmax(condition_mask, axis=1) # 标记无匹配项的行 no_match = ~condition_mask.any(axis=1) first_match_idx[no_match] = -1 # 计算行间距,无匹配项设为None column_c = np.where(first_match_idx != -1, row_indices[first_match_idx] - row_indices, None) # 将结果合并回原DataFrame df = df.with_columns(pl.Series("Column C", column_c, dtype=pl.Int64)) print(df)
方案选择建议
- 中小规模数据(<10万行):优先方案一,代码简洁易维护。
- 大规模数据(>100万行):选择方案二,若内存不足可进一步拆分数据分块处理,平衡内存占用与计算速度。
内容的提问来源于stack exchange,提问作者Levi Santos
相关产品推荐
相关产品推荐

