多周期股票数据集合并内存不足 Dask实现方案咨询
方案结论
你提出的「分频率存储独立Parquet文件、ML流水线启动时通过Dask DataFrame导入完成合并」的思路完全可行,是适配你当前场景的最优解:既不需要用到并行计算能力,也能完美解决全量加载占满内存的问题,同时支持磁盘持久化、按需动态加载。
具体实现步骤
- 预处理存储阶段
先给四个频率的数据集列名加上对应周期前缀(如1分钟频列加1m_、日线加1d_)避免列名冲突,逐个将单数据集写入独立Parquet文件,存储时指定时间戳列为索引,全程不要做全量内存合并。如果单数据集行数较多,存储时可以按年/月做时间分区,后续加载时可以直接裁剪需要的时间范围,加载效率更高。 - Dask懒加载合并阶段
Dask DataFrame的所有操作默认是懒执行模式:读取、合并操作不会立刻把数据载入内存,只有实际触发计算(如取批次、做聚合、写入磁盘)时才会按分块加载对应数据到内存,全程不会出现全量数据占满内存的问题。
最小实现代码如下:
后续对接ML流水线时,直接通过import dask.dataframe as dd # 逐文件懒加载,无实际内存占用 df_1m = dd.read_parquet("./stock_data/1min.parquet", index="trade_time") df_5m = dd.read_parquet("./stock_data/5min.parquet", index="trade_time") df_15m = dd.read_parquet("./stock_data/15min.parquet", index="trade_time") df_1d = dd.read_parquet("./stock_data/1d.parquet", index="trade_time") # 按时间索引做列向合并,仍然是懒执行 merged_ddf = dd.concat([df_1m, df_5m, df_15m, df_1d], axis=1, join="left") # 如果需要持久化合并后的全量数据集到磁盘,直接调用to_parquet即可,分块写入不会占满内存 # merged_ddf.to_parquet("./stock_data/merged_all.parquet")merged_ddf.partitions[i].compute()逐块取出对应分区的Pandas DataFrame喂入模型即可,不需要全量加载数据。
非固定频率floor报错修复方案
1B(工作日)属于非固定频率,Pandas原生的floor()方法仅支持固定时间差的频率,直接替换原有floor逻辑即可,通用兼容写法不需要单独做频率分支判断:
# 原报错语句:data_index = higher_resolution_index.floor(data_freq).drop_duplicates() # 替换为以下写法,固定/非固定频率均兼容,效果和原floor逻辑完全一致 data_index = higher_resolution_index.dt.to_period(data_freq).dt.start_time.drop_duplicates()
如果需要自定义节假日、特殊休市日的工作日规则,只需要在初始化BDay偏移时传入自定义holidays列表,再通过自定义映射函数对齐时间即可,不需要修改其他逻辑。
额外说明:如果你完全不想引入Dask的依赖,也可以给Parquet文件手写分块迭代器,每次按固定行数读取批次数据喂入模型,但Dask已经封装好了索引对齐、分区裁剪、分块读写的完整逻辑,开发成本远低于手写迭代器,不需要额外调整就能适配你的需求。
内容的提问来源于stack exchange,提问作者wildcat89
相关产品推荐
相关产品推荐

