如何基于其他列最新记录计算Polars DataFrame的Rolling mean?
Polars 按执行时间计算带历史补充的滚动均值实现方案
核心需求
基于包含execution_time(执行时间)和valid_time(预测时间)的Polars DataFrame,为每行计算valid_time过去n小时内的均值:
- 同一
valid_time可能对应多个execution_time,需取**小于等于当前execution_time的最大execution_time**对应的value补充数据 - 当窗口内数据不足n个点时,返回
None
实现思路
- 补全历史数据:生成所有
execution_time与valid_time的组合,对缺失的value用更早执行时间的对应值向前填充,确保每个(execution_time, valid_time)都有可用数据 - 分组计算滚动均值:按
execution_time分组,对组内valid_time排序后,用固定窗口大小的滚动均值函数计算结果,不足n个点时返回None
完整代码实现
import polars as pl import datetime # 原始输入数据 df = pl.DataFrame({ 'execution_time': [datetime.datetime(2023,1,1,0)]*10 + [datetime.datetime(2023,1,1,6)]*6, 'valid_time': [datetime.datetime(2023,1,1,i) for i in range(10)] + [datetime.datetime(2023,1,1,i) for i in range(6,12)], 'value': range(16) }) n = 2 # 窗口覆盖的小时数(对应n个整点数据) # 步骤1:生成全量(execution_time, valid_time)组合并补全缺失值 unique_executions = df.select('execution_time').unique().sort('execution_time') unique_valid_times = df.select('valid_time').unique().sort('valid_time') # 生成所有可能的执行时间与有效时间组合 full_combinations = unique_executions.join(unique_valid_times, how='cross') # 左连接原始数据,匹配对应value full_combinations = full_combinations.join( df, on=['execution_time', 'valid_time'], how='left' ) # 按valid_time分组,按execution_time排序后向前填充缺失的value filled_df = full_combinations.sort(['valid_time', 'execution_time']).with_columns( pl.col('value').forward_fill().over('valid_time').alias('filled_value') ) # 步骤2:按execution_time分组计算滚动均值 grouped = filled_df.sort(['execution_time', 'valid_time']).group_by('execution_time').agg( pl.col('valid_time'), pl.col('filled_value'), # 窗口大小为n,不足n个点时返回None pl.col('filled_value').rolling_mean(window_size=n, min_periods=n).alias('rolling_mean') ).explode(['valid_time', 'filled_value', 'rolling_mean']) # 合并回原始数据,保留原始value列 final_result = df.join( grouped.select('execution_time', 'valid_time', 'rolling_mean'), on=['execution_time', 'valid_time'], how='left' ) # 打印结果 print(final_result)
结果验证
运行代码后得到的结果完全符合需求:
- 执行时间为
2023-01-01 00:00:00的行:valid_time=00:00:00因数据不足返回None,其余行取当前组内前n个整点的均值 - 执行时间为
2023-01-01 06:00:00的行:valid_time=06:00:00自动补充2023-01-01 00:00:00组的valid_time=05:00:00数据计算均值,其余行取当前组内的前n个整点均值
扩展说明
当n大于执行时间间隔(如n=7)时,代码依然有效:forward_fill会自动从所有更早的执行时间中补充对应valid_time的value,滚动均值会包含所有符合窗口要求的历史补全数据。
内容的提问来源于stack exchange,提问作者SergioGM
相关产品推荐
相关产品推荐

