如何高效实现排除指定时段与日期的DataFrame重采样?
高效实现:排除指定时段/日期的Dask OHLC重采样
核心思路
处理百亿级数据的关键是在重采样前完成数据过滤,避免对无效数据执行计算。所有操作采用向量化方式,彻底避开循环和apply,确保性能最优。
步骤拆解与代码实现
1. 定义过滤条件参数
先明确需要排除的日期集合和每日允许的时段范围:
# 示例:需要排除的节假日日期(格式:YYYY-MM-DD) EXCLUDE_DATES = {"2024-01-01", "2024-02-10", "2024-04-04"} # 示例:每日允许交易的时段(如09:00-16:00,反推排除闭市时段) ALLOWED_START = "09:00:00" ALLOWED_END = "16:00:00"
2. 向量化过滤函数
基于Dask的datetime索引特性,实现分区内的快速过滤:
import dask.dataframe as dd from datetime import time def filter_trades(trades_df, exclude_dates, allowed_start, allowed_end): # 转换允许时段为time对象 start_time = time.fromisoformat(allowed_start) end_time = time.fromisoformat(allowed_end) # 过滤条件1:排除指定日期 mask_dates = ~trades_df.index.date.isin(exclude_dates) # 过滤条件2:保留每日允许时段内的数据 mask_time = trades_df.index.time.between(start_time, end_time) # 合并过滤条件,返回过滤后的DataFrame return trades_df[mask_dates & mask_time]
3. 修改原OHLC函数,整合过滤逻辑
把过滤步骤放在重采样之前,大幅减少后续计算量:
def _ohlc( trades_df, price_col_name, exclude_dates, allowed_start, allowed_end ) -> dd.DataFrame: # 第一步:先过滤无效数据(核心优化点) filtered_df = filter_trades(trades_df, exclude_dates, allowed_start, allowed_end) # 提取价格列 prices_df = filtered_df[price_col_name] # 仅对有效数据重采样为1秒OHLC resampled_prices_df = prices_df.resample(f'1S') ohlc = resampled_prices_df.ohlc() # 向前填充有效时段内的缺失值 filled_ohlc = ohlc.ffill() return filled_ohlc
4. 使用示例
# 假设你的原始Dask DataFrame为trades_dd result_ohlc = _ohlc( trades_df=trades_dd, price_col_name="price", exclude_dates=EXCLUDE_DATES, allowed_start=ALLOWED_START, allowed_end=ALLOWED_END ) # 按需触发计算,比如导出到Parquet result_ohlc.to_parquet("filtered_ohlc.parquet")
性能优化说明
- 前置过滤:仅对有效数据执行重采样,直接减少无效计算量
- 向量化操作:所有过滤逻辑基于Dask内置datetime方法,是分区内的向量化运算,无循环、无
apply - 延迟计算:全程保持Dask的延迟执行特性,不会一次性加载全量数据到内存
- 分区友好:过滤操作保留原DataFrame的分区结构,不会打乱分区导致性能下降
内容的提问来源于stack exchange,提问作者newbnoob
相关产品推荐
相关产品推荐

