内网分布式计算实现和部署方案
1. 实现目标:
现有内网机器分为两个容器,compute1和compute2,这两个容器目前能够访问一个共享盘根目录,但由于历史遗留问题,compute1目前单独保存一份与compute2高度相似的工作区进行因子计算和回测。这不仅不方便两机器之间的协调,也浪费了大量硬盘资源,因此需要进行分布式计算改造,本文主要探讨Ray和Spark两种分布式系统的特点和适用性,以及给出完整的factor_backtest改造框架。
2. Ray和Spark实现路径探讨:
| 模块 | 当前结构 | Ray 适配 | Spark 适配 | 关键点 |
|---|---|---|---|---|
| 日频 FactorStore 增量更新 | execution graph → batch → subprocess | 很高 | 中低 | 几乎可以直接替换 batch launcher |
| 日频全量因子计算 | module/factor batch + Polars | 很高 | 中低 | 保留现有 Polars 内核 |
| Trade Raw → shared parent | ticker hash partition + merge | 很高 | 高 | 本身已经接近 MapReduce |
| Minute Factor 计算 | 30 日 batch,共享 loader,逐日 wide parquet | 高 | 低 | 分布式单位必须是 date batch |
| Minute Factor Backtest | date-first,跨日 state | 中 | 低 | 日期不能直接并行 |
| 调度/监控 | subprocess + cgroup guard | 很高 | 中 | Ray Dashboard 可直接补足监控 |
由上述表格可见,Ray对于本仓库的适配性更高,并且可以在改动较小的基础上实现完整的分布式计算框架。
- 日频框架:对于日频因子计算和回测,目前采取的是
{factor_logic}.py保存在factor/factors,计算文件持久化在data/factors,每因子每日一个parquet并由manifest.json维护的策略。合理的分布式粒度是对每一个因子类分配不同的node做计算,此结构天生适配目前的方案,事实上现在只是在手工+半自动做Ray的工作。 - 日频增量框架:Compute1机器的日频因子和各类缓存更新已经跑通,Compute2还未跑但预计debug的次数不会很多,目前采用的是计算graph+subprocess的更新模式,且在缓存中维护各种因子可能需要的日级别(甚至分钟级别、Trade Agg级别)的中间结果,目前复用性较差,由于有graph的存在,该计算几乎可直接提交给Ray进行。
- 缓存落盘与增量更新:当前在
/mnt/data/0/zjx/factor_backtest/cache内维护诸多缓存,大小约3~5TB,但增量更新开销不大,只由一台机器完成即可。 - 分钟Level因子计算:目前的计算是Light和Standard计算开销的因子一次读取30天的元数据,接着每天做聚合降频成1m的因子。若要分布式计算,应首先对任务做调度改造,使用不同频率的数据的因子划分为不同等级,然后在每个机器上均匀分配任务,在此基础上最大化共享父表并建立可复用缓存;在落盘时,每个机器计算完成一批因子,生成parquet后先拼接到一个parent parquet,然后在每个日期完成后由一个committer将parquet拼接,此过程应与整个因子计算独立;增量更新层面,暂时只更新过去一个一个交易日及以前的数据,不做分钟频次的更新。
- 分钟Level因子回测:在因子落盘后,按Factor增量平均分给多个机器然后计算落盘,此实现由于已经物化所以比较简单,计算起来也很快。
3. 实现路径与后续扩展资源方式
目前测试后发现compute1和2的tcp互通且都可以访问某一路径,因此在此基础上做Ray是自然且容易实现的,依据上述框架即可改造已有的factor_backtest进行分布式计算(甚至是可以只新增Ray调度脚本控制两个机器)
新增因子文件、因子回测结果和数据缓存由于都与现有结果独立,完全可以做增量分布式计算;而新增机器的话只需维护Ray调度器即可。
- 本文链接: https://otter-garden.cn/posts/2026/08/分布式计算方案实现与探讨-适用于因子计算和回测的架构/
- 版权声明: 本博客所有文章除特别声明外,均默认采用 CC BY-NC 4.0 许可协议。