第二十七章:高级话题八:做市策略的并行计算与加速
做市策略的回测,说白了就是一场与时间的赛跑。你想想看,一个策略在单核上跑一遍要3小时,参数优化要跑1000次——那就是3000小时,125天。等你调完参数,市场风格早就变了。
我最早做高频做市回测时,就吃过这个亏。一个简单的价差策略,参数网格稍微密一点,跑完要整整一个周末。后来我痛定思痛,开始研究并行计算。今天就把这些经验掰开揉碎讲给你听。
为什么做市策略特别需要并行加速?
做市策略的回测有几个特点,天然适合并行:
- 参数空间巨大:价差阈值、库存上限、撤单时间、报价偏移量……随便组合就是几千种
- 独立性强:不同参数组合的回测互不依赖,可以同时跑
- 计算密集:每笔订单的撮合、库存计算、PnL统计,CPU消耗不小
我个人的习惯是,先把策略拆成「可并行」和「不可并行」两部分。可并行的部分,比如参数扫描、蒙特卡洛模拟,直接扔给多核处理。不可并行的部分,比如依赖前序状态的时序逻辑,就老老实实单线程。
并行计算的三种主流方案
做市策略加速,我试过三种方案。这里直接上对比:
| 方案 | 适用场景 | 学习成本 | 加速比 | 我的评价 |
|---|---|---|---|---|
| 多进程(multiprocessing) | CPU密集型,参数扫描 | 低 | 接近线性 | 最推荐,简单粗暴 |
| 多线程(threading) | I/O密集型,数据加载 | 低 | 受GIL限制 | 做市回测不推荐 |
| 分布式(Ray/Dask) | 超大规模,集群计算 | 高 | 可扩展 | 团队协作时用 |
嗯,这里要注意。Python的多线程因为全局解释器锁(GIL),做CPU密集型的回测加速效果很差。我曾经试过用多线程跑参数优化,8核机器只跑出了1.2倍的加速——还不如不开。
实战:用多进程加速参数扫描
先看一个最典型的场景:我们要对做市策略的「价差阈值」和「库存上限」两个参数做网格搜索。
这是单进程版本:
def backtest_single(spread_threshold, inventory_limit):
# 模拟回测逻辑
pnl = run_market_making(spread_threshold, inventory_limit)
return pnl
# 参数网格
params = [(0.001, 100), (0.001, 200), (0.002, 100), (0.002, 200)]
results = []
for spread, inv in params:
pnl = backtest_single(spread, inv)
results.append((spread, inv, pnl))
这段代码跑起来,CPU利用率只有25%(4核机器)。剩下的75%在摸鱼。改成多进程:
from multiprocessing import Pool
def backtest_wrapper(args):
spread, inv = args
return backtest_single(spread, inv)
if __name__ == '__main__':
params = [(0.001, 100), (0.001, 200), (0.002, 100), (0.002, 200)]
with Pool(processes=4) as pool:
results = pool.map(backtest_wrapper, params)
# results 顺序与 params 一致
for i, (spread, inv) in enumerate(params):
print(f"参数: {spread}, {inv} -> PnL: {results[i]}")
就这么简单,加速比直接拉到3.8倍。为什么不是4倍?进程间通信和内存复制有开销,这是正常的。
Pool(processes=cpu_count()-1) 留一个核给系统,防止机器卡死。我曾经贪心把8核全占满,结果回测跑到一半系统无响应,数据都没保存——血的教训。
进阶:共享内存与数据分片
当参数组合数量达到上万级别时,每个进程都复制一份完整的历史数据,内存会爆炸。我遇到过最夸张的一次,32GB内存直接OOM(内存溢出)。
解决方案是「数据分片 + 共享内存」:
from multiprocessing import shared_memory
import numpy as np
# 创建共享内存
data = np.array(historical_ticks) # 假设这是你的历史tick数据
shm = shared_memory.SharedMemory(create=True, size=data.nbytes)
shared_data = np.ndarray(data.shape, dtype=data.dtype, buffer=shm.buf)
shared_data[:] = data[:]
def worker(param, shm_name, data_shape, data_dtype):
# 每个worker attach到共享内存
existing_shm = shared_memory.SharedMemory(name=shm_name)
shared_data = np.ndarray(data_shape, dtype=data_dtype, buffer=existing_shm.buf)
# 只处理自己负责的时间段
start_idx, end_idx = get_chunk(param, len(shared_data))
chunk = shared_data[start_idx:end_idx]
result = run_backtest_on_chunk(chunk, param)
existing_shm.close()
return result
这样做的好处是:无论开多少个进程,历史数据只在内存中存一份。我实测过,100个进程同时回测,内存占用只比单进程多了不到200MB。
用Ray做分布式加速
当单机多核不够用时,就要上分布式了。Ray是我用得最多的框架,原因很简单——API设计得跟写普通Python函数一样。
import ray
ray.init(address='auto') # 连接集群
@ray.remote
def distributed_backtest(spread, inventory_limit):
# 这里的代码和单机版一模一样
pnl = run_market_making(spread, inventory_limit)
return {'spread': spread, 'inventory': inventory_limit, 'pnl': pnl}
# 提交任务
futures = [distributed_backtest.remote(s, i) for s, i in params]
# 获取结果(自动负载均衡)
results = ray.get(futures)
Ray会自动把任务分发到集群中的各个节点。我曾在20台机器、160核的集群上跑过做市策略的参数优化,原本需要跑3天的任务,2小时就搞定了。
不过说实话,大部分个人做市商用不到分布式。我建议你先用多进程把单机性能榨干,实在不够再上Ray。
并行计算的避坑指南
这些年踩过的坑,我总结成几条:
- 随机数种子要独立:每个worker必须设置不同的seed,否则所有worker产生相同的结果,并行等于白跑
- 日志要加进程ID:不然你根本分不清哪条日志是哪个worker打的
- 结果合并要小心:不同worker返回的结果顺序可能乱,建议用字典或DataFrame带索引
- 避免频繁I/O:每个worker都写文件会导致磁盘争抢,改成攒一批再写
核心思路总结:
做市策略的并行加速,本质上是「把独立的任务拆开,同时执行」。多进程适合单机参数扫描,共享内存解决大数据量问题,Ray搞定跨机器分布式。别追求花哨的技术,先把多进程用熟,就能解决80%的性能问题。
这张图展示了我最常用的并行计算流程。数据先分片,然后多个Worker同时处理不同的参数组合,最后合并结果。整个流程就像一条流水线——数据进来,参数出去,中间全是并行。
我个人建议,刚开始做并行加速时,先拿4个参数组合练手。跑通了再扩展到100个、1000个。别一上来就搞分布式,容易把自己绕晕。
好了,并行计算这块就聊到这儿。记住一句话:能用多进程解决的,就别上分布式;能用共享内存解决的,就别搞网络通信。简单,才是最好的。