第二十七章:高级话题八:做市策略的并行计算与加速

做市策略的回测,说白了就是一场与时间的赛跑。你想想看,一个策略在单核上跑一遍要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。

注意:共享内存里的数据是只读的!千万别在worker里修改共享数据,否则会出现数据竞争。我有个同事就犯过这个错,回测结果每次跑都不一样,排查了整整两天。

用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%的性能问题。

做市策略并行计算架构图 历史Tick数据 数据分片 + 共享内存 Worker 1 参数组 A1-B1 Worker 2 参数组 A2-B2 Worker 3 参数组 A3-B3 结果合并与排序 最优参数组合

这张图展示了我最常用的并行计算流程。数据先分片,然后多个Worker同时处理不同的参数组合,最后合并结果。整个流程就像一条流水线——数据进来,参数出去,中间全是并行。

我个人建议,刚开始做并行加速时,先拿4个参数组合练手。跑通了再扩展到100个、1000个。别一上来就搞分布式,容易把自己绕晕。

好了,并行计算这块就聊到这儿。记住一句话:能用多进程解决的,就别上分布式;能用共享内存解决的,就别搞网络通信。简单,才是最好的。