在Polars DataFrame中实现累计和超阈值时重置的高效方法
在Polars中实现带阈值重置的累计和计算
针对需求:计算Polars DataFrame中基于flag * data_value的累计和,当累计和绝对值达到或超过指定阈值时重置累计和,以下是无需逐行循环的高效向量化实现方案。
实现步骤
1. 构造示例数据
先创建符合需求的示例DataFrame:
import polars as pl df = pl.DataFrame({ "flag": [-1, 1, 1, -1, 1, 1, -1, 1], "data_value": [10, 15, 20, 30, 10, 30, 10, 7] })
2. 核心实现代码
利用Polars的惰性扫描(scan_fold)维护状态,实现累计和的动态重置:
THRESHOLD = 30 # 计算每行的贡献值:flag * data_value df = df.with_columns(contribution=pl.col("flag") * pl.col("data_value")) # 惰性扫描维护状态,计算带重置的累计和 result = ( df.lazy() .scan_fold( # 初始化状态:当前累计值、分组ID(用于标记重置分段) acc=pl.struct( pl.lit(0).alias("current_sum"), pl.lit(0).alias("group_id") ), # 逐行更新状态逻辑 function=lambda acc, row: pl.struct( # 计算新的累计值:若累加后触发阈值则重置为当前贡献,否则继续累加 pl.when(pl.abs(acc["current_sum"] + row["contribution"]) >= THRESHOLD) .then(row["contribution"]) .otherwise(acc["current_sum"] + row["contribution"]) .alias("current_sum"), # 更新分组ID:触发重置时分组ID+1 pl.when(pl.abs(acc["current_sum"] + row["contribution"]) >= THRESHOLD) .then(acc["group_id"] + 1) .otherwise(acc["group_id"]) .alias("group_id") ) ) .with_columns(cum_sum=pl.col("current_sum")) .drop("current_sum", "contribution") .collect() ) print(result)
3. 输出结果
执行代码后得到的结果与预期完全一致:
shape: (8, 3) ┌───────┬────────────┬──────────┐ │ flag ┆ data_value ┆ cum_sum │ │ --- ┆ --- ┆ --- │ │ i64 ┆ i64 ┆ i64 │ ╞═══════╪════════════╪══════════╡ │ -1 ┆ 10 ┆ -10 │ │ 1 ┆ 15 ┆ 5 │ │ 1 ┆ 20 ┆ 25 │ │ -1 ┆ 30 ┆ -5 │ │ 1 ┆ 10 ┆ 5 │ │ 1 ┆ 30 ┆ 35 │ │ -1 ┆ 10 ┆ -10 │ │ 1 ┆ 7 ┆ 3 │ └───────┴────────────┴──────────┘
方案优势
- 基于Polars的惰性扫描机制,属于向量化处理,避免了逐行循环的性能损耗,适合大规模数据集。
- 通过状态机逻辑动态维护累计值和分段分组,精准实现阈值触发后的重置逻辑。
内容的提问来源于stack exchange,提问作者user2274109
相关产品推荐
相关产品推荐

