寻找高效实现两个DataFrame逐行配对计算的优化方案
优化DataFrame逐行配对计算的高效方案
问题背景
需要处理两个大规模DataFrame的逐行配对计算,这类操作本质为O(n²)复杂度,现有实现方法速度无法满足需求,需寻找更高效的替代方案。
已尝试的低效方法
itertools.product遍历行对
扩展性差,逐行处理开销大:for (_,x), (_,y) in product(X, Y): func(x, y)双重apply嵌套调用
速度最慢,pandas的apply本身开销高,嵌套后放大性能问题:df1.apply(lambda x: df2.apply(lambda y: func(x, y), axis=1), axis=1)cross merge后apply处理
虽避免了嵌套,但apply逐行处理仍未摆脱性能瓶颈:def apply_func(row, df_dict): start = row["ts_min_x"] - row["ts_min_y"] end = row["ts_max_x"] - row["ts_max_y"] start_ts = max([row["ts_min_x"], row["ts_min_y"]]) end_ts = min([row["ts_max_x"], row["ts_max_y"]]) duration = (end_ts - start_ts).total_seconds() if row["id_x"] not in df_dict: df_dict[row["id_x"]] = {} if row["id_y"] not in df_dict.get(row["id_x"]): df_dict[row["id_x"]][row["id_y"]] = [] df_dict[row["id_x"]][row["id_y"]].append((start, end, start_ts, duration)) merged = df1.merge(df2, how="cross") merged.apply(lambda x: apply_func(x, df_dict), axis=1)
高效实现方案
1. 向量化操作(针对数值类计算)
利用numpy广播机制,直接对整列进行矩阵运算,避免逐行循环:
# 以value列求和场景为例 import numpy as np import pandas as pd df1 = pd.DataFrame([ {"value": 1, "timestamp": "1970-01-01T00:00:01z"}, {"value": 2, "timestamp": "1970-01-01T00:00:02z"} ]) df2 = pd.DataFrame([ {"value": 3, "timestamp": "1970-01-01T00:00:03z"}, {"value": 4, "timestamp": "1970-01-01T00:00:04z"} ]) # 转成numpy数组并调整形状实现广播 values1 = df1['value'].to_numpy().reshape(-1, 1) values2 = df2['value'].to_numpy().reshape(1, -1) # 直接矩阵相加后展平 results = (values1 + values2).flatten().tolist() print(results) # [4, 5, 5, 6]
2. Numba并行加速(针对自定义复杂逻辑)
将数据转成numpy数组,用Numba的JIT编译和并行循环替代pandas的apply,大幅降低Python循环开销:
import numba import pandas as pd import numpy as np df1 = pd.DataFrame([ {"id": 1, "active": True, "ts_min": pd.Timestamp(1), "ts_max": pd.Timestamp(10)}, {"id": 2, "active": False, "ts_min": pd.Timestamp(10), "ts_max": pd.Timestamp(12)}, {"id": 3, "active": True, "ts_min": pd.Timestamp(12), "ts_max": pd.Timestamp(19)}, ]) df2 = pd.DataFrame([ {"id": 4, "active": True, "ts_min": pd.Timestamp(4), "ts_max": pd.Timestamp(9)}, {"id": 5, "active": True, "ts_min": pd.Timestamp(9), "ts_max": pd.Timestamp(16)}, {"id": 6, "active": False, "ts_min": pd.Timestamp(16), "ts_max": pd.Timestamp(19)}, ]) # 提取核心数据转成numpy数组 df1_arr = df1[['id', 'ts_min', 'ts_max']].to_numpy() df2_arr = df2[['id', 'ts_min', 'ts_max']].to_numpy() @numba.njit(parallel=True) def compute_pairs(df1_arr, df2_arr): n1, n2 = df1_arr.shape[0], df2_arr.shape[0] results = [] # 并行遍历df1的行 for i in numba.prange(n1): id_x, ts_min_x, ts_max_x = df1_arr[i] for j in range(n2): id_y, ts_min_y, ts_max_y = df2_arr[j] start = ts_min_x - ts_min_y end = ts_max_x - ts_max_y start_ts = max(ts_min_x, ts_min_y) end_ts = min(ts_max_x, ts_max_y) duration = (end_ts - start_ts).total_seconds() if end_ts >= start_ts else 0 results.append((id_x, id_y, start, end, start_ts, duration)) return results # 计算并整理成目标字典结构 raw_results = compute_pairs(df1_arr, df2_arr) df_dict = {} for id_x, id_y, start, end, start_ts, duration in raw_results: if id_x not in df_dict: df_dict[id_x] = {} if id_y not in df_dict[id_x]: df_dict[id_x][id_y] = [] df_dict[id_x][id_y].append((start, end, start_ts, duration))
3. 提前过滤减少计算基数
如果存在筛选条件(如active=True),先过滤掉不需要的行,直接减少配对数量:
# 仅保留active为True的行,降低n²的计算规模 df1_filtered = df1[df1['active']] df2_filtered = df2[df2['active']] # 再对过滤后的DataFrame执行后续计算 merged = df1_filtered.merge(df2_filtered, how="cross")
内容的提问来源于stack exchange,提问作者Martin Lange
相关产品推荐
相关产品推荐

