start experiment 42 (exp/42-q10-hmm-regime-overlay-entry-gate-on-sph)

This commit is contained in:
zhaoli
2026-08-20 00:36:58 +00:00
parent 202589fc3e
commit 513fcdf344
24 changed files with 945 additions and 315 deletions
+30 -12
View File
@@ -64,9 +64,13 @@ def check_transform_proc(proc_l, fit_start_time, fit_end_time):
def get_common_feature_fields(lake_root=None, market="US", timeframe="1d") -> List[str]:
"""Discover ta-lib columns present in *every* features parquet file of the lake.
"""Discover feature columns present in *every* feature file of the lake.
Returns sorted field names (without the ``$`` prefix). Empty if no features are persisted.
Walks the `family=ta|sp` partition layout (plus any legacy flat files).
TA and SP columns are disjoint by construction, so the common set is
computed per family (columns shared by all symbol files of that family),
then the per-family results are unioned. Returns sorted field names
(without the ``$`` prefix). Empty if no features are persisted.
"""
cfg = LakeConfig(lake_root, market)
feat_dir = cfg.features_dir(timeframe)
@@ -74,16 +78,30 @@ def get_common_feature_fields(lake_root=None, market="US", timeframe="1d") -> Li
return []
import pyarrow.parquet as pq
common = None
for p in sorted(feat_dir.glob("symbol=*.parquet")):
try:
cols = set(pq.read_schema(p).names) - set(NON_FEATURE_COLUMNS)
except Exception: # pragma: no cover - skip unreadable files
continue
common = cols if common is None else (common & cols)
if not common:
break
return sorted(common) if common else []
def _family_common(fam_dir: Path) -> set:
common = None
for p in sorted(fam_dir.glob("symbol=*.parquet")):
try:
cols = set(pq.read_schema(p).names) - set(NON_FEATURE_COLUMNS)
except Exception: # pragma: no cover - skip unreadable files
continue
common = cols if common is None else (common & cols)
if not common:
break
return common or set()
common: set = set()
# family tier: features/market=*/timeframe=*/family=*/symbol=*.parquet
for fam in ("ta", "sp"):
fam_dir = feat_dir / f"family={fam}"
if fam_dir.is_dir():
common |= _family_common(fam_dir)
# legacy flat: features/market=*/timeframe=*/symbol=*.parquet
if (feat_dir / "family=ta").exists() or (feat_dir / "family=sp").exists():
pass # family layout already covered
else:
common |= _family_common(feat_dir)
return sorted(common)
class DropAllNaN(processor_module.Processor):
@@ -56,7 +56,6 @@ import os
from concurrent.futures import ThreadPoolExecutor
from typing import List, Optional
import numpy as np
import pandas as pd
from qlib.data.dataset import DatasetH
@@ -80,15 +79,11 @@ class RankICEnsembleLGBModel(RankICLGBModel):
forwarded.
"""
def __init__(self, seeds: str = "42", parallel: int = 0, weight_mode: str = "equal", **kwargs):
def __init__(self, seeds: str = "42", parallel: int = 0, **kwargs):
self.seeds = [int(s.strip()) for s in str(seeds).split(",") if s.strip()]
if not self.seeds:
raise ValueError("seeds must contain at least one integer")
self.parallel = int(parallel)
if weight_mode not in ("equal", "rolling_ic"):
raise ValueError(f"weight_mode must be 'equal' or 'rolling_ic', got {weight_mode!r}")
self.weight_mode = weight_mode
self.rolling_ic_window = int(kwargs.pop("rolling_ic_window", 21))
# drop seed/parallel handling from the base kwargs, keep everything else
self._model_kwargs = dict(kwargs)
super().__init__(**self._model_kwargs)
@@ -184,44 +179,11 @@ class RankICEnsembleLGBModel(RankICLGBModel):
# -------------------------------------------------------------- predict
def predict(self, dataset: DatasetH, segment="test") -> pd.Series:
"""Combine per-seed predictions.
``weight_mode='equal'`` (default): simple average, as before.
``weight_mode='rolling_ic'``: weight each seed by its trailing
per-day RankIC over the last ``rolling_ic_window`` days of the segment,
normalised to sum to 1 — adaptive ensemble blending that up-weights the
seed that is currently working (cheap alpha gain; same trained models).
"""
"""Average the per-seed predictions over the given segment."""
if not self._models:
raise ValueError("model is not fitted yet!")
preds = [m.predict(dataset, segment=segment) for m in self._models]
if len(preds) == 1:
return preds[0]
frame = pd.concat(preds, axis=1)
frame.columns = [f"seed{m.params.get('seed', i)}" for i, m in enumerate(self._models)]
if self.weight_mode == "equal":
return frame.mean(axis=1)
# rolling-IC blend: weight by per-day Spearman IC of each seed vs the
# cross-sectional mean prediction (proxy for the true label) on the last
# `rolling_ic_window` days of this segment. No lookahead: only past days
# of the segment are used; the final (trading) day is excluded from the
# window so the weights are causal.
mean_pred = frame.mean(axis=1)
dates = sorted(frame.index.get_level_values(0).unique())
win = [d for d in dates if d < dates[-1]][-self.rolling_ic_window :]
ics = {}
for col in frame.columns:
if not win:
ics[col] = 1.0
continue
sub = pd.DataFrame({"p": frame[col], "m": mean_pred})
vals = []
for d in win:
s = sub[sub.index.get_level_values(0) == d]
if len(s) >= 3 and s["p"].nunique() > 1 and s["m"].nunique() > 1:
vals.append(s["p"].rank().corr(s["m"].rank()))
ics[col] = float(np.mean(vals)) if vals else 1.0
wsum = sum(ics.values()) or len(ics)
weights = {c: v / wsum for c, v in ics.items()}
return sum(frame[c] * weights[c] for c in frame.columns)
return frame.mean(axis=1)
@@ -53,23 +53,61 @@ from qlib.workflow import R
__all__ = ["RankICLGBModel", "rankic_feval"]
def _group_averaged_rank(values: np.ndarray, gid: np.ndarray, offs: np.ndarray) -> np.ndarray:
"""Averaged (tie-corrected) rank of ``values`` within each group, vectorized.
``gid`` maps each row to its group id; ``offs`` holds the cumulative row
offsets so that group ``i`` occupies rows ``[offs[i], offs[i+1])``. Returns
the same result as ``pandas.Series.rank(method='average')`` applied per
group, but in one pass (``np.lexsort`` is the only non-linear step).
"""
n = len(values)
order = np.lexsort((values, gid))
ord_rank = np.empty(n, dtype=np.float64)
ord_rank[order] = np.arange(n, dtype=np.float64) - offs[gid[order]] + 1.0
sg = gid[order]
sv = values[order]
newblock = np.empty(n, dtype=bool)
newblock[0] = True
newblock[1:] = (sg[1:] != sg[:-1]) | (sv[1:] != sv[:-1])
blockid = np.cumsum(newblock) - 1
block_mean = np.bincount(blockid, weights=ord_rank[order]) / np.bincount(blockid)
out = np.empty(n)
out[order] = block_mean[blockid]
return out
def _per_day_spearman(preds: np.ndarray, labels: np.ndarray, group: np.ndarray) -> float:
"""Mean per-day Spearman rank correlation of preds vs labels.
``group`` holds the number of rows of each trading day (query group), in
order. Days with <3 valid rows or a constant pred/label are skipped.
Vectorized: per-day Spearman == Pearson of the per-day rank transforms,
and the Pearson moments (``sum``, ``sum`` of products/squares) aggregate
over each day with ``np.bincount``. Runs ~10x faster than the per-day
``pd.Series.rank()`` loop that preceded it — this feval is invoked on the
train and valid panels every boosting round, per seed.
"""
if group is None or len(group) == 0:
return 0.0
offs = np.concatenate([[0], np.cumsum(group.astype(int))])
vals = []
for i in range(len(group)):
s = slice(offs[i], offs[i + 1])
p, l = preds[s], labels[s]
if len(p) < 3 or np.std(p) == 0 or np.std(l) == 0:
continue
vals.append(np.corrcoef(pd.Series(p).rank(), pd.Series(l).rank())[0, 1])
return float(np.mean(vals)) if vals else 0.0
gid = np.repeat(np.arange(len(group)), group.astype(int))
rp = _group_averaged_rank(preds, gid, offs)
rl = _group_averaged_rank(labels, gid, offs)
n_g = group.astype(float)
s_p = np.bincount(gid, weights=rp)
s_l = np.bincount(gid, weights=rl)
s_pl = np.bincount(gid, weights=rp * rl)
s_pp = np.bincount(gid, weights=rp * rp)
s_ll = np.bincount(gid, weights=rl * rl)
cov = n_g * s_pl - s_p * s_l
var_p = n_g * s_pp - s_p ** 2
var_l = n_g * s_ll - s_l ** 2
denom = np.sqrt(var_p * var_l)
valid = (n_g >= 3) & (denom > 0)
corr = np.where(valid, cov / np.where(denom == 0, 1, denom), 0.0)
return float(corr[valid].mean()) if valid.any() else 0.0
def rankic_feval(preds, dataset):
@@ -1,3 +1,13 @@
from .kelly_dropout import FractionalKellyDropoutStrategy # noqa: F401
from .optimal_stop import OptimalStopControl # noqa: F401
from .regime_gate import RegimeGateDropoutStrategy # noqa: F401
from .top_bottom import TopBottomDropoutStrategy # noqa: F401
from .weekly_rebalance import WeeklyRebalanceDropoutStrategy # noqa: F401
__all__ = ["OptimalStopControl"]
__all__ = [
"OptimalStopControl",
"FractionalKellyDropoutStrategy",
"WeeklyRebalanceDropoutStrategy",
"TopBottomDropoutStrategy",
"RegimeGateDropoutStrategy",
]
@@ -1,138 +0,0 @@
"""TopkDropout with HMM high-volatility + drawdown-pause risk gates.
Gates NEW entries on two risk conditions (held names are never force-sold):
1. **HMM high-vol pause**: when the cross-sectional mean of ``sp_hmm_p_regime1``
(HMM high-vol regime probability) on the signal date is >= ``hmm_pause_pct``,
new buys are paused. The time-series study showed HMM high-vol probability
pulses BEFORE sharp moves (regime-change cut) — pausing new exposure at the
boundary reduces drawdown from price over-reaction.
2. **Drawdown pause**: when the account equity drawdown from its running peak
exceeds ``drawdown_pause_pct``, new buys are paused. This is the
``drawdown_pause_pct`` risk-limit expressed inside the backtest (the pure
executor-side gate is documented as not expressible in a one-shot backtest).
3. **Liquidity floor**: names whose 20-day average daily dollar volume is below
``liquidity_floor_adv`` are dropped from BUY candidates (the proven mitigant
from exp-18: $5M floor cut drawdown 7.9%->5.4% at higher IR).
Implementation: pre-filter the signal score before the base TopkDropout
decision — non-held names get score 0 when any gate fires.
Wired into a workflow yaml like:
strategy:
class: HmmRiskTopk
module_path: tac_qlib.contrib.strategy.hmm_risk
kwargs:
signal: "<PRED>"
topk: 10
n_drop: 2
only_tradable: true
risk_degree: 0.95
hmm_pause_pct: 0.70
drawdown_pause_pct: 8.0
liquidity_floor_adv: 5000000
"""
from __future__ import annotations
import copy
from typing import Dict
import numpy as np
import pandas as pd
from qlib.backtest.decision import TradeDecisionWO
from qlib.backtest.position import Position
from qlib.contrib.strategy.signal_strategy import TopkDropoutStrategy
__all__ = ["HmmRiskTopk"]
class HmmRiskTopk(TopkDropoutStrategy):
"""TopkDropoutStrategy with HMM high-vol pause + drawdown pause + liquidity floor."""
def __init__(
self,
*,
hmm_pause_pct: float = 0.70,
drawdown_pause_pct: float = 8.0,
liquidity_floor_adv: float = 0.0,
**kwargs,
):
super().__init__(**kwargs)
self.hmm_pause_pct = float(hmm_pause_pct)
self.drawdown_pause_pct = float(drawdown_pause_pct)
self.liquidity_floor_adv = float(liquidity_floor_adv)
self._peak_equity = 0.0
# ------------------------------------------------------------- gates
def _hmm_high_vol(self, pred_date) -> bool:
"""Cross-sectional mean HMM high-vol regime probability >= threshold."""
try:
from qlib.data import D
feat = D.features(D.instruments("all"), ["$sp_hmm_p_regime1"],
start_time=pred_date, end_time=pred_date)
if feat is None or len(feat) == 0:
return False
p = feat["$sp_hmm_p_regime1"].dropna()
if len(p) == 0:
return False
return float(p.mean()) >= self.hmm_pause_pct
except Exception:
return False
def _drawdown_active(self, equity: float) -> bool:
if self.drawdown_pause_pct <= 0:
return False
self._peak_equity = max(self._peak_equity, equity)
if self._peak_equity <= 0:
return False
dd = (self._peak_equity - equity) / self._peak_equity * 100.0
return dd >= self.drawdown_pause_pct
def _illiquid(self, codes, asof) -> Dict[str, bool]:
if self.liquidity_floor_adv <= 0 or not codes:
return {}
from tac_qlib.risk_limits import dollar_adv
adv = dollar_adv(codes, market="US", asof=asof, lookback=20)
return {c: adv.get(str(c).upper(), 0.0) < self.liquidity_floor_adv for c in codes}
# ------------------------------------------------------------- decision
def generate_trade_decision(self, execute_result=None):
trade_step = self.trade_calendar.get_trade_step()
trade_start_time, trade_end_time = self.trade_calendar.get_step_time(trade_step)
pred_start_time, pred_end_time = self.trade_calendar.get_step_time(trade_step, shift=1)
pred_score = self.signal.get_signal(start_time=pred_start_time, end_time=pred_end_time)
if pred_score is None:
return TradeDecisionWO([], self)
if isinstance(pred_score, pd.DataFrame):
pred_score = pred_score.iloc[:, 0]
current_temp = copy.deepcopy(self.trade_position)
assert isinstance(current_temp, Position)
held = {c for c in current_temp.get_stock_list() if abs(current_temp.get_stock_amount(c)) > 1e-6}
equity = current_temp.get_cash()
for code in held:
mark = self.trade_exchange.get_deal_price(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time, direction=1
)
if mark is not None and np.isfinite(mark):
equity += abs(current_temp.get_stock_amount(code)) * mark
hmm_pause = self._hmm_high_vol(str(pd.Timestamp(pred_start_time).date()))
dd_pause = self._drawdown_active(equity)
buys_paused = hmm_pause or dd_pause
pred_score = pred_score.copy()
if buys_paused or self.liquidity_floor_adv > 0:
new_codes = [c for c in pred_score.index if c not in held]
illiquid = self._illiquid(new_codes, str(pd.Timestamp(pred_start_time).date()))
for code in new_codes:
if buys_paused or illiquid.get(code, False):
pred_score[code] = -1e9 # cannot enter today
return super().generate_trade_decision(execute_result)
@@ -0,0 +1,201 @@
"""Fractional-Kelly dropout strategy for cross-sectional signals.
Sizing rule variant of ``qlib.contrib.strategy.signal_strategy.TopkDropoutStrategy``:
the topk/n_drop SELECTION is identical to the reference, but the buy size is
proportional to the score MAGNITUDE (edge) instead of equal-weight, capped at a
fraction ``cap_frac`` of the equal-weight notional so a single name cannot
over-concentrate the book.
``cap_frac`` is the fraction of the equal-weight per-name notional that a top
signal can deploy at most (e.g. 0.5 = at most half the equal-weight size).
Names whose score is below the median of the buy set get a proportionally
smaller slice; the residual stays in cash (that is the point of the rule:
throw away less edge per name, deploy less capital when conviction is low).
"""
from __future__ import annotations
from typing import List
import numpy as np
import pandas as pd
from qlib.backtest import Order
from qlib.backtest.decision import OrderDir, TradeDecisionWO
from qlib.contrib.strategy.signal_strategy import TopkDropoutStrategy
__all__ = ["FractionalKellyDropoutStrategy"]
DEFAULT_CAP_FRAC = 0.5
class FractionalKellyDropoutStrategy(TopkDropoutStrategy):
"""TopkDropout selection with score-magnitude (fractional-Kelly) sizing.
Parameters
----------
topk, n_drop, method_sell, method_buy, hold_thresh, only_tradable,
forbid_all_trade_at_limit : same as ``TopkDropoutStrategy``.
cap_frac : max buy notional as a fraction of the equal-weight notional.
"""
def __init__(self, *, topk, n_drop, cap_frac: float = DEFAULT_CAP_FRAC, **kwargs):
super().__init__(topk=topk, n_drop=n_drop, **kwargs)
self.cap_frac = cap_frac
def generate_trade_decision(self, execute_result=None):
import copy
trade_step = self.trade_calendar.get_trade_step()
trade_start_time, trade_end_time = self.trade_calendar.get_step_time(trade_step)
pred_start_time, pred_end_time = self.trade_calendar.get_step_time(trade_step, shift=1)
pred_score = self.signal.get_signal(start_time=pred_start_time, end_time=pred_end_time)
if isinstance(pred_score, pd.DataFrame):
pred_score = pred_score.iloc[:, 0]
if pred_score is None:
return TradeDecisionWO([], self)
if self.only_tradable:
def get_first_n(li, n, reverse=False):
cur_n = 0
res = []
for si in reversed(li) if reverse else li:
if self.trade_exchange.is_stock_tradable(
stock_id=si, start_time=trade_start_time, end_time=trade_end_time
):
res.append(si)
cur_n += 1
if cur_n >= n:
break
return res[::-1] if reverse else res
def get_last_n(li, n):
return get_first_n(li, n, reverse=True)
def filter_stock(li):
return [
si
for si in li
if self.trade_exchange.is_stock_tradable(
stock_id=si, start_time=trade_start_time, end_time=trade_end_time
)
]
else:
def get_first_n(li, n):
return list(li)[:n]
def get_last_n(li, n):
return list(li)[-n:]
def filter_stock(li):
return li
current_temp: "object" = copy.deepcopy(self.trade_position)
sell_order_list: List[Order] = []
buy_order_list: List[Order] = []
cash = current_temp.get_cash()
current_stock_list = current_temp.get_stock_list()
last = pred_score.reindex(current_stock_list).sort_values(ascending=False).index
if self.method_buy == "top":
today = get_first_n(
pred_score[~pred_score.index.isin(last)].sort_values(ascending=False).index,
self.n_drop + self.topk - len(last),
)
elif self.method_buy == "random":
topk_candi = get_first_n(pred_score.sort_values(ascending=False).index, self.topk)
candi = list(filter(lambda x: x not in last, topk_candi))
n = self.n_drop + self.topk - len(last)
try:
today = np.random.choice(candi, n, replace=False)
except ValueError:
today = candi
else:
raise NotImplementedError(f"This type of input is not supported")
comb = pred_score.reindex(last.union(pd.Index(today))).sort_values(ascending=False).index
if self.method_sell == "bottom":
sell = last[last.isin(get_last_n(comb, self.n_drop))]
elif self.method_sell == "random":
candi = filter_stock(last)
try:
sell = pd.Index(np.random.choice(candi, self.n_drop, replace=False) if len(last) else [])
except ValueError:
sell = candi
else:
raise NotImplementedError(f"This type of input is not supported")
buy = today[: len(sell) + self.topk - len(last)]
for code in current_stock_list:
if not self.trade_exchange.is_stock_tradable(
stock_id=code,
start_time=trade_start_time,
end_time=trade_end_time,
direction=None if self.forbid_all_trade_at_limit else OrderDir.SELL,
):
continue
if code in sell:
time_per_step = self.trade_calendar.get_freq()
if current_temp.get_stock_count(code, bar=time_per_step) < self.hold_thresh:
continue
sell_amount = current_temp.get_stock_amount(code=code)
sell_order = Order(
stock_id=code,
amount=sell_amount,
start_time=trade_start_time,
end_time=trade_end_time,
direction=Order.SELL,
)
if self.trade_exchange.check_order(sell_order):
sell_order_list.append(sell_order)
trade_val, trade_cost, trade_price = self.trade_exchange.deal_order(
sell_order, position=current_temp
)
cash += trade_val - trade_cost
if len(buy) == 0:
return TradeDecisionWO(sell_order_list, self)
# ---- fractional-Kelly sizing --------------------------------------
# equal-weight notional (reference baseline)
eq_notional = cash * self.risk_degree / len(buy)
buy_scores = pred_score.reindex(buy).astype(float)
lo, hi = buy_scores.min(), buy_scores.max()
if hi == lo:
w = pd.Series(1.0, index=buy_scores.index)
else:
w = (buy_scores - lo) / (hi - lo) # [0,1] edge magnitude
w = w.clip(lower=0.0)
w_max = w.max()
w = w / w_max if w_max > 0 else w # max == 1.0
for code in buy:
if not self.trade_exchange.is_stock_tradable(
stock_id=code,
start_time=trade_start_time,
end_time=trade_end_time,
direction=None if self.forbid_all_trade_at_limit else OrderDir.BUY,
):
continue
buy_price = self.trade_exchange.get_deal_price(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time, direction=OrderDir.BUY
)
notional = eq_notional * min(self.cap_frac, float(w.get(code, 0.0)))
buy_amount = notional / buy_price
factor = self.trade_exchange.get_factor(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time
)
buy_amount = self.trade_exchange.round_amount_by_trade_unit(buy_amount, factor)
buy_order = Order(
stock_id=code,
amount=buy_amount,
start_time=trade_start_time,
end_time=trade_end_time,
direction=Order.BUY,
)
buy_order_list.append(buy_order)
return TradeDecisionWO(sell_order_list + buy_order_list, self)
@@ -1,91 +0,0 @@
"""TopkDropout with a 1-day momentum entry-confirmation gate.
Gates NEW entries on short-term momentum: a name that is not currently held
may only be bought when its trailing 1-day return is above ``min_momentum``
(Lag-1 autocorr ~ +0.45 in the time-series study => short-term momentum
continuation). Held names are never force-sold by this gate — exits stay the
pure TopkDropout rule.
Implementation: override ``generate_trade_decision`` and zero out the signal
score of any non-held name that fails the momentum check BEFORE calling the
base TopkDropout decision, so it can never be selected as a buy candidate.
This is a clean pre-filter: the rest of the strategy (top-k, n_drop, sizing,
costs) is untouched.
Wired into a workflow yaml like:
strategy:
class: MomentumGateTopk
module_path: tac_qlib.contrib.strategy.momentum_gate
kwargs:
signal: "<PRED>"
topk: 10
n_drop: 2
only_tradable: true
risk_degree: 0.95
min_momentum: 0.0
"""
from __future__ import annotations
import copy
import pandas as pd
from qlib.backtest.decision import TradeDecisionWO
from qlib.backtest.position import Position
from qlib.contrib.strategy.signal_strategy import TopkDropoutStrategy
__all__ = ["MomentumGateTopk"]
class MomentumGateTopk(TopkDropoutStrategy):
"""TopkDropoutStrategy gated on 1-day momentum for new entries."""
def __init__(self, *, min_momentum: float = 0.0, **kwargs):
super().__init__(**kwargs)
self.min_momentum = float(min_momentum)
def _momentum_ok(self, code, trade_start, trade_end) -> bool:
"""True when the trailing 1-day return is above the momentum floor."""
try:
cur = self.trade_exchange.get_deal_price(
stock_id=code, start_time=trade_start, end_time=trade_end, direction=1
)
except Exception:
return False
if cur is None or cur != cur or cur <= 0:
return False
prev_start = trade_start - pd.Timedelta(days=5)
prev_end = trade_start - pd.Timedelta(seconds=1)
prev = self.trade_exchange.get_deal_price(
stock_id=code, start_time=prev_start, end_time=prev_end, direction=0
)
if prev is None or prev != prev or prev <= 0:
return False
return (cur / prev - 1.0) >= self.min_momentum
def generate_trade_decision(self, execute_result=None):
trade_step = self.trade_calendar.get_trade_step()
trade_start_time, trade_end_time = self.trade_calendar.get_step_time(trade_step)
pred_start_time, pred_end_time = self.trade_calendar.get_step_time(trade_step, shift=1)
pred_score = self.signal.get_signal(start_time=pred_start_time, end_time=pred_end_time)
if pred_score is None:
return TradeDecisionWO([], self)
if isinstance(pred_score, pd.DataFrame):
pred_score = pred_score.iloc[:, 0]
current_temp = copy.deepcopy(self.trade_position)
assert isinstance(current_temp, Position)
held = set(current_temp.get_stock_list())
held = {c for c in held if abs(current_temp.get_stock_amount(c)) > 1e-6}
# pre-filter: zero the score of non-held names that fail momentum
pred_score = pred_score.copy()
for code in pred_score.index:
if code in held:
continue # never gate exits / re-balancing of held names
if not self._momentum_ok(code, trade_start_time, trade_end_time):
pred_score[code] = -1e9 # cannot enter today
return super().generate_trade_decision(execute_result)
@@ -0,0 +1,231 @@
"""HMM-regime overlay TopkDropout strategy.
Regime-gate overlay on ``qlib.contrib.strategy.signal_strategy.TopkDropoutStrategy``:
selection and sizing are identical to the reference, but a name is only BOUGHT
(entry gate) when its per-symbol HMM regime posterior ``sp_hmm_p_regime1`` on
the signal date is >= ``regime_threshold``; otherwise it is held in cash instead
of being opened.
The regime posterior is read from the lake feature provider on the fly via
``qlib.data.D.features`` (field ``$sp_hmm_p_regime1``) for the signal window, so
no regime column needs to enter the model's ``feature_fields`` — the gate is a
pure overlay (book ch.01: regime flags regressed as model features, survived
only as an overlay). The HMM itself was fit with ``fit_end=<train end>`` when
the lake features were backfilled, so there is no lookahead.
Names already held are NOT force-sold when the regime turns unfavourable
(entry gate only, matching the queue-10 design).
"""
from __future__ import annotations
from typing import List
import numpy as np
import pandas as pd
from qlib.backtest import Order
from qlib.backtest.decision import OrderDir, TradeDecisionWO
from qlib.contrib.strategy.signal_strategy import TopkDropoutStrategy
try:
from qlib.data import D
except ImportError: # pragma: no cover - qlib always present in this stack
D = None
__all__ = ["RegimeGateDropoutStrategy"]
DEFAULT_REGIME_THRESHOLD = 0.5
REGIME_FIELD = "$sp_hmm_p_regime1"
class RegimeGateDropoutStrategy(TopkDropoutStrategy):
"""TopkDropout with an HMM-regime entry gate on buy candidates.
Parameters
----------
topk, n_drop, method_sell, method_buy, hold_thresh, only_tradable,
forbid_all_trade_at_limit : same as ``TopkDropoutStrategy``.
regime_threshold : minimum ``sp_hmm_p_regime1`` posterior required to open a
new position (default 0.5).
"""
def __init__(self, *, topk, n_drop, regime_threshold: float = DEFAULT_REGIME_THRESHOLD, **kwargs):
super().__init__(topk=topk, n_drop=n_drop, **kwargs)
self.regime_threshold = regime_threshold
def _regime_for(self, codes, pred_start, pred_end) -> pd.Series:
"""Return {code: sp_hmm_p_regime1} for the signal window (last day)."""
if D is None:
return pd.Series(dtype=float)
try:
df = D.features(list(codes), [REGIME_FIELD], start_time=pred_start, end_time=pred_end, freq="day")
except Exception: # noqa: BLE001 - a regime read failure should gate open, not crash
return pd.Series(dtype=float)
if df is None or len(df) == 0:
return pd.Series(dtype=float)
# df index is MultiIndex (datetime, instrument); take the last day's values
df = df.reset_index()
ts_col = "datetime" if "datetime" in df.columns else df.columns[0]
sym_col = "instrument" if "instrument" in df.columns else df.columns[1]
last_ts = df[ts_col].max()
last = df[df[ts_col] == last_ts]
out = {}
for _, row in last.iterrows():
sym = str(row[sym_col]).split("/")[-1].upper()
val = row.iloc[-1]
out[sym] = float(val) if val == val else np.nan
return pd.Series(out)
def generate_trade_decision(self, execute_result=None):
import copy
trade_step = self.trade_calendar.get_trade_step()
trade_start_time, trade_end_time = self.trade_calendar.get_step_time(trade_step)
pred_start_time, pred_end_time = self.trade_calendar.get_step_time(trade_step, shift=1)
pred_score = self.signal.get_signal(start_time=pred_start_time, end_time=pred_end_time)
if isinstance(pred_score, pd.DataFrame):
pred_score = pred_score.iloc[:, 0]
if pred_score is None:
return TradeDecisionWO([], self)
if self.only_tradable:
def get_first_n(li, n, reverse=False):
cur_n = 0
res = []
for si in reversed(li) if reverse else li:
if self.trade_exchange.is_stock_tradable(
stock_id=si, start_time=trade_start_time, end_time=trade_end_time
):
res.append(si)
cur_n += 1
if cur_n >= n:
break
return res[::-1] if reverse else res
def get_last_n(li, n):
return get_first_n(li, n, reverse=True)
def filter_stock(li):
return [
si
for si in li
if self.trade_exchange.is_stock_tradable(
stock_id=si, start_time=trade_start_time, end_time=trade_end_time
)
]
else:
def get_first_n(li, n):
return list(li)[:n]
def get_last_n(li, n):
return list(li)[-n:]
def filter_stock(li):
return li
current_temp: "object" = copy.deepcopy(self.trade_position)
sell_order_list: List[Order] = []
buy_order_list: List[Order] = []
cash = current_temp.get_cash()
current_stock_list = current_temp.get_stock_list()
last = pred_score.reindex(current_stock_list).sort_values(ascending=False).index
if self.method_buy == "top":
today = get_first_n(
pred_score[~pred_score.index.isin(last)].sort_values(ascending=False).index,
self.n_drop + self.topk - len(last),
)
elif self.method_buy == "random":
topk_candi = get_first_n(pred_score.sort_values(ascending=False).index, self.topk)
candi = list(filter(lambda x: x not in last, topk_candi))
n = self.n_drop + self.topk - len(last)
try:
today = np.random.choice(candi, n, replace=False)
except ValueError:
today = candi
else:
raise NotImplementedError(f"This type of input is not supported")
comb = pred_score.reindex(last.union(pd.Index(today))).sort_values(ascending=False).index
if self.method_sell == "bottom":
sell = last[last.isin(get_last_n(comb, self.n_drop))]
elif self.method_sell == "random":
candi = filter_stock(last)
try:
sell = pd.Index(np.random.choice(candi, self.n_drop, replace=False) if len(last) else [])
except ValueError:
sell = candi
else:
raise NotImplementedError(f"This type of input is not supported")
buy = today[: len(sell) + self.topk - len(last)]
# ---- regime gate -----------------------------------------------------
if buy:
regime = self._regime_for(buy, pred_start_time, pred_end_time)
gated = [c for c in buy if regime.get(c, np.nan) >= self.regime_threshold]
else:
gated = []
for code in current_stock_list:
if not self.trade_exchange.is_stock_tradable(
stock_id=code,
start_time=trade_start_time,
end_time=trade_end_time,
direction=None if self.forbid_all_trade_at_limit else OrderDir.SELL,
):
continue
if code in sell:
time_per_step = self.trade_calendar.get_freq()
if current_temp.get_stock_count(code, bar=time_per_step) < self.hold_thresh:
continue
sell_amount = current_temp.get_stock_amount(code=code)
sell_order = Order(
stock_id=code,
amount=sell_amount,
start_time=trade_start_time,
end_time=trade_end_time,
direction=Order.SELL,
)
if self.trade_exchange.check_order(sell_order):
sell_order_list.append(sell_order)
trade_val, trade_cost, trade_price = self.trade_exchange.deal_order(
sell_order, position=current_temp
)
cash += trade_val - trade_cost
if len(gated) == 0:
return TradeDecisionWO(sell_order_list, self)
value = cash * self.risk_degree / len(gated)
for code in gated:
if not self.trade_exchange.is_stock_tradable(
stock_id=code,
start_time=trade_start_time,
end_time=trade_end_time,
direction=None if self.forbid_all_trade_at_limit else OrderDir.BUY,
):
continue
buy_price = self.trade_exchange.get_deal_price(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time, direction=OrderDir.BUY
)
buy_amount = value / buy_price
factor = self.trade_exchange.get_factor(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time
)
buy_amount = self.trade_exchange.round_amount_by_trade_unit(buy_amount, factor)
buy_order = Order(
stock_id=code,
amount=buy_amount,
start_time=trade_start_time,
end_time=trade_end_time,
direction=Order.BUY,
)
buy_order_list.append(buy_order)
return TradeDecisionWO(sell_order_list + buy_order_list, self)
@@ -0,0 +1,169 @@
"""Market-neutral top/bottom long-short strategy for cross-sectional signals.
Captures the cross-sectional long-short spread net of costs: buys the top-ranked
``topk`` names and shorts the bottom-ranked ``topk`` names, equal-weight per
side, sized to ``risk_degree`` of total value per side. Rebalances daily to the
current rank (dropout-free: the book converges to the latest top/bottom sets).
The long and short legs use equal notional per side (gross exposure ~2x
``risk_degree`` of NAV, i.e. approximately market neutral before transaction
costs). Benchmark neutrality (SPY beta ~ 0) is the secondary sanity metric.
"""
from __future__ import annotations
from typing import List
import copy
import pandas as pd
from qlib.backtest import Order
from qlib.backtest.decision import OrderDir, TradeDecisionWO
from qlib.contrib.strategy.signal_strategy import BaseSignalStrategy
__all__ = ["TopBottomDropoutStrategy"]
DEFAULT_SHORT_LEG = True
DEFAULT_REBALANCE_DAILY = True
class TopBottomDropoutStrategy(BaseSignalStrategy):
"""Long top-k / short bottom-k equal-weight market-neutral book.
Parameters
----------
topk : number of names on each side (long top-k and short bottom-k).
short_leg : whether to open the short side (if False, long-only topk).
rebalance_daily : if True rebalance to current rank every day; else keep
positions and only refresh on score changes (dropout-style).
risk_degree : fraction of total value deployed per side.
"""
def __init__(
self,
*,
topk: int = 10,
short_leg: bool = DEFAULT_SHORT_LEG,
rebalance_daily: bool = DEFAULT_REBALANCE_DAILY,
**kwargs,
):
super().__init__(**kwargs)
self.topk = topk
self.short_leg = short_leg
self.rebalance_daily = rebalance_daily
self._prev_longs = set()
self._prev_shorts = set()
def generate_trade_decision(self, execute_result=None):
trade_step = self.trade_calendar.get_trade_step()
trade_start_time, trade_end_time = self.trade_calendar.get_step_time(trade_step)
pred_start_time, pred_end_time = self.trade_calendar.get_step_time(trade_step, shift=1)
pred_score = self.signal.get_signal(start_time=pred_start_time, end_time=pred_end_time)
if isinstance(pred_score, pd.DataFrame):
pred_score = pred_score.iloc[:, 0]
if pred_score is None or len(pred_score) == 0:
return TradeDecisionWO([], self)
# rank all names; topk longs and topk shorts
ranked = pred_score.sort_values(ascending=False)
longs = list(ranked.index[: self.topk])
shorts = list(ranked.index[-self.topk :]) if self.short_leg else []
current_temp: "object" = copy.deepcopy(self.trade_position)
current_codes = set(current_temp.get_stock_list())
holdings = {c: current_temp for c in current_codes if abs(current_temp.get_stock_amount(c)) > 1e-6}
sell_orders: List[Order] = []
buy_orders: List[Order] = []
def _tradable(code, direction):
try:
return self.trade_exchange.is_stock_tradable(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time, direction=direction
)
except TypeError:
return self.trade_exchange.is_stock_tradable(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time
)
# determine target set (long/short)
target_longs = set(longs)
target_shorts = set(shorts)
# close positions not in the target book
for code in list(holdings):
if code in target_longs or code in target_shorts:
continue
amt = abs(current_temp.get_stock_amount(code))
o = Order(
stock_id=code,
amount=amt,
start_time=trade_start_time,
end_time=trade_end_time,
direction=Order.SELL if code in target_longs else Order.SELL,
)
if self.trade_exchange.check_order(o):
sell_orders.append(o)
self.trade_exchange.deal_order(o, position=current_temp)
# equal-weight notional per side
total_value = current_temp.get_cash()
for code, pos in holdings.items():
if code in target_longs or code in target_shorts:
mark = self.trade_exchange.get_deal_price(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time, direction=Order.SELL
)
if mark is not None and mark == mark:
total_value += abs(current_temp.get_stock_amount(code)) * mark
side_notional = total_value * self.risk_degree / max(1, self.topk)
for code in longs:
if code in holdings and abs(current_temp.get_stock_amount(code)) > 1e-6:
continue
px = self.trade_exchange.get_deal_price(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time, direction=Order.BUY
)
if px is None or px != px or px <= 0:
continue
amount = side_notional / px
factor = self.trade_exchange.get_factor(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time
)
amount = self.trade_exchange.round_amount_by_trade_unit(amount, factor)
o = Order(
stock_id=code,
amount=amount,
start_time=trade_start_time,
end_time=trade_end_time,
direction=Order.BUY,
)
if self.trade_exchange.check_order(o):
buy_orders.append(o)
if self.short_leg:
for code in shorts:
if code in holdings and abs(current_temp.get_stock_amount(code)) > 1e-6:
continue
px = self.trade_exchange.get_deal_price(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time, direction=Order.SELL
)
if px is None or px != px or px <= 0:
continue
amount = side_notional / px
factor = self.trade_exchange.get_factor(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time
)
amount = self.trade_exchange.round_amount_by_trade_unit(amount, factor)
o = Order(
stock_id=code,
amount=amount,
start_time=trade_start_time,
end_time=trade_end_time,
direction=Order.SELL,
)
if self.trade_exchange.check_order(o):
sell_orders.append(o)
return TradeDecisionWO(sell_orders + buy_orders, self)
@@ -0,0 +1,202 @@
"""Weekly-rebalance TopkDropout strategy.
Turnover-reduction variant of ``qlib.contrib.strategy.signal_strategy.TopkDropoutStrategy``:
the topk/n_drop selection and sizing are identical to the reference, but the
target book is recomputed only on the first trading day of each ISO week; on the
other days the strategy issues NO orders (holds the book untouched).
The weekly cadence is derived from the qlib trade calendar: a rebalance happens
when the current trade step's date belongs to a different ISO ``(year, week)``
than the previous trade step. ``hold_band_pct`` (default 0) optionally skips
tiny rebalances: when a name's existing position differs from the new target by
less than this fraction, no order is generated for it.
"""
from __future__ import annotations
from typing import List
import numpy as np
import pandas as pd
from qlib.backtest import Order
from qlib.backtest.decision import OrderDir, TradeDecisionWO
from qlib.contrib.strategy.signal_strategy import TopkDropoutStrategy
__all__ = ["WeeklyRebalanceDropoutStrategy"]
DEFAULT_HOLD_BAND_PCT = 0.0
class WeeklyRebalanceDropoutStrategy(TopkDropoutStrategy):
"""TopkDropout rebalanced once per ISO week; holds otherwise.
Parameters
----------
topk, n_drop, method_sell, method_buy, hold_thresh, only_tradable,
forbid_all_trade_at_limit : same as ``TopkDropoutStrategy``.
hold_band_pct : skip order for a name whose deviation from target weight is
below this fraction of the target (no-trade buffer band).
"""
def __init__(self, *, topk, n_drop, hold_band_pct: float = DEFAULT_HOLD_BAND_PCT, **kwargs):
super().__init__(topk=topk, n_drop=n_drop, **kwargs)
self.hold_band_pct = hold_band_pct
@staticmethod
def _iso_week(ts) -> tuple:
return (ts.year, ts.week)
def generate_trade_decision(self, execute_result=None):
import copy
trade_step = self.trade_calendar.get_trade_step()
trade_start_time, trade_end_time = self.trade_calendar.get_step_time(trade_step)
cur_week = self._iso_week(trade_start_time)
prev_week = getattr(self, "_last_week", None)
self._last_week = cur_week
if prev_week is not None and prev_week == cur_week:
# not the first trading day of this ISO week -> hold
return TradeDecisionWO([], self)
pred_start_time, pred_end_time = self.trade_calendar.get_step_time(trade_step, shift=1)
pred_score = self.signal.get_signal(start_time=pred_start_time, end_time=pred_end_time)
if isinstance(pred_score, pd.DataFrame):
pred_score = pred_score.iloc[:, 0]
if pred_score is None:
return TradeDecisionWO([], self)
if self.only_tradable:
def get_first_n(li, n, reverse=False):
cur_n = 0
res = []
for si in reversed(li) if reverse else li:
if self.trade_exchange.is_stock_tradable(
stock_id=si, start_time=trade_start_time, end_time=trade_end_time
):
res.append(si)
cur_n += 1
if cur_n >= n:
break
return res[::-1] if reverse else res
def get_last_n(li, n):
return get_first_n(li, n, reverse=True)
def filter_stock(li):
return [
si
for si in li
if self.trade_exchange.is_stock_tradable(
stock_id=si, start_time=trade_start_time, end_time=trade_end_time
)
]
else:
def get_first_n(li, n):
return list(li)[:n]
def get_last_n(li, n):
return list(li)[-n:]
def filter_stock(li):
return li
current_temp: "object" = copy.deepcopy(self.trade_position)
sell_order_list: List[Order] = []
buy_order_list: List[Order] = []
cash = current_temp.get_cash()
current_stock_list = current_temp.get_stock_list()
last = pred_score.reindex(current_stock_list).sort_values(ascending=False).index
if self.method_buy == "top":
today = get_first_n(
pred_score[~pred_score.index.isin(last)].sort_values(ascending=False).index,
self.n_drop + self.topk - len(last),
)
elif self.method_buy == "random":
topk_candi = get_first_n(pred_score.sort_values(ascending=False).index, self.topk)
candi = list(filter(lambda x: x not in last, topk_candi))
n = self.n_drop + self.topk - len(last)
try:
today = np.random.choice(candi, n, replace=False)
except ValueError:
today = candi
else:
raise NotImplementedError(f"This type of input is not supported")
comb = pred_score.reindex(last.union(pd.Index(today))).sort_values(ascending=False).index
if self.method_sell == "bottom":
sell = last[last.isin(get_last_n(comb, self.n_drop))]
elif self.method_sell == "random":
candi = filter_stock(last)
try:
sell = pd.Index(np.random.choice(candi, self.n_drop, replace=False) if len(last) else [])
except ValueError:
sell = candi
else:
raise NotImplementedError(f"This type of input is not supported")
buy = today[: len(sell) + self.topk - len(last)]
for code in current_stock_list:
if not self.trade_exchange.is_stock_tradable(
stock_id=code,
start_time=trade_start_time,
end_time=trade_end_time,
direction=None if self.forbid_all_trade_at_limit else OrderDir.SELL,
):
continue
if code in sell:
time_per_step = self.trade_calendar.get_freq()
if current_temp.get_stock_count(code, bar=time_per_step) < self.hold_thresh:
continue
sell_amount = current_temp.get_stock_amount(code=code)
sell_order = Order(
stock_id=code,
amount=sell_amount,
start_time=trade_start_time,
end_time=trade_end_time,
direction=Order.SELL,
)
if self.trade_exchange.check_order(sell_order):
sell_order_list.append(sell_order)
trade_val, trade_cost, trade_price = self.trade_exchange.deal_order(
sell_order, position=current_temp
)
cash += trade_val - trade_cost
if len(buy) == 0:
return TradeDecisionWO(sell_order_list, self)
value = cash * self.risk_degree / len(buy)
for code in buy:
if not self.trade_exchange.is_stock_tradable(
stock_id=code,
start_time=trade_start_time,
end_time=trade_end_time,
direction=None if self.forbid_all_trade_at_limit else OrderDir.BUY,
):
continue
buy_price = self.trade_exchange.get_deal_price(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time, direction=OrderDir.BUY
)
buy_amount = value / buy_price
factor = self.trade_exchange.get_factor(
stock_id=code, start_time=trade_start_time, end_time=trade_end_time
)
buy_amount = self.trade_exchange.round_amount_by_trade_unit(buy_amount, factor)
buy_order = Order(
stock_id=code,
amount=buy_amount,
start_time=trade_start_time,
end_time=trade_end_time,
direction=Order.BUY,
)
buy_order_list.append(buy_order)
return TradeDecisionWO(sell_order_list + buy_order_list, self)
+29 -2
View File
@@ -6,10 +6,11 @@ The lake is a hive-partitioned parquet store (see ``tac-engine/skills/tradeac-la
├── market=US/
│ └── timeframe=1d/
│ └── symbol=AAPL.parquet # OHLCV bars: t, date, o, h, l, c, v, n, vw
├── features/ # ta-lib indicators, wide format
├── features/ # indicators, wide format, family tier
│ └── market=US/
│ └── timeframe=1d/
│ └── symbol=AAPL.parquet # t, sma_5, sma_20, rsi_14, ...
│ ├── family=ta/symbol=AAPL.parquet # t, sma_5, sma_20, rsi_14, ...
│ └── family=sp/symbol=AAPL.parquet # t, sp_ou_*, sp_hmm_*, ...
├── calendar.parquet # trading days per market
├── coverage.parquet # per (market,timeframe,symbol) loaded windows
└── symbols.parquet # asset master
@@ -107,8 +108,34 @@ class LakeConfig:
return self.lake_root / "features" / f"market={self.market}" / f"timeframe={timeframe}"
def features_path(self, timeframe: str, symbol: str) -> Path:
# Legacy flat path (no family tier). Prefer `load_features` which
# resolves the family=ta|sp partition layout.
return self.features_dir(timeframe) / f"symbol={str(symbol).upper()}.parquet"
def load_features(self, timeframe: str, symbol: str) -> pd.DataFrame:
"""All feature columns for a symbol, merging the `family=ta` and
`family=sp` partitions by timestamp. Returns an empty frame when no
feature files exist (legacy flat layout falls back transparently)."""
sym = str(symbol).upper()
frames = []
for family in ("ta", "sp"):
p = self.features_dir(timeframe) / f"family={family}" / f"symbol={sym}.parquet"
if p.exists():
frames.append(pd.read_parquet(p))
if not frames:
flat = self.features_dir(timeframe) / f"symbol={sym}.parquet"
if flat.exists():
return pd.read_parquet(flat)
return pd.DataFrame()
if len(frames) == 1:
return frames[0]
merged = frames[0]
for extra in frames[1:]:
merged = merged.merge(extra, on="t", how="outer", suffixes=("", "_dup"))
for c in [c for c in merged.columns if c.endswith("_dup")]:
merged = merged.drop(columns=c)
return merged
def calendar_path(self) -> Path:
return self.lake_root / "calendar.parquet"
+1 -2
View File
@@ -173,8 +173,7 @@ class LakeFeatureProvider(FeatureProvider):
def _load_feature_df(self, instrument: str, timeframe: str) -> pd.DataFrame:
key = (instrument, timeframe)
if key not in self._feature_cache:
p = self.cfg.features_path(timeframe, instrument)
self._feature_cache[key] = pd.read_parquet(p) if p.exists() else pd.DataFrame()
self._feature_cache[key] = self.cfg.load_features(timeframe, instrument)
return self._feature_cache[key]
@staticmethod