如何用向量化方法查找两个DataFrame的区间重叠并计算加权平均?
问题与解决方案
问题描述
现有两个DataFrame:df1(万级规模):
id start end 1 0 15 1 15 30 1 30 45 2 0 15 2 15 30 2 30 45
df2(百万级规模):
id start end 1 0 1.1 1 1.1 11.4 1 11.4 34 1 34 46 2 0 1.5 2 1.5 20 2 20 30
需求:
- 对
df1的每一行,找到df2中满足同id且区间重叠的行,重叠判定规则:df1.start <= df2.end AND df1.end > df2.start,预期输出为长度等于len(df1)的结构,每行存储符合条件的df2行索引(示例预期结果:[0,1,2],[2],[2,3],[4,5],[5,6],[6])。 - 最终需基于重叠量占比作为权重,计算
df2中quantity列的加权平均值,需兼顾内存效率,避免全量笛卡尔积导致的内存溢出。
向量化实现方案(低内存友好)
核心思路是按id分组处理,避免跨id的无效计算,同时在组内利用numpy广播实现向量化匹配,既保留向量化的高效性,又控制内存占用。
步骤1:数据预处理
给df2添加原始索引列,方便后续记录匹配到的行索引:
import pandas as pd import numpy as np # 构造示例数据(含quantity列) df1 = pd.DataFrame({ 'id': [1,1,1,2,2,2], 'start': [0,15,30,0,15,30], 'end': [15,30,45,15,30,45] }) df2 = pd.DataFrame({ 'id': [1,1,1,1,2,2,2], 'start': [0,1.1,11.4,34,0,1.5,20], 'end': [1.1,11.4,34,46,1.5,20,30], 'quantity': [10,20,30,40,50,60,70] }) df2['idx'] = df2.index # 保存df2原始行索引
步骤2:定义组内处理函数
函数完成两个核心任务:匹配重叠行的索引、计算加权平均值,全程用向量化操作:
def process_id_group(group_data): df1_group = group_data['df1'] df2_group = group_data['df2'] # 1. 向量化生成重叠条件矩阵(shape: [len(df1_group), len(df2_group)]) start_mask = df1_group['start'].values[:, np.newaxis] <= df2_group['end'].values end_mask = df1_group['end'].values[:, np.newaxis] > df2_group['start'].values overlap_mask = start_mask & end_mask # 提取每个df1行对应的df2索引列表 df1_group['matched_df2_indices'] = [ df2_group['idx'].values[mask].tolist() for mask in overlap_mask ] # 2. 计算加权平均(基于重叠长度占比) # 计算df1每行的区间总长度 df1_interval_len = df1_group['end'] - df1_group['start'] # 向量化计算两两区间的重叠长度 overlap_start = np.maximum(df1_group['start'].values[:, np.newaxis], df2_group['start'].values) overlap_end = np.minimum(df1_group['end'].values[:, np.newaxis], df2_group['end'].values) overlap_len = np.maximum(0, overlap_end - overlap_start) # 计算权重(重叠长度/df1该行区间长度) weights = overlap_len / df1_interval_len.values[:, np.newaxis] # 计算加权后的quantity值 weighted_qty = weights * df2_group['quantity'].values # 求和得到加权平均(处理无重叠的情况) total_weight = weights.sum(axis=1) df1_group['weighted_avg_quantity'] = np.where( total_weight > 0, weighted_qty.sum(axis=1) / total_weight, np.nan ) return df1_group
步骤3:分组执行处理
按id分组,将同id的df1和df2数据传入处理函数:
# 合并df1和df2并按id分组 combined = pd.concat([ df1.assign(data_type='df1'), df2.assign(data_type='df2') ]).groupby('id') # 对每个id组执行处理 final_result = combined.apply(lambda g: process_id_group({ 'df1': g[g['data_type'] == 'df1'].drop('data_type', axis=1), 'df2': g[g['data_type'] == 'df2'].drop('data_type', axis=1) })).reset_index(drop=True)
结果展示
final_result包含两列关键数据:
matched_df2_indices:对应预期的df2索引列表weighted_avg_quantity:计算好的加权平均值
示例输出片段:
id start end matched_df2_indices weighted_avg_quantity 0 1 0 15 [0, 1, 2] 25.533333 1 1 15 30 [2] 30.000000 2 1 30 45 [2, 3] 32.608696 ...
内存优化说明
- 分组规避全量笛卡尔积:百万级df2按id拆分后,每个组的规模大幅缩小,避免了万×百万级的超大矩阵(10^10量级),内存占用可控。
- 按需取舍索引存储:如果不需要保存
matched_df2_indices,可直接删除相关代码,进一步节省内存(仅保留加权平均计算逻辑即可)。 - 大组拆分处理:若存在单个id占比极高的情况(比如某id对应10万行df2数据),可将该组内的df1进一步拆分批次处理,降低单批次矩阵的规模。
内容的提问来源于stack exchange,提问作者24n8
相关产品推荐
相关产品推荐

