配置并行化
PyBroker 使用 Joblib 并行计算指标、训练模型以及优化参数。
设置工作进程
set_parallel 用于更新全局 Joblib 配置。n_jobs 参数指定工作任务的数量:-1 (默认值)表示使用所有可用的 CPU 核心,1 表示按顺序运行。可以使用 get_parallel_config 读取当前设置:
[1]:
from pybroker import set_parallel, get_parallel_config
# Use a fixed number of workers.
set_parallel(n_jobs=4)
print(get_parallel_config())
# Or disable parallel execution entirely.
set_parallel(n_jobs=1)
print(get_parallel_config())
ParallelConfig(n_jobs=4, backend='loky', parallel=None)
ParallelConfig(n_jobs=1, backend='loky', parallel=None)
并行计算指标
指标是按**每个股票代码**计算的:给定股票代码的所有指标都会被归入一个任务,每个股票代码分派一个任务。向 backtest、walkforward 或 optimize 传入 parallel_indicators=True,会让这些任务在已配置的工作进程之间运行(默认值为 False)。
为了实际演示这一点,让我们回测一个移动平均线交叉策略:
[2]:
import numpy as np
import pybroker
from pybroker import Strategy, YFinance, sumv
pybroker.enable_data_source_cache("parallelization")
set_parallel(n_jobs=-1)
def sma(bar_data, period):
return sumv(bar_data.close, period) / period
sma_20 = pybroker.indicator("sma_20", sma, period=20)
def sma_cross(ctx):
sma_vals = ctx.indicator("sma_20")
if np.isnan(sma_vals[-1]):
return
pos = ctx.long_pos()
if not pos and ctx.close[-1] > sma_vals[-1]:
ctx.buy_shares = 100
elif pos and ctx.close[-1] < sma_vals[-1]:
ctx.sell_all_shares()
yfinance = YFinance()
strategy = Strategy(yfinance, start_date="1/1/2021", end_date="1/1/2026")
strategy.add_execution(sma_cross, ["V", "MA", "AXP"], indicators=sma_20)
result = strategy.backtest(parallel_indicators=True, warmup=20)
result.metrics_df.head()
Backtesting: 2021-01-01 00:00:00 to 2026-01-01 00:00:00
Loading bar data...
[*********************100%***********************] 3 of 3 completed
Loaded bar data: 0:00:00
Computing indicators...
100% (3 of 3) |##########################| Elapsed Time: 0:00:01 Time: 0:00:01
Test split: 2021-01-04 00:00:00 to 2025-12-31 00:00:00
100% (1255 of 1255) |####################| Elapsed Time: 0:00:00 Time: 0:00:00
Finished backtest: 0:00:02
[2]:
| name | value | |
|---|---|---|
| 0 | trade_count | 218 |
| 1 | initial_market_value | 100000.0 |
| 2 | end_market_value | 116119.9 |
| 3 | total_pnl | 15837.68 |
| 4 | unrealized_pnl | 282.22 |
使用 IndicatorSet 进行独立的指标计算时,同样接受相同的 parallel_indicators 标志:
[3]:
from pybroker import IndicatorSet
df = yfinance.query(
["V", "MA", "AXP"], start_date="1/1/2021", end_date="1/1/2026"
)
ind_set = IndicatorSet()
ind_set.add(sma_20)
ind_set(df, parallel_indicators=True).tail()
Loaded cached bar data.
Computing indicators...
100% (3 of 3) |##########################| Elapsed Time: 0:00:01 Time: 0:00:01
[3]:
| symbol | date | sma_20 | |
|---|---|---|---|
| 3760 | V | 2025-12-24 | 339.050002 |
| 3761 | V | 2025-12-26 | 340.110501 |
| 3762 | V | 2025-12-29 | 341.119000 |
| 3763 | V | 2025-12-30 | 342.280499 |
| 3764 | V | 2025-12-31 | 343.334999 |
并行模型训练
模型训练默认按顺序运行。向 backtest 或 walkforward 传入 parallel_models=True,会让每个模型在已配置的工作进程中以各自独立的任务进行训练。
本示例改编自 训练模型 文档中的线性回归模型:
[4]:
from sklearn.linear_model import LinearRegression
from pybroker.indicator import close_minus_ma
cmma_20 = close_minus_ma("cmma_20", lookback=20, atr_length=20)
def train_slr(symbol, train_data, test_data):
# Predict the next bar's return given the 20-day CMMA.
prev_close = train_data["close"].shift(1)
daily_returns = (train_data["close"] - prev_close) / prev_close
train_data["pred"] = daily_returns.shift(-1)
train_data = train_data.dropna()
model = LinearRegression()
model.fit(train_data[["cmma_20"]], train_data[["pred"]])
# Return the trained model and columns to use as input data.
return model, ["cmma_20"]
model_slr = pybroker.model("slr", train_slr, indicators=[cmma_20])
def hold_long(ctx):
if not ctx.long_pos():
if ctx.preds("slr")[-1] > 0:
ctx.buy_shares = 100
elif ctx.preds("slr")[-1] < 0:
ctx.sell_all_shares()
model_strategy = Strategy(yfinance, start_date="1/1/2021", end_date="1/1/2026")
model_strategy.add_execution(hold_long, ["V", "MA", "AXP"], models=model_slr)
result = model_strategy.walkforward(
windows=2,
train_size=0.5,
lookahead=1,
warmup=20,
parallel_models=True,
)
result.metrics_df.head()
Backtesting: 2021-01-01 00:00:00 to 2026-01-01 00:00:00
Loaded cached bar data.
Computing indicators...
100% (3 of 3) |##########################| Elapsed Time: 0:00:00 Time: 0:00:00
Train split: 2021-01-05 00:00:00 to 2022-08-31 00:00:00
Finished training models: 0:00:01
Test split: 2022-09-01 00:00:00 to 2024-05-01 00:00:00
100% (418 of 418) |######################| Elapsed Time: 0:00:00 Time: 0:00:00
Train split: 2022-09-01 00:00:00 to 2024-05-01 00:00:00
Finished training models: 0:00:01
Test split: 2024-05-02 00:00:00 to 2025-12-31 00:00:00
100% (418 of 418) |######################| Elapsed Time: 0:00:00 Time: 0:00:00
Finished backtest: 0:00:03
[4]:
| name | value | |
|---|---|---|
| 0 | trade_count | 52 |
| 1 | initial_market_value | 100000.0 |
| 2 | end_market_value | 146555.0 |
| 3 | total_pnl | 16058.0 |
| 4 | unrealized_pnl | 30497.0 |
使用 Ray 作为后端
Ray 可以将同样的工作分布到多个核心甚至整个集群上。使用 pip install ray 进行安装,然后通过 register_ray (来自 ray.util.joblib)将其注册到 Joblib:
[5]:
import ray
from ray.util.joblib import register_ray
ray.init(num_cpus=2, include_dashboard=False, ignore_reinit_error=True)
register_ray()
2026-08-11 13:30:35,114 INFO worker.py:2024 -- Started a local Ray instance.
注册 Ray 之后,向 set_parallel 传入 backend="ray",即可将其作为后端使用:
[6]:
set_parallel(backend="ray", n_jobs=-1)
print(get_parallel_config())
ParallelConfig(n_jobs=-1, backend='ray', parallel=None)
PyBroker 现在会将 Ray 后端用于所有并行任务。例如,以 parallel_indicators=True 调用 backtest:
[7]:
result = strategy.backtest(parallel_indicators=True, warmup=20)
print(f"Total return: {result.metrics.total_return_pct:.2f}%")
ray.shutdown()
Backtesting: 2021-01-01 00:00:00 to 2026-01-01 00:00:00
Loaded cached bar data.
Computing indicators...
100% (3 of 3) |##########################| Elapsed Time: 0:00:02 Time: 0:00:02
Test split: 2021-01-04 00:00:00 to 2025-12-31 00:00:00
100% (1255 of 1255) |####################| Elapsed Time: 0:00:00 Time: 0:00:00
Finished backtest: 0:00:02
Total return: 15.84%
我们将在下一篇文档中探讨 参数优化,在某些场景下,参数优化同样可以并行化。