如何基于组内时间差条件对Pandas列进行高效前向填充
带时间限制的分组前向填充(向量化高效实现)
问题背景
我有一个按PATIENT_ID、ENCNTR_ID分组的DataFrame,包含Event_timestamp列和多列需要前向填充的生命体征数据。要求仅当组内当前行的时间戳与最后一个原始有效值(非填充值)的时间差小于设定阈值(如60分钟)时,才对当前行的空值进行前向填充。
当前使用自定义apply逐行循环的方式实现,但面对1.23亿行数据、20个待填充列时,处理速度极慢,急需向量化方法提升效率。
当前低效实现代码
自定义分组填充函数
def ffill_across_episodes(group, targetCol,maxMinsDiff): curPat = 0 curEnc = 0 lastValidEntry = np.nan hasValidEntry = False validTimestamp = None for index, row in group.iterrows(): #print(f"Target col timestamp = {row['Event_timestamp']}, value = {row[targetCol]}") if curPat!=row['PATIENT_ID'] or curEnc != row['ENCNTR_ID']: # 切换患者/就诊记录,重置状态 #print(f"change patient or encounter") curPat = row['PATIENT_ID'] curEnc = row['ENCNTR_ID'] if np.isnan(row[targetCol]): hasValidEntry=False #print(f"set NON valid prior entry") else: #print(f"set valid prior entry") hasValidEntry=True lastValidEntry=row[targetCol] validTimestamp = row['Event_timestamp'] else: # 同一患者同一就诊记录 if np.isnan(row[targetCol]): # 当前值为空,尝试填充 if hasValidEntry: #print(f"has valid prior entry. Timediff is: {( row['Event_timestamp']-validTimestamp).total_seconds()/60:.2f} mins") if ( row['Event_timestamp']-validTimestamp).total_seconds()/60 < maxMinsDiff: group.at[index, targetCol] = lastValidEntry #print(f"set index={index} entry to {lastValidEntry}") else: # 当前值为有效值,更新状态 #print(f"set valid prior entry") hasValidEntry=True lastValidEntry=row[targetCol] validTimestamp = row['Event_timestamp'] return group
填充执行逻辑
for aVar in VITALS_COLS: if aVar in df.columns: print(f"Fwd filling {aVar} column by patient, encounter and across episode boundaries") maxMinsDiff = MAX_HR_INTERVAL_BTW_VITALS*60 df_updated = df.groupby(GROUP_BY_ENCOUNTER_COLS).apply(ffill_across_episodes,aVar,maxMinsDiff) else: print(f"error: {aVar} is not a column within the dataframe") df = df_updated df_updated=None
向量化解决方案
核心思路
完全避免逐行循环,利用Pandas内置的分组、向量化函数实现:
- 为每个分组标记原始有效值的位置,向前传播有效值的数值和对应的时间戳
- 计算每行与最近原始有效值的时间差,仅对时间差小于阈值的空值进行填充
- 批量处理所有待填充列,减少重复分组开销
高效实现代码
import pandas as pd import numpy as np # 预定义参数(需与原代码保持一致) GROUP_BY_ENCOUNTER_COLS = ['PATIENT_ID', 'ENCNTR_ID'] max_mins_diff = MAX_HR_INTERVAL_BTW_VITALS * 60 # 批量处理所有待填充列 for col in VITALS_COLS: if col not in df.columns: print(f"error: {col} is not a column within the dataframe") continue # 1. 标记原始有效值,传播最近有效值的数值和时间戳 df[f'{col}_is_valid'] = df[col].notna() df[f'{col}_last_valid_val'] = df.groupby(GROUP_BY_ENCOUNTER_COLS)[col].ffill() df[f'{col}_last_valid_ts'] = df.groupby(GROUP_BY_ENCOUNTER_COLS)['Event_timestamp'].transform( lambda x: x.where(df[f'{col}_is_valid']).ffill() ) # 2. 计算时间差(分钟),生成填充掩码 time_diff = (df['Event_timestamp'] - df[f'{col}_last_valid_ts']).dt.total_seconds() / 60 fill_mask = df[col].isna() & (time_diff < max_mins_diff) # 3. 执行填充并清理临时列 df.loc[fill_mask, col] = df.loc[fill_mask, f'{col}_last_valid_val'] df.drop(columns=[f'{col}_is_valid', f'{col}_last_valid_val', f'{col}_last_valid_ts'], inplace=True) print(f"Completed forward fill for {col}")
性能优化说明
- 完全向量化:所有操作均使用Pandas内置的向量化函数,相比逐行循环,速度可提升10~100倍
- 减少分组次数:每个列仅需一次分组操作,避免原代码中多次
groupby.apply的重复开销 - 内存可控:临时列随用随删,避免内存占用过高
额外性能建议
- 确保
Event_timestamp为datetime64类型:df['Event_timestamp'] = pd.to_datetime(df['Event_timestamp']) - 为分组列设置索引并排序:
df = df.set_index(GROUP_BY_ENCOUNTER_COLS).sort_index(),可大幅加速分组操作 - 若内存不足,可使用Dask或PySpark进行分布式处理,适配超大规模数据集
内容的提问来源于stack exchange,提问作者WaterBoy
相关产品推荐
相关产品推荐

