如何基于40天全量数据快速计算DataFrame前几小时的滚动中位数?
问题
我用Python计算DataFrame的滚动中位数,脚本每15分钟执行一次。原全量计算方式耗时过长,我仅需要计算最近4小时的滚动中位数,但计算必须基于40天的全量历史数据。
原全量计算代码如下:
df[column+'_ROLLING_MEDIAN']=df.groupby(grouping)[column].transform(lambda x: x.rolling(window=window_size, min_periods=minimum, center=False).median())
我尝试自定义函数仅计算目标时段的结果,但运行速度比原生滚动函数更慢:
def rolling_function_med(series): counter = 0 for row_index in series.index: if row_index[2] >= day_cutoff: series[row_index] = series.iloc[counter:].median() counter+=1 else: break series.iloc[counter+1:] = None return series df[column+'_ROLLING_MEDIAN']=df.groupby(grouping)[column].transform(rolling_function_med )
是否有更高效的实现方式?附可复现代码:
import pandas as pd from datetime import datetime, timedelta data = {'SERVICE_NAME': {0: 'Foo', 1: 'Foo', ...}, 'HOUR': {0:0,1:0,...}, ...} df = pd.DataFrame(data) df = df.set_index(['SERVICE_NAME','HOUR','TRANSACTION_DATETIME_GROUPED_CST','MINUTE_GROUPING']) min_time = min(df.index.get_level_values(level = 2)) day_cutoff = (min_time + timedelta(days=40) - timedelta(hours=4)) # 自定义函数同上 column = 'SUM_EXEC_TIME_IN_MILLIS_RATE' grouping = ['SERVICE_NAME','HOUR', 'MINUTE_GROUPING'] df[column+'_ROLLING_MEDIAN']=df.groupby(grouping)[column].transform(rolling_function_med ) print(df)
高效实现方案
核心优化思路
自定义函数慢的根源是手动循环遍历行,完全没有利用pandas的向量化运算优势。我们可以基于原生rolling方法,结合掩码过滤仅保留目标时段的结果,既保证计算效率,又满足仅输出最近4小时数据的需求。
优化代码示例
def optimized_rolling_median(series): # 提取时间索引列(对应TRANSACTION_DATETIME_GROUPED_CST层级) time_idx = series.index.get_level_values(2) # 生成目标时段的掩码:仅保留最近4小时的行 target_mask = time_idx >= day_cutoff # 用原生rolling计算全量滚动中位数(C实现,速度远快于手动循环) full_rolling_med = series.rolling( window=window_size, min_periods=minimum, center=False ).median() # 仅保留目标时段的结果,其余设为None full_rolling_med[~target_mask] = None return full_rolling_med # 应用优化后的函数 df[column+'_ROLLING_MEDIAN'] = df.groupby(grouping)[column].transform(optimized_rolling_median)
额外提速建议
- 确保数据按时间排序:每个分组内的数据必须按
TRANSACTION_DATETIME_GROUPED_CST升序排列,否则滚动窗口计算逻辑会出错,还可能影响性能。可在分组前执行排序:df = df.sort_index(level=['SERVICE_NAME','HOUR','MINUTE_GROUPING','TRANSACTION_DATETIME_GROUPED_CST']) - 简化索引层级:如果部分索引层级不影响分组或时间判断,可考虑删除,降低索引操作的开销。
- 多进程加速:若数据量极大,单进程pandas仍不够快,可尝试用Dask DataFrame实现并行计算,其API与pandas兼容,能充分利用多核心资源。
内容的提问来源于stack exchange,提问作者Michael Runyan
相关产品推荐
相关产品推荐

