start experiment 22 (exp/22-re-run-experiment-16s-5-day-rankic-ensem)
This commit is contained in:
@@ -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):
|
||||
|
||||
Binary file not shown.
Binary file not shown.
Binary file not shown.
@@ -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"
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user