大规模时间序列重采样性能优化求助:Pandas提速方案咨询
大规模时间序列重采样性能优化问题
问题背景
我一直在用Pandas做时间序列重采样,结果符合预期但性能跟不上需求。手里是分钟级数据,需要转成5min, 15min, 3H等多种频率。小数据集下Pandas重采样没问题,但处理1000个标的10天数据时性能骤降。
已尝试方案
- 基于
numpy数组实现重采样,性能反而更差(可能实现方式不对) - Pandas中用
resample+apply结合Cython自定义聚合函数,性能不如原生resample+agg - 用
numba.jit优化,没提升 - 拆分数据用多进程重采样,性能提升约50%但开销和计算需求大幅增加
观察结果
- 移除Pandas聚合里的
sum操作,性能小幅提升15%左右 - 自己写的numpy实现远慢于预期
测试代码及输出
测试代码
from datetime import datetime from time import time import numpy as np import pandas as pd symbols = 1000 start = datetime(2023, 1, 1) end = datetime(2023, 1, 2) data_cols = ['open', 'high', 'low', 'close', 'volume'] agg_props = {'open': 'first', 'high': 'max', 'low': 'min', 'close': 'last', 'volume': 'sum'} base, sample = '1min', '5min' def pandas_resample(df: pd.DataFrame): df.sort_values(['sid', 'timestamp'], inplace=True) df = df.set_index('timestamp') re_df = df.groupby('sid').resample(sample, label='left', closed='left').agg(agg_props).reset_index() return re_df def numpy_resample(arr): intervals = pd.date_range(arr[:, 1].min(), arr[:, 1].max(), freq=sample) intervals = list(zip(intervals[:-1], intervals[1:])) # chunk_dates(data_df.index.min(), data_df.index.max(), interval=self.freq, as_range=True) data = [] groups = np.unique(arr[:, 0]) for _group in groups: group_data = arr[arr[:, 0] == _group, :] for _start, _end in intervals: # print(_start) _data_filter = (_start <= group_data[:, 1]) & (group_data[:, 1] < _end) _interval_data = group_data[_data_filter] _interval_agg = [_group, _start] _skip = len(_interval_data) == 0 for _val, _key in [['open', 2], ['high', 3], ['low', 4], ['close', 5], ['volume', 6]]: # print(_key) _col_data = _interval_data[:, _key] if not _skip: if _val in ['open']: _key_val = _col_data[0] if _val in ['high']: _key_val = _col_data.max() if _val in ['low']: _key_val = _col_data.min() if _val in ['close']: _key_val = _col_data[-1] if _val in ['volume']: _key_val = _col_data.sum() else: _key_val = None _interval_agg.append(_key_val) data.append(_interval_agg) return data if __name__ == '__main__': timestamps = pd.date_range(start, end, freq=base) candles = pd.DataFrame({'timestamp': pd.DatetimeIndex(timestamps), **{_col: np.random.randint(50, 150, len(timestamps)) for _col in data_cols}}) symbol_id = pd.DataFrame({'sid': np.random.randint(1000, 2000, symbols)}) candles['id'] = 1 symbol_id['id'] = 1 data_df = symbol_id.merge(candles, on='id').drop(columns=['id']) print(len(data_df), "\n", data_df.head(3)) st1 = time() resampled_df = pandas_resample(data_df.copy()) print('pandas', time() - st1) st2 = time() numpy_resample(data_df.values) print('numpy', time() - st2)
输出结果
pandas 3.5319528579711914 numpy 93.10612797737122
优化方案及思路
1. 优化Pandas原生重采样细节
- 提前排序避免重复操作:把
sort_values移到数据预处理阶段,不要在每次重采样时重复执行,减少冗余开销。 - 替换
groupby.resample为groupby+Grouper:直接用时间分组器作为groupby参数,性能更优:def optimized_pandas_resample(df: pd.DataFrame): # 假设df已按sid和timestamp预处理排序 re_df = df.groupby( ['sid', pd.Grouper(key='timestamp', freq=sample, label='left', closed='left')] ).agg(agg_props).reset_index() return re_df - 压缩数据类型:将
sid转为整数类型,数值列用int32/float32替代默认的int64/float64,减少内存占用和IO开销。
2. 用Dask实现并行重采样
Dask兼容Pandas API,自动分块并行处理大规模数据,开销远低于手动多进程:
import dask.dataframe as dd def dask_resample(df: pd.DataFrame): ddf = dd.from_pandas(df, npartitions=8) # 根据CPU核心数调整分区数 ddf = ddf.groupby('sid').resample( sample, label='left', closed='left', on='timestamp' ).agg(agg_props).reset_index() return ddf.compute()
3. 优化NumPy实现(修复低效问题)
原NumPy实现嵌套循环过多,可通过向量化计算区间ID减少循环开销:
def optimized_numpy_resample(arr): # arr列顺序:sid, timestamp, open, high, low, close, volume sid_col = arr[:, 0].astype(np.int64) ts_col = arr[:, 1].astype('datetime64[ns]') # 计算每个时间戳所属的重采样区间起始时间 epoch = np.datetime64('1970-01-01') ts_epoch = (ts_col - epoch).astype(np.int64) interval_len = 5 * 60 * 10**9 # 5分钟对应的纳秒数 interval_ids = ts_epoch // interval_len interval_starts = interval_ids * interval_len + epoch result = [] for sid in np.unique(sid_col): mask = sid_col == sid sid_interval_ids = interval_ids[mask] sorted_idx = np.argsort(sid_interval_ids) # 按区间分组拆分数据 splits = np.split( np.arange(len(sorted_idx)), np.where(np.diff(sid_interval_ids[sorted_idx]))[0] + 1 ) for split in splits: interval_start = interval_starts[sid_interval_ids[sorted_idx[split[0]]]] # 向量化聚合 agg_open = arr[mask, 2][sorted_idx][split[0]] agg_high = arr[mask, 3][sorted_idx][split].max() agg_low = arr[mask, 4][sorted_idx][split].min() agg_close = arr[mask, 5][sorted_idx][split[-1]] agg_volume = arr[mask, 6][sorted_idx][split].sum() result.append([sid, interval_start, agg_open, agg_high, agg_low, agg_close, agg_volume]) return result
4. 用Vaex处理超大规模数据
Vaex支持内存映射,无需加载全量数据到内存,同时支持向量化并行计算:
import vaex def vaex_resample(df: pd.DataFrame): vdf = vaex.from_pandas(df) vdf['interval_start'] = vdf.timestamp.dt.floor(sample) agg_df = vdf.groupby(['sid', 'interval_start']).agg({ 'open': 'first', 'high': 'max', 'low': 'min', 'close': 'last', 'volume': 'sum' }) return agg_df.to_pandas_df()
内容的提问来源于stack exchange,提问作者ap14
相关产品推荐
相关产品推荐

