banner
NEWS LETTER

分布式计算方案实现与探讨-适用于因子计算和回测的架构

Scroll down

内网分布式计算实现和部署方案

1. 实现目标:

现有内网机器分为两个容器,compute1和compute2,这两个容器目前能够访问一个共享盘根目录,但由于历史遗留问题,compute1目前单独保存一份与compute2高度相似的工作区进行因子计算和回测。这不仅不方便两机器之间的协调,也浪费了大量硬盘资源,因此需要进行分布式计算改造,本文主要探讨RaySpark两种分布式系统的特点和适用性,以及给出完整的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调度器即可。

其他文章
目录导航 置顶
  1. 1. 内网分布式计算实现和部署方案
    1. 1.1. 1. 实现目标:
    2. 1.2. 2. Ray和Spark实现路径探讨:
    3. 1.3. 3. 实现路径与后续扩展资源方式
请输入关键词进行搜索