如何提升多个带datetime索引的DataFrame逐行并行迭代的运行速度?
性能优化方案
原代码核心瓶颈
- 迭代次数过高:迭代次数等于所有DataFrame的总行数,存在大量重复计算
- 冗余操作多:每次循环都需要遍历所有DataFrame、对时间列表全量排序、判断处理状态,无效开销大
- 没有利用pandas向量化特性:逐次loc切片的效率远低于批量拆分操作
优化实现思路
方案一:全量预拆分(推荐,性能提升最明显)
本质上你需要的所有切片分界点,就是19个DataFrame的所有时间索引去重排序后的结果,直接预先生成分界点后批量拆分即可,完全不需要逐轮while循环:
import pandas as pd import datetime import numpy as np from concurrent.futures import ProcessPoolExecutor # 测试数据同原代码 df1 = pd.DataFrame({'A': range(9)}) df1.index = [pd.Timestamp('20130101 09:00:00'), pd.Timestamp('20130101 09:01:00'), pd.Timestamp('20130101 09:30:00'), pd.Timestamp('20130101 09:44:00'), pd.Timestamp('20130101 09:50:00'), pd.Timestamp('20130101 10:16:00'), pd.Timestamp('20130101 10:47:00'), pd.Timestamp('20130101 10:53:00'), pd.Timestamp('20130101 11:22:00')] df2 = pd.DataFrame({'B': range(9)}) df2.index = [pd.Timestamp('20130101 09:00:00'), pd.Timestamp('20130101 09:01:00'), pd.Timestamp('20130101 09:04:00'), pd.Timestamp('20130101 09:05:00'), pd.Timestamp('20130101 09:09:00'), pd.Timestamp('20130101 10:10:00'), pd.Timestamp('20130101 10:15:00'), pd.Timestamp('20130101 10:16:00'), pd.Timestamp('20130101 11:18:00')] db_dict = {"a": df1, "b": df2} start_time = datetime.datetime.now() # 1. 合并所有df的时间索引,生成全局唯一排序分界点 all_timestamps = set() for df in db_dict.values(): all_timestamps.update(df.index.tolist()) # 按需求补初始10秒的分界逻辑 sorted_times = sorted(all_timestamps) bins = [sorted_times[0], sorted_times[0] + datetime.timedelta(seconds=10)] + sorted_times[1:] bins = sorted(list(set(bins))) # 去重避免空区间 # 2. 批量拆分每个df的所有切片 all_slices = [] for name, df in db_dict.items(): # 给每行打区间标签 df['bin_idx'] = pd.cut(df.index, bins=bins, labels=False, include_lowest=True) # 按标签分组得到所有切片 df_groups = df.groupby('bin_idx') for bin_idx, group in df_groups: slice_start = bins[bin_idx] slice_end = bins[bin_idx+1] all_slices.append((name, slice_start, slice_end, group.drop('bin_idx', axis=1))) # 3. 并行处理所有切片,这里的process_func替换为你实际的处理逻辑 def process_func(slice_info): name, start, end, slice_df = slice_info # 你的处理逻辑写在这里 return f"{name} 区间 {start}~{end} 处理完成,共{len(slice_df)}行" with ProcessPoolExecutor() as executor: results = list(executor.map(process_func, all_slices)) # 输出耗时 print(datetime.datetime.now() - start_time)
这个方案的性能提升可以达到几十到上百倍,完全避免了原代码的高频循环开销。
方案二:保留原迭代逻辑的小优化
如果你的业务逻辑必须保留逐轮迭代的写法,可以做以下小改动降低开销:
- 用
set代替list存储complete_list,判断成员存在的时间复杂度从O(n)降到O(1) - 用最小堆维护
time_dict_end的取值,每次取最小时间的复杂度从O(nlogn)(全量排序)降到O(1),插入新值的复杂度为O(logn) - 提前缓存每个df的索引numpy数组,避免每次
searchsorted都转换类型
内容的提问来源于stack exchange,提问作者Max2603
相关产品推荐
相关产品推荐

