pybroker.model 模块
包含与模型相关的功能。
- class CachedModel(model: Any, input_cols: tuple[str] | None, lag_columns: tuple[str, ...] | None = None)[源代码]
基类:
NamedTuple存储已缓存的模型数据。
- lag_columns
训练时用于构建滞后特征的列名称,按特征块顺序排列。如果模型未使用
lags训练,或加载的是本字段出现之前缓存的模型,则为None。
- class IntervalBoundModel(source: ModelSource, intervals: frozenset[TimeframeInterval])[源代码]
基类:
NamedTuple绑定到一个或多个压缩区间的
ModelSource,由ModelSource.intervals()返回,并传递给pybroker.strategy.Strategy.add_execution()的models参数。- intervals: frozenset[int | Literal['daily', 'weekly', 'monthly', 'quarterly', 'yearly'] | str]
该模型据以训练的规范化
TimeframeInterval。可能包含字面量'base'表示基础时间框架。
- source: ModelSource
已绑定的
ModelSource。
- class LagSeriesKey(symbol: str, column: str, lag: int, interval: str | None = None)[源代码]
基类:
object完整历史滞后序列的内部缓存键。
- class ModelInput(columns: tuple[str, ...], arrays: dict[str, ndarray], dates: ndarray, lag_features: ndarray | None = None, lags: int | None = None, lag_columns: tuple[str, ...] | None = None)[源代码]
基类:
object内部的、numpy 支持的模型输入,带有可选的滞后特征元数据。
不属于公共 API。面向用户的代码会收到通过
to_dataframe()具体化的pandas.DataFrame实例,任何滞后特征矩阵都会作为单独的numpy.ndarray参数显式传递。- drop_lag_warmup() ModelInput[源代码]
丢弃滞后特征尚未定义的前导行。
仅裁剪每个品种自身的预热区间。若丢弃所有存在 NaN 滞后特征的行,每当滞后列稀疏时,也会一并移除序列中间的行 —— 通过
pybroker.scope.register_columns()注册的事件列,除了触发的那些 K 线外都是 NaN,而对于指数、外汇和 CFD 序列,全 NaN 的volume列也是常见情况。这样做要么会悄然缩小训练集,要么会将其清空。池化(pooled)输入按每个品种堆叠一个数据块,因此每个数据块都在各自的偏移处带有自己的预热期。仅裁剪矩阵前端,会使第一个数据块之后的每个数据块仍保留其 NaN 行,而这会被
sklearn.fit直接拒绝。
- lag_warmup_len() int[源代码]
返回有多少前导行的滞后特征未定义。
这些行不能交给估计器:sklearn 会直接拒绝 NaN,而一个能容忍 NaN 的 predict_fn 则会悄悄地根据尚不存在的特征生成一个预测结果。
- select_columns(columns: tuple[str, ...]) ModelInput[源代码]
返回一个限定在
columns范围内的视图。
- slice(end_index: int | None = None) ModelInput[源代码]
返回一个共享底层数组内存的行切片。
- class ModelLoader(name: str, load_fn: Callable[[...], Any | tuple[Any, Iterable[str]]], indicator_names: Iterable[str], input_data_fn: Callable[[DataFrame], DataFrame] | None, predict_fn: Callable[[Any, DataFrame | ndarray[tuple[Any, ...], dtype[_ScalarT]]], ndarray[tuple[Any, ...], dtype[_ScalarT]]] | None, pooled: bool, kwargs: dict[str, Any], lags: int | None = None, lag_cols: tuple[str, ...] = (), per_bar: bool = False)[源代码]
基类:
ModelSource加载一个预训练模型。
- 参数:
name -- 模型名称。
load_fn -- 用于加载并返回预训练模型的
Callable[[symbol: str, train_start_date: datetime, train_end_date: datetime, ...], DataFrame]。该函数应返回一个已训练的模型实例,或一个元组,其中包含已训练的模型实例,以及进行预测时用作模型输入的列名称Iterable。indicator_names -- 用作该模型特征的
pybroker.indicator.Indicator名称的Iterable。input_data_fn -- 用于预处理进行预测时传递给模型的输入数据的
Callable[[DataFrame], DataFrame]。如果设置了该值,input_data_fn会以包含全部测试数据的pandas.DataFrame被调用。predict_fn -- 用于覆盖调用模型默认
predict函数的Callable[[Model, DataFrame], ndarray]。如果设置了该值,predict_fn会以已训练的模型和包含全部测试数据的pandas.DataFrame被调用。当设置了lags时,调用时会改为传入滞后特征矩阵(numpy.ndarray)而不是 DataFrame。pooled -- 如果为
True,该模型每次执行会使用合并后的多品种数据训练一次。默认为False。kwargs -- 传递给
load_fn的关键字参数dict。
- class ModelSource(name: str, indicator_names: Iterable[str], input_data_fn: Callable[[DataFrame], DataFrame] | None, predict_fn: Callable[[Any, DataFrame | ndarray[tuple[Any, ...], dtype[_ScalarT]]], ndarray[tuple[Any, ...], dtype[_ScalarT]]] | None, pooled: bool, kwargs: dict[str, Any], lags: int | None = None, lag_cols: tuple[str, ...] = (), per_bar: bool = False)[源代码]
基类:
object模型源的基类。模型源要么通过训练,要么通过加载预训练模型来提供一个模型实例。
- 参数:
name -- 模型名称。
indicator_names -- 用作该模型特征的
pybroker.indicator.Indicator名称的Iterable。input_data_fn -- 用于预处理进行预测时传递给模型的输入数据的
Callable[[DataFrame], DataFrame]。如果设置了该值,input_data_fn会以包含全部测试数据的pandas.DataFrame被调用。predict_fn -- 用于覆盖调用模型默认
predict函数的Callable[[Model, DataFrame], ndarray]。如果设置了该值,predict_fn会以已训练的模型和包含全部测试数据的pandas.DataFrame被调用。当设置了lags时,调用时会改为传入滞后特征矩阵(numpy.ndarray)而不是 DataFrame。lags --
lag_cols中每一列要包含的滞后值数量,生成一个形状为(n_rows, len(lag_cols) * (lags + 1))的特征矩阵,该矩阵会以lag_train/lag_test的形式传递给train_fn,并替代输入 DataFrame 传递给predict_fn,而不是作为列添加。lag_cols -- 要计算滞后值的列。默认为训练数据的数据列,不包括
date和symbol;只有在此处指定名称的指标才会被计算滞后值。per_bar -- 如果为
True,predict_fn会按每根 K 线调用一次,输入会被截断至包含当前 K 线在内的各行。pooled -- 如果为
True,该模型每次执行会使用合并后的多品种数据训练一次。默认为False。kwargs -- 附加关键字参数的
dict。
- intervals(*intervals: int | Literal['daily', 'weekly', 'monthly', 'quarterly', 'yearly'] | str) IntervalBoundModel[源代码]
将此模型绑定到一个或多个压缩区间,以便与
pybroker.strategy.Strategy.add_execution()搭配使用。已绑定的模型只会在所列出的区间的压缩 K 线上训练,连同注册给它的任何指标一起。绑定会替代默认的基础时间框架训练;在
intervals中包含字面量'base'即可同时在基础时间框架上训练该模型。逐区间的预测结果通过pybroker.context.IntervalContext.preds()读取。绑定的区间会自动通过pybroker.context.ExecContext.interval()提供,无需在add_execution()的intervals参数中再次声明它们:trend = pybroker.model("trend", train_fn, indicators=[sma_10]) strategy.add_execution( fn, "SPY", models=trend.intervals("base", "weekly") )
只有可训练的模型才支持区间绑定。对预训练模型(
ModelLoader)调用本方法会引发ValueError。- 参数:
intervals -- 训练此模型所使用的一个或多个
TimeframeInterval,每一个都必须严格粗于回测数据的基础 K 线间隔,或使用字面量'base'表示基础时间框架。- 返回:
将此模型绑定到
intervals的IntervalBoundModel。
- prepare_input_data(df: DataFrame) DataFrame[源代码]
为进行预测时传递给模型准备输入数据的
pandas.DataFrame。如果设置了该值,则使用input_data_fn对输入数据进行预处理。如果为False,则使用df中的指标列作为输入特征。
- class ModelTrainer(name: str, train_fn: Callable[[...], Any | tuple[Any, Iterable[str]]], indicator_names: Iterable[str], input_data_fn: Callable[[DataFrame], DataFrame] | None, predict_fn: Callable[[Any, DataFrame | ndarray[tuple[Any, ...], dtype[_ScalarT]]], ndarray[tuple[Any, ...], dtype[_ScalarT]]] | None, pooled: bool, kwargs: dict[str, Any], lags: int | None = None, lag_cols: tuple[str, ...] = (), per_bar: bool = False)[源代码]
基类:
ModelSource训练一个模型。
- 参数:
name -- 模型名称。
train_fn -- 当
pooled为False时,为Callable[[symbol: str, train_data: DataFrame, test_data: DataFrame, ...], DataFrame]。当pooled为True时,为Callable[[symbols: Sequence[str], train_data: DataFrame, test_data: DataFrame, ...], DataFrame]。当设置了lags时,train_fn还会额外接收lag_train=和lag_test=关键字参数,其中包含与train_data/test_data逐行对齐的滞后特征矩阵,且该函数必须同时接受这两个参数。该函数应返回一个已训练的模型实例,或一个元组,其中包含已训练的模型实例,以及进行预测时用作模型输入的列名称Iterable。indicator_names -- 用作该模型特征的
pybroker.indicator.Indicator名称的Iterable。input_data_fn -- 用于预处理进行预测时传递给模型的输入数据的
Callable[[DataFrame], DataFrame]。如果设置了该值,input_data_fn会以包含全部测试数据的pandas.DataFrame被调用。predict_fn -- 用于覆盖调用模型默认
predict函数的Callable[[Model, DataFrame], ndarray]。如果设置了该值,predict_fn会以已训练的模型和包含全部测试数据的pandas.DataFrame被调用。当设置了lags时,调用时会改为传入滞后特征矩阵(numpy.ndarray)而不是 DataFrame。pooled -- 如果为
True,该模型每次执行会使用合并后的多品种数据训练一次。默认为False。kwargs -- 传递给
train_fn的关键字参数dict。
- __call__(symbol: str, train_data: DataFrame, test_data: DataFrame, *, lag_train: ndarray[tuple[Any, ...], dtype[_ScalarT]] | None = None, lag_test: ndarray[tuple[Any, ...], dtype[_ScalarT]] | None = None) Any | tuple[Any, Iterable[str]][源代码]
按品种训练模型。
- 参数:
symbol -- 该模型的股票代码(模型按品种训练)。
train_data -- 训练数据。
test_data -- 测试数据。
lag_train -- 与
train_data逐行对齐的滞后特征矩阵。当模型使用lags注册时,会以lag_train=传递给train_fn。lag_test -- 与
test_data逐行对齐的滞后特征矩阵。当模型使用lags注册时,会以lag_test=传递给train_fn。
- 返回:
已训练的模型。
- train_pooled(symbols: Sequence[str], train_data: DataFrame, test_data: DataFrame, *, lag_train: ndarray[tuple[Any, ...], dtype[_ScalarT]] | None = None, lag_test: ndarray[tuple[Any, ...], dtype[_ScalarT]] | None = None) Any | tuple[Any, Iterable[str]][源代码]
使用合并后的多品种数据训练模型。
- 参数:
symbols -- 该池化分组的股票代码,按升序排列,与
train_data和test_data中品种数据块出现的顺序一致。列表中的某个品种在某个数据帧中可能没有任何行,例如当lags因滞后预热而丢弃了其所有行时。train_data -- 包含
symbol列的训练数据。test_data -- 包含
symbol列的测试数据。lag_train -- 与
train_data逐行对齐的滞后特征矩阵。当模型使用lags注册时,会以lag_train=传递给train_fn。lag_test -- 与
test_data逐行对齐的滞后特征矩阵。当模型使用lags注册时,会以lag_test=传递给train_fn。
- 返回:
已训练的模型。
- class ModelsMixin[源代码]
基类:
object实现与模型相关功能的 Mixin。
- train_models(model_syms: Iterable[ModelSymbol], train_data: DataFrame, test_data: DataFrame, indicator_data: Mapping[IndicatorSymbol, Series], cache_date_fields: CacheDateFields, parallel_models: bool = False, pooled_model_groups: Mapping[tuple[str, int], frozenset[str]] | None = None, interval_data: IntervalData | None = None, *, history_store: SymbolArrayStore | None = None, train_store: SymbolArrayStore | None = None, test_store: SymbolArrayStore | None = None, lookahead: int = 1) dict[ModelSymbol, TrainedModel][源代码]
为给定的
pybroker.common.ModelSymbol组合训练模型。- 参数:
model_syms -- 待训练模型的
pybroker.common.ModelSymbol组合的Iterable。train_data -- 训练数据的
pandas.DataFrame。test_data -- 测试数据的
pandas.DataFrame。indicator_data -- 将
pybroker.common.IndicatorSymbol组合映射到pybroker.indicator.Indicator值pandas.Series的Mapping。cache_date_fields -- 用于为缓存数据生成键的日期字段。
parallel_models -- 如果为
True,ModelTrainer模型会使用多个进程并行训练。默认为False。pooled_model_groups -- 将
(model_name, execution_id)组合映射到用于池化训练的品种frozenset[str]的Mapping。默认为None。lookahead -- 预测目标未来的 K 线数量,以每个模型所拟合的时间框架的 K 线为单位:绑定到某个区间的模型,会在其训练行和测试行之间保留
lookahead根压缩 K 线。默认为1。
- 返回:
将每个
pybroker.common.ModelSymbol组合映射到pybroker.common.TrainedModel的dict。
- apply_lags_to_model_input(model_input: ModelInput, lag_columns: tuple[str, ...], lags: int, lag_cache: dict[LagSeriesKey, ndarray], symbol: str, history_dates: ndarray, interval: str | None = None) ModelInput[源代码]
为
model_input附加滞后特征元数据。
- apply_lags_to_model_input_pooled(model_input: ModelInput, lag_columns: tuple[str, ...], lags: int, lag_cache: dict[LagSeriesKey, ndarray], history_dates_by_symbol: dict[str, ndarray], symbols: Iterable[str], interval: str | None = None) ModelInput[源代码]
为池化的
model_input附加滞后特征元数据。
- apply_prepare_input_data(model_input: ModelInput, prepare_fn: Callable[[DataFrame], DataFrame]) ModelInput[源代码]
对
model_input应用一个仅接受 DataFrame 的准备函数。
- build_lag_feature_matrix(symbol: str, columns: tuple[str, ...], lags: int, row_dates: ndarray, history_dates: ndarray, lag_cache: dict[LagSeriesKey, ndarray], interval: str | None = None) ndarray[源代码]
根据 numpy 数组构建一个滞后展开的特征矩阵。
返回一个形状为
(len(row_dates), len(columns) * (lags + 1))的矩阵,按每列一个连续数据块排列,每个数据块依次存放该列的当前值,以及滞后 1 到lags的值。
- build_lag_feature_matrix_pooled(sym_col: ndarray, columns: tuple[str, ...], lags: int, row_dates: ndarray, history_dates_by_symbol: dict[str, ndarray], lag_cache: dict[LagSeriesKey, ndarray], symbols: Iterable[str], interval: str | None = None) ndarray[源代码]
为池化的多品种数据构建一个滞后展开的特征矩阵。
- cached_stacked_lags(cache: dict[LagSeriesKey, ndarray], symbol: str, col: str, lags: int, interval: str | None = None) ndarray | None[源代码]
返回一个深度足以满足
lags的已缓存堆叠数组,否则返回None。
- compute_lag_series_cache(df: DataFrame, symbols: Iterable[str], columns: tuple[str, ...], lags: int) dict[LagSeriesKey, ndarray][源代码]
为日线/基础 K 线计算完整历史的滞后数组。
- history_date_offset(history_dates: ndarray, row_dates: ndarray) int[源代码]
返回
row_dates在history_dates中的起始索引。
- merge_interval_lag_series_cache(cache: dict[LagSeriesKey, ndarray], symbols: Iterable[str], columns: tuple[str, ...], lags: int, interval: str, bars_by_symbol, arrays_by_symbol=None) dict[LagSeriesKey, ndarray][源代码]
将完整历史的区间滞后数组添加到
cache中。pybroker.interval.CompressedBars仅包含数据列,因此arrays_by_symbol会补充压缩 K 线无法提供的列 —— 尤其是指标值 —— 且在给定时会被优先查询。
- merge_lag_series_cache(cache: dict[LagSeriesKey, ndarray], history_df: DataFrame, symbols: Iterable[str], columns: tuple[str, ...], lags: int, history_dates: dict[str, ndarray] | None = None) dict[LagSeriesKey, ndarray][源代码]
将
columns的完整历史滞后数组添加到cache中。
- merge_lag_series_cache_from_arrays(cache: dict[LagSeriesKey, ndarray], symbol: str, columns: tuple[str, ...], lags: int, history_dates: ndarray, column_arrays: Mapping[str, ndarray]) None[源代码]
添加根据 numpy 列数据构建的完整历史滞后数组。
- merge_lag_series_cache_from_store(cache: dict[LagSeriesKey, ndarray], store: SymbolArrayStore, symbols: Iterable[str], columns: tuple[str, ...], lags: int, history_dates: dict[str, ndarray] | None = None, indicators: tuple[str, ...] = (), indicator_data: Mapping[IndicatorSymbol, Series] | None = None) dict[LagSeriesKey, ndarray][源代码]
从
pybroker.scope.SymbolArrayStore添加完整历史滞后数组。pybroker.scope.SymbolArrayStore仅包含数据列,因此指标值会从覆盖各品种完整历史的indicator_data中,对齐到该存储的日期上。否则,即使指标是模型输入的一列,对其计算滞后值也会失败。
- model(name: str, fn: Callable[[...], Any | tuple[Any, Iterable[str]]], indicators: Iterable[Indicator] | None = None, lags: int | None = None, lag_cols: Iterable[str | Indicator] | None = None, per_bar: bool = False, input_data_fn: Callable[[DataFrame], DataFrame] | None = None, predict_fn: Callable[[Any, DataFrame | ndarray[tuple[Any, ...], dtype[_ScalarT]]], ndarray[tuple[Any, ...], dtype[_ScalarT]]] | None = None, pretrained: bool = False, pooled: bool = False, **kwargs) ModelSource[源代码]
创建一个
ModelSource实例,并以name全局注册。- 参数:
name -- 用于全局引用该模型的名称。
fn -- 用于训练或加载模型实例的
Callable。若用于pooled=False的训练,则fn的签名为Callable[[symbol: str, train_data: DataFrame, test_data: DataFrame, ...], DataFrame]。若用于pooled=True的训练,则fn的签名为Callable[[symbols: Sequence[str], train_data: DataFrame, test_data: DataFrame, ...], DataFrame],其中symbols包含按升序排列的池化品种,且两个数据帧都包含一个symbol列,各品种的行按该顺序分组排列。列表中的某个品种在某个数据帧中可能没有任何行,例如当lags因滞后预热而丢弃了其所有行时。若用于加载,则fn的签名为Callable[[symbol: str, train_start_date: datetime, train_end_date: datetime, ...], DataFrame]。当设置了lags时,用于训练的fn还会额外接收lag_train=和lag_test=关键字参数,其中包含预先构建好的滞后特征矩阵,且该函数必须同时接受这两个参数。该函数应返回一个已训练的模型实例,或一个元组,其中包含已训练的模型实例,以及进行预测时用作模型输入的列名称Iterable。当只返回模型实例时,会使用训练 DataFrame 中的列用于预测。对于池化模型,推断出的预测列中会省略symbol列。indicators -- 用作该模型特征的
pybroker.indicator.Indicator的Iterable。lags --
lag_cols中每一列要包含的滞后值数量。这些滞后值会被构建为一个形状为(n_rows, len(lag_cols) * (lags + 1))的特征矩阵:按lag_cols声明顺序,每列一个连续数据块,每个数据块依次存放该列的当前值,以及滞后1到lags的值,其中滞后1是前一根 K 线的值。该矩阵会以lag_train和lag_test关键字参数的形式传递给用于训练的fn,与train_data/test_data逐行对齐,并替代输入 DataFrame 传递给predict_fn``(或模型默认的 ``predict)。由于当前 K 线的值是每个数据块的第一个特征,预测目标应为 下一根 K 线 —— 若以当前 K 线的值作为目标则会造成数据泄漏。滞后数据与模型输入分开保存,而不是作为列添加,因此输入的pandas.DataFrame永远不会被加宽或复制。滞后值是根据每个品种的完整历史计算的,因此测试窗口开头的行会使用从前一个训练窗口延续下来的真实值,而不是NaN。滞后值未定义的行仅会从训练数据中被丢弃。lag_cols -- 要计算滞后值的列名称和/或
pybroker.indicator.Indicator。此处传入的pybroker.indicator.Indicator会被添加到indicators中。声明顺序决定了特征块的顺序。默认为训练数据的数据列,不包括date和symbol;只有在此处指定名称的指标才会被计算滞后值。per_bar -- 如果为
True,predict_fn会按每根 K 线调用一次,输入会被截断至包含当前 K 线在内的各行,且必须为该根 K 线返回一个标量预测值。当设置了lags时,输入是以同样方式截断的滞后特征矩阵,最后一行为当前 K 线。适用于每根 K 线都重新拟合或更新的模型,例如 GARCH 或状态空间模型。请注意,这会使每个品种每根 K 线都调用一次模型,远慢于对整个测试窗口进行一次向量化的predict_fn调用。需要设置predict_fn,且不支持与pooled=True同时使用。input_data_fn -- 用于预处理进行预测时传递给模型的输入数据的
Callable[[DataFrame], DataFrame]。如果设置了该值,input_data_fn会以包含全部测试数据的pandas.DataFrame被调用,即使在per_bar=True时也是如此。它必须每根 K 线返回一行;增加或删除行会使预测结果与 K 线错位,并引发ValueError。对于使用lags注册的模型,它仅影响pybroker.context.ExecContext.input()的形态,因为预测是根据滞后特征矩阵进行的。predict_fn -- 用于覆盖调用模型默认
predict函数的Callable[[Model, DataFrame], ndarray]。如果设置了该值,predict_fn会以已训练的模型和包含全部测试数据的pandas.DataFrame被调用。当设置了lags时,调用时会改为传入滞后特征矩阵(numpy.ndarray)而不是 DataFrame。当per_bar=True时,predict_fn会接收包含当前 K 线在内的输入行,并且必须返回一个标量预测值。pretrained -- 如果为
True,则使用fn加载并返回一个预训练模型。如果为False,则使用fn训练并返回一个新模型。默认为False。pooled -- 如果为
True,该模型每次执行会使用合并后的多品种数据训练一次。默认为False。**kwargs -- 传递给
fn的附加参数。
- 返回:
ModelSource实例。
- model_input_from_arrays(columns: tuple[str, ...], arrays: dict[str, ndarray], dates: ndarray) ModelInput[源代码]
根据列数组构建一个
ModelInput。
- model_input_from_frame(df: DataFrame, columns: tuple[str, ...] | None = None, dates: ndarray | None = None) ModelInput[源代码]
从 DataFrame 构建一个
ModelInput,且不复制列。