考虑缺失值的Pandas滚动聚合实现方案问询
处理带缺失月度数据的分组滚动聚合(满足有效数据量阈值)
现有客户跨年度月度数据,部分客户存在缺失的月度记录(如示例中客户2缺失200102期),需要按客户分组实现以下滚动聚合需求:
max_vol_3:最近3个月度周期的volume最大值,要求窗口内至少有3个有效数据,否则返回Nonemax_vol_6:最近6个月度周期的volume最大值,要求窗口内至少有5个有效数据,否则返回Nonesum_trans_3:最近3个月度周期的num_transactions总和,要求窗口内至少有3个有效数据,否则返回None
原数据示例:
import pandas as pd df = pd.DataFrame({'cust_id': [1,1,1,1,1,1,2,2,2,2,2], 'period' : [200010,200011,200012,200101,200102,200103,200010,200011,200012,200101,200103], 'volume' : [1,2,3,4,5,6,7,8,9,10,12], 'num_transactions': [3,4,5,6,7,8,9,10,11,12,13]})
期望输出:
out = pd.DataFrame({'cust_id': [1,1,1,1,1,1,2,2,2,2,2], 'period' : [200010,200011,200012,200101,200102,200103,200010,200011,200012,200101,200103], 'max_vol_3' : [None, None, 3,4,5,6,None,None,9,10,None], 'max_vol_6' :[None,None,None,None,None,6,None,None,None,None,None], 'sum_trans_3': [None, None, 12, 15, 18, 21, None, None, 30, 33, None]})
实现步骤
核心思路是先补全每个客户的所有缺失月度,确保滚动窗口能覆盖连续的月度周期,再通过自定义函数判断窗口内有效数据量是否满足要求,最后筛选回原数据的行。
import pandas as pd # 1. 转换period为datetime格式(统一为月度最后一天) df['period'] = pd.to_datetime(df['period'], format='%Y%m') + pd.offsets.MonthEnd(0) # 2. 生成每个客户的完整月度序列,补全缺失行 all_periods = pd.date_range(df['period'].min(), df['period'].max(), freq='M') full_df = pd.DataFrame( pd.MultiIndex.from_product([df['cust_id'].unique(), all_periods], names=['cust_id', 'period']), columns=['cust_id', 'period'] ) # 合并原数据,缺失行的指标列自动填充为NaN full_df = full_df.merge(df, on=['cust_id', 'period'], how='left') # 3. 定义滚动聚合函数,判断有效数据量阈值 def calc_max_3(window): valid_count = window['volume'].count() return window['volume'].max() if valid_count >= 3 else None def calc_max_6(window): valid_count = window['volume'].count() return window['volume'].max() if valid_count >= 5 else None def calc_sum_trans_3(window): valid_count = window['num_transactions'].count() return window['num_transactions'].sum() if valid_count >= 3 else None # 4. 计算3个月窗口的指标 result_3 = (full_df .groupby('cust_id') .rolling(window=3, min_periods=1) .apply(lambda x: pd.Series({ 'max_vol_3': calc_max_3(x), 'sum_trans_3': calc_sum_trans_3(x) }), include_groups=False) .reset_index()) # 5. 计算6个月窗口的指标 result_6 = (full_df .groupby('cust_id') .rolling(window=6, min_periods=1) .apply(calc_max_6, include_groups=False) .rename(columns={'volume': 'max_vol_6'}) .reset_index()) # 6. 合并所有指标结果 combined_result = result_3.merge(result_6[['cust_id', 'period', 'max_vol_6']], on=['cust_id', 'period'], how='left') # 7. 筛选回原数据的行,并还原period为原数字格式 final_out = df.merge(combined_result[['cust_id', 'period', 'max_vol_3', 'max_vol_6', 'sum_trans_3']], on=['cust_id', 'period'], how='left') final_out['period'] = final_out['period'].dt.strftime('%Y%m').astype(int) # 查看最终结果 print(final_out)
代码说明
- 补全缺失月度:确保每个客户的时间序列是连续的月度,避免仅按现有行数量计算滚动窗口,保证“3个月度周期”的逻辑准确。
- 自定义聚合函数:在滚动窗口内统计非NaN的有效数据量,只有满足阈值时才返回聚合结果,否则返回
None。 - 筛选原数据行:最终结果仅保留原数据中存在的月度记录,完全匹配需求输出格式。
内容的提问来源于stack exchange,提问作者Nick
相关产品推荐
相关产品推荐

