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

百万级任务DataFrame处理优化:替代逐行迭代的高效方案

百万行任务数据集高效匹配下一个符合条件任务的方案

逐行apply的时间复杂度为O(n²),对百万级数据完全不可行。以下是两种基于Pandas的高效方案,核心思路是分组+排序+二分查找,将时间复杂度降至O(N log N),大幅提升处理速度。

前置预处理

先确保时间列是datetime类型,并将countries和functions转换为集合(加速in操作):

import pandas as pd
import numpy as np
import bisect

# 转换时间列类型
df['time_closed'] = pd.to_datetime(df['time_closed'])
df['time_assigned'] = pd.to_datetime(df['time_assigned'])

# 转换为集合,提升in操作效率
df['countries_set'] = df['countries'].apply(set)
df['functions_set'] = df['functions'].apply(set)

方案一:分组+二分查找+纯Python循环

该方案无需额外依赖,适配大多数场景:

def process_employee_tasks(group):
    # 按任务分配时间升序排序,保证能找到最小的time_assigned
    sorted_group = group.sort_values('time_assigned').reset_index(drop=True)
    
    # 提取数组化数据,避免循环中重复索引查找
    time_assigned_arr = sorted_group['time_assigned'].values
    country_arr = sorted_group['country_code'].values
    function_arr = sorted_group['function'].values
    time_closed_arr = sorted_group['time_closed'].values
    # 预计算每个任务的10天时间上限
    upper_time_arr = time_closed_arr + np.timedelta64(10, 'D')
    countries_sets = sorted_group['countries_set'].values
    functions_sets = sorted_group['functions_set'].values
    
    next_task_indices = []
    for idx in range(len(sorted_group)):
        lower_time = time_closed_arr[idx]
        upper_time = upper_time_arr[idx]
        
        # 用二分查找快速定位时间范围边界
        start_pos = bisect.bisect_right(time_assigned_arr, lower_time)  # 第一个time_assigned > 当前任务time_closed
        end_pos = bisect.bisect_right(time_assigned_arr, upper_time)    # 最后一个time_assigned <= 当前任务time_closed+10天
        
        # 无符合时间条件的任务
        if start_pos >= end_pos:
            next_task_indices.append(-1)
            continue
        
        # 在时间范围内找第一个符合国家、职能条件的任务(已排序,第一个就是最小time_assigned)
        match_idx = -1
        for j in range(start_pos, end_pos):
            if country_arr[j] in countries_sets[idx] and function_arr[j] in functions_sets[idx]:
                match_idx = j
                break
        next_task_indices.append(match_idx)
    
    # 匹配下一个任务数据,添加前缀区分原任务
    next_tasks = sorted_group.iloc[next_task_indices].reset_index(drop=True).add_prefix('next_')
    result = pd.concat([sorted_group, next_tasks], axis=1)
    
    # 对无匹配的行填充NaN
    no_match_mask = result['next_employee_id'].isna()
    result.loc[no_match_mask, result.columns.str.startswith('next_')] = np.nan
    
    return result

# 按员工分组处理,禁用group_keys避免生成额外列
final_result = df.groupby('employee_id', group_keys=False).apply(process_employee_tasks)

方案二:Numba加速循环(更快)

若要进一步提升速度,用Numba将循环编译为机器码,处理百万行数据可压缩至数分钟内:

from numba import jit

# Numba编译的匹配函数,加速循环逻辑
@jit(nopython=True)
def find_matching_task(start_pos, end_pos, country_arr, function_arr, target_countries, target_functions):
    for j in range(start_pos, end_pos):
        if country_arr[j] in target_countries and function_arr[j] in target_functions:
            return j
    return -1

def process_employee_tasks_numba(group):
    sorted_group = group.sort_values('time_assigned').reset_index(drop=True)
    
    # 转换为Numba支持的数组类型
    time_assigned_arr = sorted_group['time_assigned'].values.astype('datetime64[ns]')
    country_arr = sorted_group['country_code'].values
    function_arr = sorted_group['function'].values
    time_closed_arr = sorted_group['time_closed'].values.astype('datetime64[ns]')
    upper_time_arr = time_closed_arr + np.timedelta64(10, 'D')
    countries_sets = sorted_group['countries_set'].values
    functions_sets = sorted_group['functions_set'].values
    
    next_task_indices = []
    for idx in range(len(sorted_group)):
        lower_time = time_closed_arr[idx]
        upper_time = upper_time_arr[idx]
        
        start_pos = bisect.bisect_right(time_assigned_arr, lower_time)
        end_pos = bisect.bisect_right(time_assigned_arr, upper_time)
        
        if start_pos >= end_pos:
            next_task_indices.append(-1)
            continue
        
        # 调用Numba加速的匹配函数
        match_idx = find_matching_task(start_pos, end_pos, country_arr, function_arr, countries_sets[idx], functions_sets[idx])
        next_task_indices.append(match_idx)
    
    next_tasks = sorted_group.iloc[next_task_indices].reset_index(drop=True).add_prefix('next_')
    result = pd.concat([sorted_group, next_tasks], axis=1)
    
    no_match_mask = result['next_employee_id'].isna()
    result.loc[no_match_mask, result.columns.str.startswith('next_')] = np.nan
    
    return result

final_result_numba = df.groupby('employee_id', group_keys=False).apply(process_employee_tasks_numba)

方案三:超大数据量用Dask/PySpark

若数据量远超内存(如10亿行),可使用Dask(兼容Pandas接口)或PySpark做分布式处理,核心逻辑与上述一致,仅将分组、排序、匹配操作分布式执行。

内容的提问来源于stack exchange,提问作者geoabram

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 06:48:12