pandas如何基于历史值聚合前序数据 实现时序分组均值编码
时间依赖历史均值特征高效实现方案
问题场景
需要在单数据集内按实体分组,基于严格的时间先后顺序计算两类无泄露特征:
- 时间依赖的目标变量均值编码
- 时间依赖的其他特征历史平均值
要求方案可适配内存无法装载的超大规模数据集,支持Dask、Vaex等外存/并行计算框架。
可复现测试基准
测试数据集构造代码
import pandas as pd import numpy as np df = pd.DataFrame({ 'customer_id': list(np.arange(0, 3)) * 5, 'date': pd.date_range('2019-01-01', '2019-01-15'), 'feature_1': np.linspace(0, 100, 15), 'feature_2': ['hi', 'hi', 'hi', 'bye', 'bye'] * 3, 'target': [0, 1, 0, 1, 1] * 3 }) df.head()
数据集前5行输出:
customer_id date feature_1 feature_2 target 0 0 2019-01-01 0.000000 hi 0 1 1 2019-01-02 7.142857 hi 1 2 2 2019-01-03 14.285714 hi 0 3 0 2019-01-04 21.428571 bye 1 4 1 2019-01-05 28.571429 bye 1
已完成前置步骤
按customer_id分组,为每条记录生成同组上一条记录的日期列prev_date:
df['prev_date'] = df.groupby('customer_id')['date'].transform(lambda x: x.shift())
以customer_id=2的样本为例,处理后数据如下:
customer_id date feature_1 feature_2 target prev_date 2 2 2019-01-03 14.285714 hi 0 NaT 5 2 2019-01-06 35.714286 hi 0 2019-01-03 8 2 2019-01-09 57.142857 bye 1 2019-01-06 11 2 2019-01-12 78.571429 hi 1 2019-01-09 14 2 2019-01-15 100.000000 bye 1 2019-01-12
原有错误实现
# 错误写法 df.groupby('customer_id').apply(lambda x: x[x['date'] > x['prev_date']]['target'].mean())
该写法问题:一是逐行过滤的时间复杂度为O(n²),数据量超过10万行就会出现明显卡顿;二是逻辑上仅返回分组整体均值,没有逐行计算截止到当前行之前的历史聚合结果,完全不符合需求。
customer_id=2对应的预期输出参考:
分场景实现方案
核心逻辑统一为:分组内按时间升序排列→做扩展窗口累计聚合→聚合结果整体下移1位排除当前行,避免数据泄露,全程使用向量化运算,无逐行循环。
1. Pandas 中小数据集实现
性能比apply逐行计算高100倍以上,百万级数据可秒级跑完:
# 必须先保证分组内严格按日期升序 df = df.sort_values(['customer_id', 'date']).reset_index(drop=True) # 计算目标变量时间依赖均值编码 df['target_hist_mean'] = df.groupby('customer_id')['target']\ .transform(lambda x: x.expanding().mean().shift()) # 计算数值特征历史均值 df['feature1_hist_mean'] = df.groupby('customer_id')['feature_1']\ .transform(lambda x: x.expanding().mean().shift()) # 计算类别特征历史取值占比(可扩展为多类别均值编码) df['feature2_hi_hist_ratio'] = df.groupby('customer_id')['feature_2']\ .transform(lambda x: (x == 'hi').expanding().mean().shift())
处理后customer_id=2的结果完全匹配预期:第一条无历史数据值为NaN,第二条历史均值为第一条target值0,第三条为前两条target均值0,第四条为前三条均值≈0.333,第五条为前四条均值0.5。
2. Dask 超大数据集实现
API和Pandas高度兼容,支持并行计算、外存调度,不需要把全量数据加载到内存,适合十亿级以内数据集:
import dask.dataframe as dd # 读取分区存储的大数据(支持parquet/csv等格式) ddf = dd.read_parquet('your_large_dataset.parquet') # 按分组键设置索引,保证同组数据落在同一分区,再按日期排序 ddf = ddf.set_index('customer_id').map_partitions( lambda x: x.sort_values('date') ).reset_index() # 历史特征计算逻辑和Pandas完全一致 ddf['target_hist_mean'] = ddf.groupby('customer_id')['target']\ .transform(lambda x: x.expanding().mean().shift(), meta=('target', 'float64')) ddf['feature1_hist_mean'] = ddf.groupby('customer_id')['feature_1']\ .transform(lambda x: x.expanding().mean().shift(), meta=('feature_1', 'float64')) # 触发计算落盘 ddf.to_parquet('dataset_with_hist_features.parquet')
3. Vaex 超大数据集实现
基于内存映射技术,支持百亿级数据秒级聚合,内存占用仅为Pandas的1%:
import vaex # 读取外存数据集(支持hdf5/parquet/arrow等格式) vx_df = vaex.open('your_large_dataset.hdf5') # 按分组和时间排序 vx_df = vx_df.sort(['customer_id', 'date']) # 分组计算累计和、累计计数,移位后计算均值 vx_df['target_cumsum'] = vx_df.groupby('customer_id', sort=True).agg({'target': 'cumsum'}) vx_df['target_cumcnt'] = vx_df.groupby('customer_id', sort=True).agg({'target': 'cumcount'}) vx_df['target_cumsum_prev'] = vx_df.groupby('customer_id')['target_cumsum'].shift(1) vx_df['target_cumcnt_prev'] = vx_df.groupby('customer_id')['target_cumcnt'].shift(1) vx_df['target_hist_mean'] = vx_df['target_cumsum_prev'] / vx_df['target_cumcnt_prev'] # 其他特征计算逻辑完全一致 vx_df['f1_cumsum'] = vx_df.groupby('customer_id', sort=True).agg({'feature_1': 'cumsum'}) vx_df['f1_cumsum_prev'] = vx_df.groupby('customer_id')['f1_cumsum'].shift(1) vx_df['feature1_hist_mean'] = vx_df['f1_cumsum_prev'] / vx_df['target_cumcnt_prev']
注意:所有时间依赖特征计算必须保证分组内严格按时间升序排列,且聚合结果移位1位,绝对不能包含当前行数据,否则会造成模型训练的数据泄露。
内容的提问来源于stack exchange,提问作者Matt Elgazar
相关产品推荐
相关产品推荐

