使用Python Polars实现高效滑动窗口加权内积(内存优化)
问题描述
我有一个权重向量:
weight_vec = pl.Series("weights", [0.125, 0.0625, 0.03125])
同时有一个可包含最多m个变量的DataFrame(示例仅保留两个变量):
df = pl.DataFrame( { "row_index": [0, 1, 2, 3, 4], "var1": [1, 2, 3, 4, 5], "var2": [6, 7, 8, 9, 10], } )
该DataFrame的观测行数可能极大(数千万行)。
需求
- 针对每个变量的每个观测值x_i(i为行索引,范围[0,...,4]),将x_i转换为从当前观测值开始的连续n个值(即[x_i, x_i+1, ..., x_i+n-1])与权重向量的点积,其中n为权重向量的长度,且n会随权重向量定义不同而变化。当行索引接近末尾(例如:最大索引 - 行索引 + 1 < n)时,对应值设为None。
数值计算示例:- var1在索引0处的值:
0.125*1 + 0.0625*2 + 0.03125*3 = 0.34375 - var2在索引2处的值:
0.125*8 + 0.0625*9 + 0.03125*10 = 1.875
- var1在索引0处的值:
- 可假设DataFrame的行数始终大于等于权重向量长度,以保证至少有一个有效结果。
期望结果
shape: (5, 3) ┌───────────┬─────────┬─────────┐ │ row_index ┆ var1 ┆ var2 │ │ --- ┆ --- ┆ --- │ │ i64 ┆ f64 ┆ f64 │ ╞═══════════╪═════════╪═════════╡ │ 0 ┆ 0.34375 ┆ 1.4375 │ │ 1 ┆ 0.5625 ┆ 1.65625 │ │ 2 ┆ 0.78125 ┆ 1.875 │ │ 3 ┆ null ┆ null │ │ 4 ┆ null ┆ null │ └───────────┴─────────┴─────────┘
实现方案
方案一:基于Rolling窗口的点积计算
利用Polars原生滚动窗口操作,结合向量点积实现,逻辑直观且符合向量化要求:
import polars as pl # 定义权重向量和DataFrame weight_vec = pl.Series("weights", [0.125, 0.0625, 0.03125]) df = pl.DataFrame( { "row_index": [0, 1, 2, 3, 4], "var1": [1, 2, 3, 4, 5], "var2": [6, 7, 8, 9, 10], } ) n = len(weight_vec) weights = weight_vec.to_list() # 定义滑动点积函数 def sliding_dot(col: pl.Series) -> pl.Series: return col.rolling(window_size=n, min_periods=n).apply( lambda window: window.dot(weights), return_dtype=pl.Float64 ) # 应用到所有变量列,保留row_index result = df.with_columns( [sliding_dot(pl.col(col)).alias(col) for col in df.columns if col != "row_index"] ) print(result)
方案二:移位广播求和(更高性能)
通过移位操作生成多列,再横向加权求和,完全避免Python层面的循环,内存和性能更优,适合超大数据集:
import polars as pl # 定义权重向量和DataFrame weight_vec = pl.Series("weights", [0.125, 0.0625, 0.03125]) df = pl.DataFrame( { "row_index": [0, 1, 2, 3, 4], "var1": [1, 2, 3, 4, 5], "var2": [6, 7, 8, 9, 10], } ) n = len(weight_vec) weights = weight_vec.to_list() var_cols = [col for col in df.columns if col != "row_index"] # 生成每个移位后的加权列 shift_exprs = [] for i in range(n): shift_exprs += [pl.col(col).shift(-i) * weights[i] for col in var_cols] # 按变量分组求和,自动处理末尾null result = df.with_columns(shift_exprs).with_columns( [pl.sum_horizontal(pl.col(f"{col}_*")).alias(col) for col in var_cols] ).select(df.columns) print(result)
方案说明
- 两种方案均为向量化操作,避免了逐行循环,处理数千万行数据时效率远高于Python循环
- 方案一中
min_periods=n参数确保只有窗口内有足够n个值时才计算点积,不足的自动填充None,正好匹配需求中末尾行的处理逻辑 - 方案二通过移位广播实现,完全基于Polars原生表达式,性能更优,内存占用更可控
内容的提问来源于stack exchange,提问作者Kevin Li
相关产品推荐
相关产品推荐

