如何高效关联带时间戳的DataFrame与另一DataFrame的时间区间聚合结果?
高效实现时间区间聚合关联的向量化方案
需求背景
DF1包含数百万行带微秒分辨率的任意时间戳数据;DF2同样包含时间戳(通常与DF1不匹配)及测量值列。需要为DF1中每个时间戳t1,找到DF2中时间戳落在[t1-30分钟, t1]区间内的所有行,对测量值列求和后关联到DF1对应行。
当前逐行循环实现(效率低下):
# Step 1: 构建时间戳到聚合结果的字典 aggcols = [c for c in DF2.columns if c not in ["timestamp"]] # 可聚合的列 def calc_agg(df, ts_to, cols): ts_from = ts_to - pd.Timedelta(minutes=30) intvl_rows = df[(df["timestamp"] <= ts_to) & (df["timestamp"] >= ts_from)] agg = intvl_rows[cols].sum() return agg ts2ser = { ts: calc_agg(DF2, ts, aggcols) for ts in DF1["timestamp"] } # Step 2: 将聚合列添加到DF1 for c in aggcols: DF1[c] = DF1["timestamp"].apply(lambda x: ts2ser[x][c])
高效向量化方案
以下方案均基于pandas/numpy的向量化操作,性能比循环实现提升100倍以上,适合数百万级数据量的场景:
方法1:numpy二分查找+累积和(最快方案)
利用numpy的searchsorted批量定位区间边界,结合累积和计算区间和,底层操作开销极低,适合超大数据量场景。
import numpy as np # 预处理DF2:排序并计算累积和 DF2_sorted = DF2.sort_values("timestamp").reset_index(drop=True) ts2_array = DF2_sorted["timestamp"].values col_cumsums = {col: DF2_sorted[col].cumsum().values for col in aggcols} # 提取DF1时间戳并计算区间起始时间 ts1_array = DF1["timestamp"].values ts1_start = ts1_array - pd.Timedelta(minutes=30).to_numpy() # 批量定位区间左右边界 # right_idx: DF2中<=t1的最后一个元素的下一个索引 right_idx = np.searchsorted(ts2_array, ts1_array, side="right") # left_idx: DF2中>=t1-30min的第一个元素的索引 left_idx = np.searchsorted(ts2_array, ts1_start, side="left") # 计算每个区间的求和结果 for col in aggcols: cumsum_vals = col_cumsums[col] # 右边界累积和:right_idx=0时无数据,和为0 right_sum = np.where(right_idx == 0, 0, cumsum_vals[right_idx - 1]) # 左边界累积和:left_idx=0时无前置数据,和为0 left_sum = np.where(left_idx == 0, 0, cumsum_vals[left_idx - 1]) # 区间和=右累积和-左累积和,区间无数据时和为0 DF1[col] = np.where(right_idx > left_idx, right_sum - left_sum, 0)
方法2:pandas merge_asof+滚动窗口(可读性优先)
基于pandas的merge_asof定位区间边界,结合累积和计算,代码可读性更强,适合需要维护的生产环境。
# 预处理DF2:排序+计算累积和 DF2_sorted = DF2.sort_values("timestamp").reset_index(drop=True) for col in aggcols: DF2_sorted[f"{col}_cumsum"] = DF2_sorted[col].cumsum() # 给DF1添加区间起始时间列 DF1["ts_start"] = DF1["timestamp"] - pd.Timedelta(minutes=30) # 用merge_asof定位每个t1对应的DF2最大匹配行索引 df_t1_match = pd.merge_asof( DF1[["timestamp", "ts_start"]], DF2_sorted[["timestamp"]].assign(idx=DF2_sorted.index), left_on="timestamp", right_on="timestamp", direction="backward" ).rename(columns={"idx": "idx_t1"}) # 定位每个ts_start对应的DF2最小匹配行索引 df_start_match = pd.merge_asof( DF1[["ts_start"]], DF2_sorted[["timestamp"]].assign(idx=DF2_sorted.index), left_on="ts_start", right_on="timestamp", direction="forward" ).rename(columns={"idx": "idx_start"}) # 合并索引并处理边界情况 DF1_merged = df_t1_match.join(df_start_match["idx_start"]) DF1_merged["idx_start"] = DF1_merged["idx_start"].fillna(0).astype(int) DF1_merged["idx_t1"] = DF1_merged["idx_t1"].fillna(-1).astype(int) # 计算区间求和 for col in aggcols: cumsum_col = f"{col}_cumsum" # 右边界累积和 right_sum = DF2_sorted[cumsum_col].iloc[DF1_merged["idx_t1"]].values # 左边界累积和,idx_start=0时取0 left_sum = DF2_sorted[cumsum_col].iloc[DF1_merged["idx_start"] - 1].values left_sum[DF1_merged["idx_start"] == 0] = 0 # 区间无数据时和为0 right_sum[DF1_merged["idx_t1"] < DF1_merged["idx_start"]] = 0 DF1[col] = right_sum - left_sum
方法3:区间索引关联(简洁但内存占用高)
利用IntervalIndex创建DF1的时间区间,通过分组聚合实现关联,代码最简洁但中间表可能占用大量内存,适合数据量较小的场景。
# 创建DF1的时间区间索引 intervals = pd.IntervalIndex.from_arrays( DF1["timestamp"] - pd.Timedelta(minutes=30), DF1["timestamp"], closed="both" ) DF1_interval = DF1.set_index(intervals) # 给DF2的每个时间戳匹配所属的DF1区间 DF2["match_interval"] = pd.cut(DF2["timestamp"], bins=intervals) # 分组求和后关联回DF1 agg_result = DF2.groupby("match_interval")[aggcols].sum() DF1 = DF1.join(agg_result, on=intervals) # 填充无匹配数据的NaN为0 DF1[aggcols] = DF1[aggcols].fillna(0)
方案选型建议
- 超大数据量(千万级):优先选方法1,速度最快,内存开销低
- 生产环境需维护:优先选方法2,代码可读性强,逻辑清晰
- 数据量较小(百万级以内):可选方法3,代码简洁易实现
内容的提问来源于stack exchange,提问作者jpp1
相关产品推荐
相关产品推荐

