You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

寻找高效实现两个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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.13 22:45:37