配置并行化

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)

并行计算指标

指标是按**每个股票代码**计算的:给定股票代码的所有指标都会被归入一个任务,每个股票代码分派一个任务。向 backtestwalkforwardoptimize 传入 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

并行模型训练

模型训练默认按顺序运行。向 backtestwalkforward 传入 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%

我们将在下一篇文档中探讨 参数优化,在某些场景下,参数优化同样可以并行化。