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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 23:43:24