start experiment 64 (exp/64-scheduled-algo-retrain-on-2026-08-27-tac)

This commit is contained in:
zhaoli
2026-08-28 12:02:45 +00:00
parent 1142cc9eef
commit f63eda9e51
14 changed files with 191 additions and 18 deletions
+177 -4
View File
@@ -15,6 +15,9 @@ import os
from inspect import getfullargspec
from typing import List, Optional, Tuple, Union
import numpy as np
import pandas as pd
from qlib.data.dataset import processor as processor_module
from qlib.data.dataset.handler import DataHandlerLP
from qlib.utils import get_callable_kwargs
@@ -144,6 +147,174 @@ class DropAllNaN(processor_module.Processor):
return df
class BenchResidual(processor_module.Processor):
"""Subtract a benchmark instrument's forward return from the label, per datetime.
Turns the training target from an absolute-return rank into a *residual* rank:
``r_i - r_bench`` is ranked cross-sectionally by the downstream ``CSRankNorm`` /
``CSZScoreNorm`` processors instead of ``r_i`` alone. Must be inserted BEFORE any
per-date normalization so the ranking itself is computed on residual returns
(ordering flips exactly where the benchmark trends).
Stateless: ``fit`` is a no-op and the benchmark forward return is recomputed from
the lake parquet on first ``__call__``. Rows whose benchmark value is missing are
left untouched. Accepts ``fit_start_time``/``fit_end_time`` (ignored) so
``check_transform_proc`` can inject the fit window uniformly.
NOTE: under any cross-sectional normalization downstream (``CSRankNorm`` /
``CSZScoreNorm``) this processor is a mathematical no-op: subtracting the same
per-date constant preserves ranks, and z-scoring absorbs constant shifts. Use
``BenchBetaResidual`` for a target that actually reorders.
"""
def __init__(
self,
benchmark="SPY",
fields_group="label",
lake_root=None,
market="US",
timeframe=None,
freq="day",
fit_start_time=None,
fit_end_time=None,
):
self.benchmark = benchmark
self.fields_group = fields_group
self.lake_root = lake_root
self.market = market
self.timeframe = timeframe or timeframe_for_freq(freq)
self.fit_start_time = fit_start_time
self.fit_end_time = fit_end_time
self._bench_label = None
def _load_bench_label(self):
if self._bench_label is not None:
return self._bench_label
cfg = LakeConfig(self.lake_root, self.market)
p = cfg.bar_path(self.timeframe, self.benchmark)
if not p.exists():
raise FileNotFoundError(f"BenchResidual: benchmark bar file not found: {p}")
df = pd.read_parquet(p)
s = pd.Series(df["c"].astype(float).values, index=pd.to_datetime(df["t"])).sort_index()
s.index = s.index.normalize()
# mirror Ref($close,-6)/Ref($close,-1)-1 on the benchmark's own calendar
bench_label = s.shift(-6) / s.shift(-1) - 1
self._bench_label = bench_label[~bench_label.index.duplicated(keep="last")]
return self._bench_label
def fit(self, df=None):
return self
def __call__(self, df):
bl = self._load_bench_label()
cols = processor_module.get_group_columns(df, self.fields_group)
dt = df.index.get_level_values("datetime")
aligned = bl.reindex(pd.DatetimeIndex(dt.unique())).reindex(dt)
mask = aligned.notna().values
out = df.copy()
for c in cols:
vals = df[c].values
res = vals.copy()
res[mask] = np.asarray(vals[mask], dtype=float) - aligned[mask].values
out[c] = res
return out
class BenchBetaResidual(processor_module.Processor):
"""Residualize the label against a beta-scaled benchmark move: ``r_i - b_i * r_bench``.
Unlike a plain constant subtraction (see ``BenchResidual``), the name-specific rolling
beta ``b_i`` makes this survive cross-sectional normalization: in up-weeks high-beta
names lose rank, in down-weeks they gain — exactly the relative structure an absolute-
return ranking hides.
Beta is estimated from *past* data only (rolling ``window`` trading days of daily close
returns of each instrument vs the benchmark, both read up to and including ``t``), so
no lookahead enters the target. The benchmark leg uses the same horizon as the label
expression (``Ref($close,-6)/Ref($close,-1)-1`` by default via ``horizon``/``base``,
matching the yaml's 6-day label). Rows with missing beta or benchmark values keep
their raw label.
Requires ``$close`` to be present in the feature group (it always is for TACHandler).
Stateless; accepts ``fit_start_time``/``fit_end_time`` (ignored) for uniform kwargs
injection. Must be inserted BEFORE any per-date normalization processor.
"""
def __init__(
self,
benchmark="SPY",
fields_group="label",
lake_root=None,
market="US",
timeframe=None,
freq="day",
window=63,
horizon=6,
base=1,
feature_field="$close",
fit_start_time=None,
fit_end_time=None,
):
self.benchmark = benchmark
self.fields_group = fields_group
self.lake_root = lake_root
self.market = market
self.timeframe = timeframe or timeframe_for_freq(freq)
self.window = int(window)
self.horizon = int(horizon)
self.base = int(base)
self.feature_field = feature_field
self.fit_start_time = fit_start_time
self.fit_end_time = fit_end_time
self._bench = None
def _load_bench_close(self):
if self._bench is not None:
return self._bench
cfg = LakeConfig(self.lake_root, self.market)
p = cfg.bar_path(self.timeframe, self.benchmark)
if not p.exists():
raise FileNotFoundError(f"BenchBetaResidual: benchmark bar file not found: {p}")
df = pd.read_parquet(p)
s = pd.Series(df["c"].astype(float).values, index=pd.to_datetime(df["t"])).sort_index()
s.index = s.index.normalize()
self._bench = s[~s.index.duplicated(keep="last")]
return self._bench
def fit(self, df=None):
return self
def __call__(self, df):
bench = self._load_bench_close()
# benchmark forward return over the same horizon as the label expression
fwd = bench.shift(-(self.base + self.horizon - 1)) / bench.shift(-self.base) - 1
px_col = ("feature", self.feature_field)
if px_col not in df.columns:
raise KeyError(f"BenchBetaResidual: {self.feature_field} not found in features")
px = df[px_col].unstack("instrument").sort_index()
rets = px / px.shift(1) - 1
bret = bench.reindex(px.index).pct_change()
# rolling beta per instrument using data <= t (no lookahead)
cov = rets.rolling(self.window, min_periods=max(10, self.window // 2)).cov(bret)
var = bret.rolling(self.window, min_periods=max(10, self.window // 2)).var()
beta = cov.div(var, axis=0)
contrib = beta.mul(fwd.reindex(px.index), axis=0)
cols = list(processor_module.get_group_columns(df, self.fields_group))
out = df.copy()
for c in cols:
lab = df[c].unstack("instrument").reindex(px.index)
resid = lab - contrib.where(contrib.notna() & lab.notna(), 0.0)
new_vals = resid.stack()
new_vals.index.names = df.index.names
# residual where available, raw label otherwise (e.g. beta warm-up rows)
out[c] = new_vals.reindex(out.index).fillna(df[c])
return out
class TACHandler(DataHandlerLP):
"""DataHandlerLP backed by the TradeAC parquet lake.
@@ -246,10 +417,12 @@ class TACHandler(DataHandlerLP):
return get_common_feature_fields(lake_root, market, timeframe_for_freq(freq))
__all__ = ["TACHandler", "DropAllNaN", "get_common_feature_fields"]
__all__ = ["TACHandler", "DropAllNaN", "BenchResidual", "BenchBetaResidual", "get_common_feature_fields"]
# Make `DropAllNaN` resolvable by bare name from processor configs (e.g. the default
# ``infer_processors`` and workflow yamls that reference it without a ``module_path``),
# mirroring how qlib registers its own processors in ``qlib.data.dataset.processor``.
# Make `DropAllNaN`/`BenchResidual`/`BenchBetaResidual` resolvable by bare name from processor
# configs (e.g. the default ``infer_processors`` and workflow yamls that reference them without a
# ``module_path``), mirroring how qlib registers its own processors in ``qlib.data.dataset.processor``.
processor_module.DropAllNaN = DropAllNaN
processor_module.BenchResidual = BenchResidual
processor_module.BenchBetaResidual = BenchBetaResidual