Modeling Infrastructure for Point-in-Time Validation and Walk-Forward Evaluation
Summary
This shared modeling module supports quantitative research workflows by loading features, labels, temporal data, and configuration, then preparing walk-forward cross-validation folds. Its design emphasizes reproducibility and data lineage: hashes capture source files, arrays, model settings, and notebook code so cached results can be invalidated when their inputs change. It also sets random seeds and includes utilities for dataset schemas, feature selection, and fold preparation.
Several details address common evaluation errors. The conformal quantile routine uses the finite-sample rank required for split-conformal intervals, returns an unbounded interval when the calibration set cannot support the requested coverage, and refuses unrankable missing scores. A cross-sectional information-coefficient adapter computes per-date Spearman correlations and averages only defined values, returning missing when no dates qualify rather than implying a measured zero. These are infrastructure methods, not evidence that any strategy works; their usefulness depends on correct upstream labels, point-in-time data, and suitable fold design.
Key ideas
- Cache signatures should bind fitted results to source code, input artifacts, and settings that affect the fit.
- Split-conformal intervals require the prescribed finite-sample calibration rank rather than an interpolated quantile.
- When the calibration sample cannot attain the requested coverage, an unbounded interval reflects the available evidence.
- Cross-sectional IC should be computed per date and averaged only over dates with defined correlations.
- Walk-forward folds and reproducible seeds help standardize model evaluation, while correct data and labels remain essential.
Tags
Full text
# modeling.py
```py
"""Shared modeling infrastructure for Ch11+ notebooks.
Provides:
- load_modeling_dataset(): Load features + temporal + labels, join, detect schema
- load_configs(): Load model configs for a label from label YAML + presets
- prepare_cv_folds(): Preprocess data into train/val folds (impute, scale)
The cross-sectional IC itself is the library's: call
``ml4t.diagnostic.metrics.cross_sectional_ic`` against a polars frame of
(date, symbol, y_true, y_pred) directly. ``cross_sectional_ic_mean()`` below is
the adapter for callers holding aligned numpy arrays instead, which is what a
scikit-learn or Optuna objective has in hand inside a fold.
Usage:
from utils.modeling import load_modeling_dataset, load_configs, prepare_cv_folds
mds = load_modeling_dataset("etfs", "fwd_ret_21d")
configs = load_configs("etfs", "fwd_ret_21d", family="linear")
folds = prepare_cv_folds(mds.dataset.to_pandas(), mds.splits, ...)
"""
from __future__ import annotations
import hashlib
import json
import math
import os
import random
import warnings
from collections.abc import Mapping, Sequence
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any, cast
import numpy as np
import pandas as pd
import polars as pl
import yaml
from utils.artifact_specs import (
load_feature_spec,
load_label_spec,
resolve_label_buffer,
resolve_label_buffer_unit,
resolve_label_horizon,
resolve_market_semantics,
resolve_storage_path,
)
from utils.cv_splits import generate_cv_splits, make_wf_config
RANDOM_SEED = 42
MIN_TEMPORAL_DATE_COVERAGE = 0.95 # Allow short calendar-edge gaps, not missing windows.
# Burn-in a temporal model cannot emit through, excused only at the start of a
# train window. Measured on crypto_perps_funding, whose GARCH and HMM features
# carry a 90-bar rolling z-score: 97/2024 dates (4.8%) for fold 0 and 119/1935
# (6.1%) for fold 1. A stale artifact whose fold IDs have shifted presents as a
# leading gap of roughly half the window, so this bound still rejects it.
MAX_TEMPORAL_WARMUP_FRACTION = 0.10
def file_sha256(path: Path) -> str:
"""The SHA-256 digest of an input artifact, for a cache key that names its inputs."""
with Path(path).open("rb") as handle:
return hashlib.file_digest(handle, "sha256").hexdigest()
def array_sha256(array: np.ndarray) -> str:
"""Hash an array's shape, dtype and contiguous values.
Covers what a file digest cannot: filtering, symbol limits, row order, feature
order and cleaning semantics, all applied after the source files were read.
"""
contiguous = np.ascontiguousarray(array)
digest = hashlib.sha256()
digest.update(str(contiguous.dtype).encode())
digest.update(repr(contiguous.shape).encode())
digest.update(contiguous.tobytes())
return digest.hexdigest()
def canonical_sha256(payload: Mapping) -> str:
"""Hash a nested contract with stable key ordering and date serialization."""
return hashlib.sha256(json.dumps(payload, sort_keys=True, default=str).encode()).hexdigest()
def notebook_cache_signature(
notebook_py: Path,
*,
inputs: Mapping[str, str],
settings: Mapping[str, object],
) -> dict:
"""The signature a notebook stores beside cached fit results, so a stale one is caught.
A cache keyed on nothing but the file's existence is a correctness hazard, not a
convenience: change the code that produced it and the notebook silently republishes
the previous run's numbers under the new source. The published page then claims
results its own code did not compute, and nothing in the provenance stamp catches
it - that stamp binds the .py to the .ipynb, and both are the new ones.
So the signature carries the notebook's own source digest alongside the input
digests and the settings that change a fit. Any edit to the notebook invalidates
it, which is deliberately coarser than tracking the functions that matter: a
needless re-run costs time, and a missed one costs a wrong number in the book.
"""
return {
"notebook_source_sha256": file_sha256(notebook_py),
"inputs": dict(inputs),
"settings": dict(settings),
}
def conformal_quantile(scores: np.ndarray, coverage: float) -> float:
"""The split-conformal interval half-width for ``coverage``, from calibration scores.
Split conformal prediction earns its finite-sample coverage guarantee by taking
the ``ceil((n + 1) * coverage)``-th smallest calibration score, and the ceiling
has to be taken literally. The value must be a score that is actually in the
calibration set, at that rank.
``np.quantile`` cannot express this, at either of its two obvious spellings.
Its default interpolates between the two neighbouring scores, returning
something narrower than the rank by exactly the finite-sample margin the
ceiling exists to supply. ``method="higher"`` rounds up to a real score but
still maps the level onto ``n - 1`` intervals rather than ``n`` ranks, so
passing ``k / n`` selects rank ``k + 1`` for every ``k < n`` - measured on
n=100 (rank 92 for k=91) and n=1000 (rank 902 for k=901). That errs wide
rather than narrow, so it does not break the guarantee, but it is not the
interval the method defines and it is not what the surrounding prose claims.
When ``k > n`` no calibration score attains the requested coverage - a
calibration set of 5 cannot certify 90% - and the interval is unbounded.
``method="higher"`` silently returns the largest score there, which asserts a
guarantee the data cannot support.
Returns ``inf`` in that case, so the interval it produces is unbounded and
obviously so.
An infinite score is kept and ranked, not filtered: it is a real
nonconformity score, it orders above every finite one, and dropping it would
shrink ``n`` and hand back a *narrower* interval than the calibration set
supports. A NaN score has no rank at all, so it is refused rather than
dropped, for the same reason: silently removing it changes the denominator
the guarantee is computed from.
"""
if not 0 < coverage < 1:
raise ValueError(f"coverage must lie strictly between 0 and 1, got {coverage}")
ranked = np.asarray(scores, dtype=float)
n_nan = int(np.isnan(ranked).sum())
if n_nan:
raise ValueError(
f"{n_nan} of {ranked.size} calibration scores are NaN and cannot be ranked. "
"A conformal quantile is a rank in the calibration set, so dropping them would "
"change the sample size the coverage guarantee is computed from."
)
n = ranked.size
if n == 0:
return float("inf")
rank = math.ceil((n + 1) * coverage)
if rank > n:
return float("inf")
return float(np.partition(ranked, rank - 1)[rank - 1])
def seed_everything(seed: int = RANDOM_SEED) -> None:
"""Set all random seeds for full reproducibility (CPU + GPU).
Must be called before any stochastic operation. For per-fold
reproducibility, call again at the start of each fold with
``seed_everything(RANDOM_SEED + fold_id)``.
"""
random.seed(seed)
np.random.seed(seed)
os.environ["PYTHONHASHSEED"] = str(seed)
try:
import torch
torch.manual_seed(seed)
if torch.cuda.is_available():
torch.cuda.manual_seed(seed)
torch.cuda.manual_seed_all(seed)
torch.backends.cudnn.deterministic = True
torch.backends.cudnn.benchmark = False
except ImportError:
pass # torch not needed for sklearn/numpy-only models (PCA, IPCA)
from utils.paths import get_case_study_dir
# Columns that are identifiers, not features
ID_COLS = {
"date",
"timestamp",
"asset",
"symbol",
"stock_id",
"product",
"position",
"instrument_id",
}
# Meta columns that may leak or are composites (not independent features)
META_LEAK = {
"underlying_price",
"instr_mid",
"instr_bid",
"instr_ask",
"ls_signal",
"risk_adj_score",
}
@dataclass
class ModelingDataset:
"""Container for a fully-joined modeling dataset with detected schema."""
dataset: pl.DataFrame
feature_names: list[str]
label_col: str
date_col: str
entity_cols: list[str]
join_cols: list[str]
splits: list[dict[str, Any]]
label_buffer: str
# Which case study this was loaded from. Callers that persist or cache anything derived from
# the dataset need it, and it is not otherwise recoverable from the frame.
case_study_id: str = ""
cv_config: Any = None # WalkForwardConfig (optional, avoids hard import dep)
task_type: str = "regression" # "regression" or "classification"
num_classes: int = 0 # 0 for regression, 2+ for classification
class_values: list = field(default_factory=list) # sorted unique values for classification
# Per-fold temporal features, carrying a 'fold' column. Held lazily where the loader
# can: on us_equities_panel this artifact is 68.7M rows, and materialising it cost
# 8.7 GB for the whole run when every consumer wants one fold. Reach it through
# ``fold_temporal_frame``, never by indexing it directly.
temporal_by_fold: pl.LazyFrame | pl.DataFrame | pd.DataFrame | None = None
temporal_keys: list[str] = field(default_factory=list) # Join keys for temporal features
temporal_feature_names: list[str] = field(default_factory=list) # Temporal feature column names
temporal_artifact_splits: list[dict[str, Any]] = field(default_factory=list)
# The float type the design matrices are built in. Read from the case study's
# ``features.storage_dtype``; fold preparation honours it rather than pinning float64.
feature_dtype: str = "float64"
# Continuous-return label that classification predictions are scored against.
# None for regression labels. When set, the column lives in ``dataset`` and
# downstream IC computation must use it instead of the binary ``label_col``.
eval_label_col: str | None = None
# The feature list the panel carries, independent of any ``columns`` projection the load
# was asked for. ``feature_names`` describes the frame in ``dataset`` and narrows with the
# projection; this one describes the artifacts and does not, which is what the recorded
# identity has to read so a caller narrowing its own load does not re-key runs it never
# touched. Equal to ``feature_names`` whenever nothing was projected. Empty only on a
# ``ModelingDataset`` built by hand rather than by ``load_modeling_dataset``, where
# ``input_lineage`` falls back to ``feature_names`` and behaves exactly as it did before.
panel_feature_names: list[str] = field(default_factory=list)
# Inputs ``input_lineage`` is derived from: the artifact paths and the
# universe reduction, which are not otherwise recoverable from this object.
lineage_inputs: dict[str, Any] = field(default_factory=dict, repr=False)
_input_lineage: dict[str, Any] | None = field(default=None, init=False, repr=False)
@property
def input_lineage(self) -> dict[str, Any]:
"""Identity-defining input lineage, computed on first use.
A training spec that persists results includes this payload so changed
artifacts or CV windows cannot reuse an old training hash. It digests
every input artifact, and those run to gigabytes - nasdaq100's feature
parquet alone is 7 GB and takes ~12 s to hash - while only the notebooks
that pass it to ``build_training_spec`` need it. Computing it here rather
than in ``load_modeling_dataset`` keeps that cost off every other caller.
"""
if self._input_lineage is None:
if not self.lineage_inputs:
raise ValueError(
"input_lineage is unavailable: this ModelingDataset was built without "
"lineage_inputs. Construct it via load_modeling_dataset()."
)
self._input_lineage = build_modeling_input_lineage(
artifacts=self.lineage_inputs["artifacts"],
feature_names=self.panel_feature_names or self.feature_names,
splits=self.splits,
label_buffer=self.label_buffer,
task_type=self.task_type,
eval_label_col=self.eval_label_col,
max_symbols=self.lineage_inputs["max_symbols"],
symbols=self.lineage_inputs["symbols"],
feature_dtype=self.feature_dtype,
)
return self._input_lineage
def _sha256_file(path: Path) -> str:
"""Digest an identity-defining input artifact's bytes.
Stable across repeated writes, not across encodings. Two parquet files holding
identical data hash differently if they were written with a different compression
codec or row-group size, and this digest is inside `computation.feature_artifacts` and
`computation.input_data_spec.artifacts`, which are inside the hashed `computation`
block - so a codec change forks the `training_hash` of every run reading that artifact.
Two things hold that shut. `artifact_digest._PARQUET_WRITE_SETTINGS` states the
encoding these artifacts are written under rather than inheriting a library default,
and `tests/test_artifact_digest_encoding.py` records the bytes a fixed frame produces
under it, so a change arrives as a failing test rather than as a registry that has
grown two identities for one piece of work.
`artifact_digest.value_digest` digests column *values* and is what a new identity
version should use here; changing it now would re-key all 1,305 registered training
runs across the nine live registries, which is a re-derivation rather than a fix.
"""
digest = hashlib.sha256()
with path.open("rb") as src:
for chunk in iter(lambda: src.read(1024 * 1024), b""):
digest.update(chunk)
return digest.hexdigest()
def build_modeling_input_lineage(
*,
artifacts: dict[str, Path],
feature_names: list[str],
splits: list[dict[str, Any]],
label_buffer: str,
task_type: str,
eval_label_col: str | None,
max_symbols: int,
symbols: list[str] | None,
feature_dtype: str = "float64",
) -> dict[str, Any]:
"""Build the portable input identity carried by persisted training runs.
``feature_dtype`` is part of the identity because the artifacts are not. The parquet files
a case study reads are unchanged by a precision declaration, so without this a result fitted
in double precision and one fitted in single resolve to the same training identity, and the
registry serves the older one for a spec that asked for the other.
It is written into the payload only when it is not ``float64``. A key added unconditionally
changes the fingerprint of every case study, including the eight that declared nothing, and
would invalidate every training run already registered against them. Omitting the default
keeps those fingerprints exactly as they were and gives only the declaring case study a new
one. Any future default must be added the same way, for the same reason.
"""
split_fields = ("fold", "train_start", "train_end", "val_start", "val_end")
def _normalize(key: str, value: Any) -> str:
# str() on a pd.Timestamp renders "2019-01-07 00:00:00" or "...+00:00"
# depending on whether the caller's boundaries are tz-aware, so the same
# window read two ways would fingerprint differently.
if key == "fold":
return str(value)
return pd.Timestamp(value).tz_localize(None).isoformat()
normalized_splits = [
{key: _normalize(key, split[key]) for key in split_fields if split.get(key) is not None}
for split in splits
]
payload: dict[str, Any] = {
"schema_version": 1,
"artifacts": {
name: {"sha256": _sha256_file(path), "size": path.stat().st_size}
for name, path in sorted(artifacts.items())
},
"feature_names": list(feature_names),
"splits": normalized_splits,
"label_buffer": label_buffer,
"task_type": task_type,
"eval_label_col": eval_label_col,
"max_symbols": int(max_symbols),
"symbols": sorted(symbols) if symbols else None,
}
if feature_dtype != "float64":
payload["feature_dtype"] = feature_dtype
canonical = json.dumps(payload, sort_keys=True, separators=(",", ":"))
payload["fingerprint"] = hashlib.sha256(canonical.encode()).hexdigest()
return payload
# ---------------------------------------------------------------------------
# CV / Protocol Configuration
# ---------------------------------------------------------------------------
@dataclass
class WalkForwardConfig:
"""Config for walk-forward cross-validation.
Compatible with ml4t.diagnostic.splitters.WalkForwardCV.
"""
n_splits: int
train_size: str
test_size: str
embargo_td: str
label_horizon: str
timestamp_col: str = "timestamp"
calendar_id: str | None = None
test_start: str | None = None
test_end: str | None = None
def model_dump(self) -> dict:
"""Return fields as a dict (Pydantic-compatible API)."""
from dataclasses import asdict
return asdict(self)
def to_json(self, path: str | Path) -> None:
"""Serialize config to a JSON file (Pydantic-compatible API)."""
import json
path = Path(path)
path.parent.mkdir(parents=True, exist_ok=True)
with open(path, "w") as f:
json.dump(self.model_dump(), f, indent=2)
@classmethod
def from_json(cls, path: str | Path) -> WalkForwardConfig:
"""Deserialize config from a JSON file."""
import json
with open(path) as f:
data = json.load(f)
return cls(**data)
def load_protocol(case_study_id: str) -> dict:
"""Load evaluation protocol from config/setup.yaml.
Returns dict with n_splits, train_size, test_size, holdout,
leakage_guards, calendar.
"""
path = get_case_study_dir(case_study_id) / "config" / "setup.yaml"
if not path.exists():
msg = f"Setup not found: {path}"
raise FileNotFoundError(msg)
with open(path) as f:
setup = yaml.safe_load(f)
ev = setup.get("evaluation", {})
labels = setup.get("labels", {})
market_semantics = resolve_market_semantics(case_study_id, setup)
protocol = {
"n_splits": ev.get("n_splits", 5),
"train_size": ev.get("train_size", "5Y"),
"test_size": ev.get("val_size", "1Y"),
"step_size": ev.get("step_size", "1Y"),
"calendar": market_semantics.get("calendar") or ev.get("calendar", "NYSE"),
"holdout": {
"start": ev.get("holdout_start"),
"end": ev.get("holdout_end"),
},
"leakage_guards": {},
}
# Parse label horizon from buffer string (e.g. "21D", "8h", "15T", "60m", "1M")
buffer = labels.get("buffer", "21D")
if buffer.endswith("D"):
protocol["leakage_guards"]["label_horizon_days"] = int(buffer[:-1])
elif buffer.endswith(("h", "H")):
protocol["leakage_guards"]["label_horizon_hours"] = int(buffer[:-1])
elif buffer.endswith(("T",)) or (buffer.endswith("m") and not buffer.endswith("M")):
protocol["leakage_guards"]["label_horizon_minutes"] = int(buffer[:-1])
elif buffer.endswith("M"):
# Monthly: approximate as 30 calendar days.
# pd.Timedelta rejects 'M' as ambiguous; ml4t-diagnostic needs a library fix
# to support calendar-month durations natively (see ml4t-diagnostic-dev/bugs/).
protocol["leakage_guards"]["label_horizon_days"] = int(buffer[:-1]) * 30
return protocol
def _label_horizon_to_iso(leakage_guards: dict) -> str:
"""Convert label horizon to ISO 8601 duration string."""
if "label_horizon_days" in leakage_guards:
return f"P{leakage_guards['label_horizon_days']}D"
if "label_horizon_hours" in leakage_guards:
return f"PT{leakage_guards['label_horizon_hours']}H"
if "label_horizon_minutes" in leakage_guards:
return f"PT{leakage_guards['label_horizon_minutes']}M"
return "P0D"
def get_cv_config(case_study_id: str) -> WalkForwardConfig:
"""Load walk-forward CV config from case study setup.yaml.
This is the single entry point for CV configuration. Reads
from ``case_studies/{id}/config/setup.yaml``.
"""
protocol = load_protocol(case_study_id)
embargo = _label_horizon_to_iso(protocol.get("leakage_guards", {}))
holdout = protocol.get("holdout", {})
return WalkForwardConfig(
n_splits=protocol["n_splits"],
train_size=protocol["train_size"],
test_size=protocol["test_size"],
embargo_td=embargo,
label_horizon=embargo,
calendar_id=protocol.get("calendar", "NYSE"),
test_start=str(holdout["start"]) if holdout.get("start") else None,
test_end=str(holdout["end"]) if holdout.get("end") else None,
)
# Classification label prefixes
_CLASSIFICATION_PREFIXES = ("fwd_dir_", "fwd_class_", "fwd_tb_", "fwd_carry_")
def detect_label_type(label_col: str, label_series: pl.Series) -> tuple[str, int, list]:
"""Detect regression vs classification from label name and values.
Returns (task_type, num_classes, class_values).
- task_type: "regression" or "classification"
- num_classes: 0 for regression, 2+ for classification
- class_values: sorted unique values for classification, [] for regression
"""
is_classification = any(label_col.startswith(p) for p in _CLASSIFICATION_PREFIXES)
# Also classify if integer dtype with few unique values
if not is_classification and label_series.dtype in (pl.Int8, pl.Int16, pl.Int32, pl.Int64):
if label_series.drop_nulls().n_unique() <= 10:
is_classification = True
if not is_classification:
return "regression", 0, []
unique_vals = sorted(label_series.drop_nulls().unique().to_list())
return "classification", len(unique_vals), unique_vals
def get_classification_eval_label(case_study_id: str, label: str) -> str:
"""Return the continuous return that a classification label is derived from.
Classification predictions (probabilities or class-expected-values) must be
scored against the underlying continuous return, not the binary/categorical
label itself: ``Spearman(score, fwd_continuous_return)`` is the proper IC,
while ``Spearman(score, binary_label)`` collapses to ``2·(AUC − 0.5)``.
The mapping is declared per-case-study under
``labels.classification_eval_label`` in ``config/setup.yaml``. There is no
runtime inference — every classification label must be registered
explicitly.
Parameters
----------
case_study_id : str
Case study identifier (e.g., ``"us_firm_characteristics"``).
label : str
Classification label name (e.g., ``"fwd_class_1m"``).
Returns
-------
str
Continuous-return label name (e.g., ``"fwd_ret_1m"``).
Raises
------
KeyError
If ``labels.classification_eval_label[label]`` is missing from
``setup.yaml``.
"""
from utils import CASE_STUDIES_DIR
setup_path = CASE_STUDIES_DIR / case_study_id / "config" / "setup.yaml"
setup = yaml.safe_load(setup_path.read_text())
mapping = (setup.get("labels") or {}).get("classification_eval_label") or {}
if label not in mapping:
raise KeyError(
f"labels.classification_eval_label[{label!r}] not declared in "
f"case_studies/{case_study_id}/config/setup.yaml. Add an entry mapping "
f"the binary/categorical label to its source continuous-return label "
f"(e.g., {label}: fwd_ret_1m). IC for classification predictions is "
f"computed against the continuous return — there is no silent fallback."
)
return str(mapping[label])
def get_direction_labels(case_study_id: str, label: str) -> tuple[str, ...]:
"""Return the classification labels cut from a given continuous-return label.
``labels.classification_eval_label`` read backwards. It declares which continuous return
each classification label was derived from; inverting it says which direction labels exist
for a return, which is what scoring a regression model by AUC needs. No separate
declaration is added for the inverse - one mapping cannot disagree with itself.
Returns an empty tuple when the case study declares no classification label for this
return, which is the case for five of the nine.
"""
from utils import CASE_STUDIES_DIR
setup_path = CASE_STUDIES_DIR / case_study_id / "config" / "setup.yaml"
setup = yaml.safe_load(setup_path.read_text())
mapping = (setup.get("labels") or {}).get("classification_eval_label") or {}
return tuple(sorted(k for k, v in mapping.items() if v == label))
def verify_artifact_sidecars(
artifacts: Mapping[str, Path],
*,
require_sidecar: bool = True,
values: bool = True,
) -> dict[str, str]:
"""Check each input artifact against the digest sidecar its producer wrote.
``case_studies/utils/artifact_digest.py`` writes ``<artifact>.digest.json``
beside every artifact, recording the content hash, the row count and the key
columns. Until this function existed nothing read one: stage 02 and stage 03
wrote them and the chain stopped, so an upstream value change could not reach
a downstream identity. What it closes is narrow and worth stating exactly - the
registry records feature-set *names* and no digest of feature *values*, so a
model trained on corrected features and one trained on the leaky version it
replaced produce the identical training_hash unless something reads the values.
Raises on a digest that disagrees with the file. Whether a *missing* sidecar
also raises is the caller's to choose, because the two carry different
evidence. A disagreeing sidecar is a value that moved without its record
moving with it, which is the defect itself, and it is never acceptable. A
missing sidecar means the producer recorded nothing, which cannot be told
apart from a silent change but is also the state of every artifact written
before the sidecar existed.
Returns the verified digest per artifact, so the caller can carry it into the
training spec. Artifacts skipped under ``require_sidecar=False``, and every
artifact under ``values=False``, are absent from the mapping rather than
present with a null.
``values`` chooses what the check costs, and the two settings answer different
questions. With it, the parquet is read and every row hashed, which is what the
recorded content digest is - order-independent, so it moves when and only when
a value moves. Without it, only the row count is compared, which parquet
metadata answers without reading a column. The cheap check cannot see a value
change that preserves the row count, and it does see the case that a crash
creates: ``write_artifact`` writes the parquet before its sidecar, so an
interrupted regeneration leaves new data beside the previous record, and new
data almost always has a different number of rows.
``load_modeling_dataset`` runs the cheap check on every load and the full one
under ``verify_input_digests``; ``scripts/verify_artifact_sidecars.py`` runs the
full one over every case study, once, after a regeneration. That split is about
cost: hashing us_equities_panel's feature matrix is a multi-GB read plus a sort
over 68M row hashes.
Whether a *missing* sidecar raises is the caller's to choose, because presence
and disagreement carry different evidence. A missing sidecar means the producer
recorded nothing, which is the state of every artifact written before the
sidecar existed, and requiring presence before those are regenerated fails for
absence rather than for a changed value. Measured rather than assumed: on
2026-08-09 the production artifacts carried 49 sidecars out of 51, both gaps in
sp500_options, while the CI fixtures under
``ml4t/third-edition-test-data/intermediates`` carried **none** - 0 of 7 for
cme_futures, 0 of 5 for etfs, 0 of 6 for fx_pairs, 0 of 6 for us_equities_panel.
"""
from case_studies.utils.artifact_digest import read_digest, sidecar_path, value_digest
verified: dict[str, str] = {}
problems: list[str] = []
for name, path in sorted(artifacts.items()):
side = sidecar_path(path)
if not side.exists():
if require_sidecar:
problems.append(f"{name}: no digest sidecar beside {path.name}")
continue
try:
record = read_digest(path)
except (OSError, ValueError, json.JSONDecodeError) as exc:
problems.append(f"{name}: sidecar unreadable ({exc})")
continue
recorded = record.get("digest")
if not recorded:
problems.append(f"{name}: sidecar records no digest")
continue
rows = record.get("n_rows")
if not values:
if rows is None:
continue
try:
actual_rows = cast(
pl.DataFrame,
pl.scan_parquet(path).select(pl.len()).collect(),
).item()
except (OSError, pl.exceptions.PolarsError) as exc:
problems.append(f"{name}: {path.name} unreadable ({exc})")
continue
if int(rows) != int(actual_rows):
problems.append(
f"{name}: {path.name} holds {actual_rows} rows, its sidecar records {rows}"
)
continue
try:
frame = pl.read_parquet(path)
except (OSError, pl.exceptions.PolarsError) as exc:
problems.append(f"{name}: {path.name} unreadable ({exc})")
continue
actual = value_digest(frame)
if actual != recorded:
problems.append(f"{name}: {path.name} hashes {actual}, its sidecar records {recorded}")
continue
if rows is not None and int(rows) != frame.height:
problems.append(
f"{name}: {path.name} holds {frame.height} rows, its sidecar records {rows}"
)
continue
verified[name] = recorded
if problems:
raise ValueError(
"input artifacts do not match the digests their producers recorded, so a "
"training run built on them would carry an identity that does not describe "
"its inputs: " + "; ".join(problems)
)
return verified
def feature_storage_dtype(case_study_id: str) -> pl.DataType:
"""The float type a case study stores and fits its feature matrices in.
``float64`` unless ``features.storage_dtype`` in the case study's ``setup.yaml`` says
otherwise. Declared per case study because only the large panels need the narrower type:
``nasdaq100_microstructure`` is 16,098,877 rows x 88 features, which is 10.9 GB of modeling
dataset in double precision against 5.6 GB in single, while the small case studies fit
comfortably either way and narrowing them would move their numbers for no gain.
"""
# Read from the repository, not through ``get_case_study_dir``, which redirects to
# ``ML4T_OUTPUT_DIR``. Every CI job points that at a scratch directory holding no
# ``config/``, so going through the redirect would resolve every case study to the default
# and silently answer a question about what a case study declares with "nothing".
from utils import CASE_STUDIES_DIR
setup_path = CASE_STUDIES_DIR / case_study_id / "config" / "setup.yaml"
if not setup_path.exists():
return pl.Float64
declared = (yaml.safe_load(setup_path.read_text()) or {}).get("features", {})
name = (declared or {}).get("storage_dtype", "float64")
if name not in {"float32", "float64"}:
raise ValueError(
f"{case_study_id}: features.storage_dtype must be float32 or float64, got {name!r}"
)
return pl.Float32 if name == "float32" else pl.Float64
def load_modeling_dataset(
case_study_id: str,
primary_label: str,
max_symbols: int = 0,
symbols: list[str] | None = None,
verify_input_digests: bool = False,
columns: Sequence[str] | None = None,
) -> ModelingDataset:
"""Load and join features + temporal + labels for a case study.
This is the canonical data loading function for ALL Ch11+ modeling
notebooks. It handles schema detection (date vs timestamp, asset vs
stock_id vs product), temporal join-key casting, and universe reduction.
Parameters
----------
case_study_id : str
Case study identifier (e.g., "etfs", "crypto_perps_funding").
primary_label : str
Label file stem (e.g., "fwd_ret_21d").
max_symbols : int, default 0
Universe reduction for fast development. 0 = all symbols.
symbols : list of str, optional
Explicit symbol whitelist. When given, the universe is restricted to
exactly these symbols (intersected with what is available) and
``max_symbols`` is ignored. Used by tests to pin a small universe that
is guaranteed to exist in the reduced test-data (e.g. the Darts base
return series), rather than the top-by-history selection ``max_symbols``
makes — which can pick symbols absent from a sampled data set.
columns : sequence of str, optional
Feature columns to keep, projected into the scans so the panel is never
materialized at full width. ``max_symbols`` narrows the row axis and this
narrows the column axis; a caller that needs the whole universe can still
use this one. The join keys, the detected date and entity columns and the
label are always added, so the returned frame is the requested columns
plus what the join and the fold geometry need. A name that is in none of
the three artifacts raises rather than being dropped silently. Passing
nothing keeps the full width.
Returns
-------
ModelingDataset
Container with joined dataset, detected schema, and CV splits.
"""
case_dir = get_case_study_dir(case_study_id)
financial_spec = load_feature_spec(case_study_id, "financial")
temporal_spec = load_feature_spec(case_study_id, "model_based")
label_spec = load_label_spec(case_study_id, primary_label)
# Check prerequisites exist before loading
features_path = resolve_storage_path(
case_study_id, financial_spec, "features/financial.parquet"
)
label_path = resolve_storage_path(case_study_id, label_spec, f"labels/{primary_label}.parquet")
missing = []
if not features_path.exists():
missing.append(("features/financial.parquet", "03_financial_features"))
if not label_path.exists():
missing.append((f"labels/{primary_label}.parquet", "02_labels"))
if not (case_dir / "config" / "setup.yaml").exists():
missing.append(("config/setup.yaml", None))
if missing:
print(f"\n Missing prerequisites for '{case_study_id}' modeling:\n")
for path, producer in missing:
if producer is None:
print(f" {path} (canonical hand-curated file — ensure committed)")
else:
print(f" {path} (run {producer}.py first)")
first_producer = next((p for _, p in missing if p is not None), None)
if first_producer is not None:
print("\n Example:")
print(f" uv run python case_studies/{case_study_id}/{first_producer}.py\n")
raise FileNotFoundError(
f"Missing prerequisites for '{case_study_id}': " + ", ".join(p for p, _ in missing)
)
# Load artifacts. The declared storage type is applied in the scan, so a case study that
# fits in single precision never materialises the double-precision form on the way.
storage_dtype = feature_storage_dtype(case_study_id)
def _narrow(frame: pl.LazyFrame) -> pl.LazyFrame:
if storage_dtype != pl.Float64:
narrow = [n for n, t in frame.collect_schema().items() if t == pl.Float64]
if narrow:
frame = frame.with_columns([pl.col(c).cast(storage_dtype) for c in narrow])
return frame
# Scanned, not read. A reduced run has to narrow the entity axis *before* the panel is
# materialised, and it cannot choose which entities to keep until the join keys are known,
# so the collect waits until both are settled - see "Universe reduction" below.
features_lazy = _narrow(pl.scan_parquet(features_path))
temporal_path = resolve_storage_path(
case_study_id, temporal_spec, "features/model_based.parquet"
)
# Scanned, not read. This artifact is per-fold, so it is a multiple of the feature table:
# 68.7M rows on us_equities_panel against the panel's 9.9M. Every consumer wants one fold,
# so the fold predicate is pushed into the scan instead of the whole thing being held.
temporal = pl.scan_parquet(temporal_path) if temporal_path.exists() else None
if temporal is not None and storage_dtype != pl.Float64:
narrow = [n for n, t in temporal.collect_schema().items() if t == pl.Float64]
if narrow:
temporal = temporal.with_columns([pl.col(c).cast(storage_dtype) for c in narrow])
temporal_columns = temporal.collect_schema().names() if temporal is not None else []
# Labels are deliberately not narrowed. ``features.storage_dtype`` covers the design
# matrix; the label is the target IC and every metric are measured against, and
# ``gbm_fold`` states that it stays float64 whatever the design matrix is cast to.
labels_lazy = pl.scan_parquet(label_path)
label_columns = labels_lazy.collect_schema().names()
feature_columns = features_lazy.collect_schema().names()
# Auto-detect label column (the non-ID column in the label file)
label_col = [c for c in label_columns if c not in ID_COLS][0]
# Detect date column from features
feature_keys = sorted(set(feature_columns) & ID_COLS)
date_col = "timestamp" if "timestamp" in feature_keys else "date"
alt_date = "timestamp" if date_col == "date" else "date"
# Normalize date column names across DataFrames
if alt_date in label_columns and date_col not in label_columns:
labels_lazy = labels_lazy.rename({alt_date: date_col})
label_columns = [date_col if c == alt_date else c for c in label_columns]
if temporal is not None and alt_date in temporal_columns and date_col not in temporal_columns:
temporal = temporal.rename({alt_date: date_col})
temporal_columns = [date_col if c == alt_date else c for c in temporal_columns]
# Detect join columns
label_keys = sorted(set(label_columns) & ID_COLS)
join_cols = sorted(set(feature_keys) & set(label_keys))
entity_cols = [c for c in join_cols if c != date_col]
# One pass for both uses below. The sort used to call ``n_unique`` from its key function,
# which re-counts the column on every comparison.
cardinality: dict[str, int] = {}
if entity_cols:
cardinality = (
features_lazy.select([pl.col(c).n_unique().alias(c) for c in entity_cols])
.collect()
.row(0, named=True)
)
# Filter out constant entity columns (e.g. instrument_id='straddle_30d_atm')
# that break cross-sectional IC computation by collapsing all entities into one group.
# NOTE: join_cols retains ALL shared ID columns for data integrity during joins;
# entity_cols is filtered separately for IC computation only.
entity_cols = [c for c in entity_cols if cardinality[c] > 1]
# Sort by cardinality descending so the primary entity (most unique values)
# comes first. Important when downstream code uses entity_cols[0] for IC
# (e.g., CME futures: 'product' has 30 values vs 'position' has 3).
entity_cols = sorted(entity_cols, key=lambda c: cardinality[c], reverse=True)
# Column projection, pushed into the SCANS for the reason the universe reduction below is:
# a caller that projects the frame it is handed has already paid the full width.
# ``us_equities_panel``'s DML estimand reads 7 of the 74 columns the join returns and
# peaked at 17.6 GiB doing it, against 0.47 GB for the seven.
#
# The request is unioned with what the join and the fold geometry need rather than taken
# literally: without the join keys the frames cannot be joined, without the label there is
# nothing to model, and without ``fold`` the per-fold substitution has no key to select on.
# That union is why the parameter belongs here instead of in each caller.
#
# The projection must not reach the recorded identity. ``feature_names`` below is derived
# from the joined frame, and it is hashed twice into every registered run - directly as
# ``computation.feature_names`` and again inside ``input_data_spec`` - so narrowing the load
# would re-key every run whose caller passed ``columns``, while the artifacts read and the
# values retained are identical. ``build_modeling_input_lineage`` already states the rule for
# the same event in the other direction, about adding ``feature_dtype``: a change that moves
# the fingerprint of a case study whose declaration did not change invalidates every run
# registered against it. So the panel's own feature list is recorded here, before the
# projection narrows the column lists, and ``panel_feature_names`` below is what the identity
# reads. What the caller asked for narrows the load and nothing else.
panel_feature_columns = list(feature_columns)
panel_temporal_columns = list(temporal_columns)
if columns is not None:
requested = list(dict.fromkeys(columns))
known = set(feature_columns) | set(temporal_columns) | set(label_columns)
unknown = [c for c in requested if c not in known]
if unknown:
raise ValueError(
f"load_modeling_dataset({case_study_id!r}, {primary_label!r}) was asked for "
f"columns no artifact carries: {unknown}. The three artifacts carry "
f"{sorted(known)}."
)
keep_cols = set(requested) | set(join_cols) | set(feature_keys) | set(entity_cols)
keep_cols |= {date_col, label_col}
feature_columns = [c for c in feature_columns if c in keep_cols]
features_lazy = features_lazy.select(feature_columns)
if temporal is not None:
temporal_columns = [c for c in temporal_columns if c in keep_cols | {"fold"}]
temporal = temporal.select(temporal_columns)
# Universe reduction, pushed into the SCANS instead of applied to the finished panel.
#
# It used to run at the bottom of this function, after features and labels had both been
# read whole and joined, which left a reduction with almost nothing to save: five of
# nasdaq100_microstructure's 115 symbols - 4.3% of the universe - still peaked at 39.98 GB
# against the full run's 51.5 GB. A preview that costs 78% of production is not a preview,
# and it is what stopped the smoke-then-full loop from running on the two largest case
# studies at all.
#
# The universe it selects is unchanged. ``top_entities`` ranks entities by their row count
# in the finished panel, so the count is taken here on the key-only inner join of the two
# scans, which carries exactly the rows the panel carries: the temporal join is a left join
# against a frame made unique on its keys, and neither it nor the META_LEAK drop moves a row.
#
# Production runs pass ``max_symbols=0`` and no ``symbols``, so neither branch fires.
#
# Bound before the filter because the fold geometry below is derived from the label frame,
# and a reduced universe is a different timeline: ``generate_cv_splits`` reads the unique
# timestamps of whatever frame it is handed, so dropping 25 of cme_futures' 30 products
# removes four sessions and moves every fold boundary two sessions earlier. Measured
# 2026-09-06: the same call over the full universe puts fold 0's validation window at
# 2019-01-03..2020-01-02 and over a five-product preview at 2018-12-31..2019-12-30. Nothing
# downstream follows that shift - ``canonical_window`` always reads the whole label parquet -
# so the backtest loads prices for the canonical window, the preview's predictions carry two
# sessions that window does not, and ``Strategy._decision_weights`` refuses the decision
# artifact for keys outside the price grid. Before #780 the reduction ran after both frames
# were read whole, so this call already saw the full timeline; pushing it into the scans is
# what put a reduced frame here.
unreduced_labels_lazy = labels_lazy
if entity_cols:
primary_entity = entity_cols[0]
keep: list | None = None
if symbols:
keep = list(symbols)
elif max_symbols > 0:
from utils.data_quality import top_entities
keep = top_entities(
features_lazy.select(join_cols).join(
labels_lazy.select(join_cols), on=join_cols, how="inner"
),
max_symbols,
primary_entity,
)
if keep is not None:
# implode: is_in against a bare Series of the same dtype is deprecated in polars
# as ambiguous, and membership in the value set is what is meant.
keep_values = pl.Series(primary_entity, keep).implode()
features_lazy = features_lazy.filter(pl.col(primary_entity).is_in(keep_values))
labels_lazy = labels_lazy.filter(pl.col(primary_entity).is_in(keep_values))
if temporal is not None and primary_entity in temporal_columns:
temporal = temporal.filter(pl.col(primary_entity).is_in(keep_values))
features = features_lazy.collect()
labels = labels_lazy.collect()
# Join features + temporal (left join to keep all feature rows)
temporal_by_fold_pd = None
_temporal_keys = []
_temporal_feature_names = []
if temporal is not None:
temporal_schema = temporal.collect_schema()
_temporal_keys = sorted(set(temporal_columns) & set(feature_keys))
casts = {
k: features.schema[k]
for k in _temporal_keys
if temporal_schema[k] != features.schema[k]
}
if casts:
temporal = temporal.cast(casts)
if "fold" in temporal_columns:
# Per-fold temporal features. The dataset carries one fold as a placeholder so the
# schema is complete; the per-fold values are substituted at fold preparation time
# by ``fold_temporal_frame``, which reads the fold it is asked for and no more.
_temporal_feature_names = [
c for c in temporal_columns if c not in set(_temporal_keys) | {"fold"}
]
fold_ids = sorted(temporal.select(pl.col("fold").unique()).collect()["fold"].to_list())
placeholder_dedup = (
temporal.filter(pl.col("fold") == fold_ids[0])
.drop("fold")
.unique(subset=_temporal_keys, keep="last")
.collect()
)
dataset = features.join(placeholder_dedup, on=_temporal_keys, how="left", suffix="_t")
del placeholder_dedup
# Kept lazy. Materialising it here is what made a run hold every fold at once.
temporal_by_fold_pd = temporal
else:
# Fold-free: one value per key, joined straight on. A refit schedule produces this
# shape, and so does any stage that fits nothing per fold.
#
# The names are recorded here for the same reason the fold branch records them:
# they say which of the panel's columns came from the model-based artifact. They
# were not recorded before, so a fold-free artifact reported no model-based
# features at all while its columns sat in `dataset` regardless - the features
# were used, and nothing that asks which ones they are could answer. Every
# consumer of this list also requires `temporal_by_fold`, which stays None here,
# so filling it in changes no fold substitution.
_temporal_feature_names = [c for c in temporal_columns if c not in set(_temporal_keys)]
temporal_dedup = temporal.unique(subset=_temporal_keys, keep="last").collect()
dataset = features.join(temporal_dedup, on=_temporal_keys, how="left", suffix="_t")
del temporal_dedup
else:
dataset = features
# Inner-join with labels (drops rows without labels)
dataset = dataset.join(labels, on=join_cols, how="inner")
# Neither operand is read again, and until they were dropped here the function returned
# holding three panels: the joined dataset, the feature panel it was built from, and the
# labels. On the full nasdaq100_microstructure panel the feature panel alone is 5.6 GB.
del features, labels
# Drop any meta columns that leaked in
drop_cols = [c for c in dataset.columns if c in META_LEAK]
if drop_cols:
dataset = dataset.drop(drop_cols)
# Feature columns = everything except IDs and label
feature_names = [c for c in dataset.columns if c not in ID_COLS and c != label_col]
# The same list the unprojected load would have produced, reconstructed from the column
# lists captured before the projection. The joins concatenate left columns then the right
# frame's non-key columns, in source order, so the unprojected order is the three source
# lists in sequence under the same filters applied to ``feature_names`` above.
#
# The one case this reconstruction cannot express is a name carried by both the financial
# and the model-based artifact: polars would suffix the second ``_t`` and the position of
# the suffixed name is not recoverable from the source lists. No case study has one -
# checked across all eight on 2026-09-11 - so it is refused rather than guessed, and only
# when a projection is actually in play, because the unprojected path takes ``feature_names``
# itself and is exact by construction.
if columns is None:
panel_feature_names = list(feature_names)
else:
collisions = sorted(
(set(panel_feature_columns) & set(panel_temporal_columns)) - set(_temporal_keys)
)
if collisions:
raise ValueError(
f"{case_study_id}: the financial and model-based artifacts both carry "
f"{collisions}, so a projected load cannot reconstruct the panel's feature "
"order and would register an identity that differs from the unprojected "
"load's. Rename the duplicate column in one artifact, or load without "
"`columns`."
)
ordered = [
*panel_feature_columns,
*[c for c in panel_temporal_columns if c not in set(_temporal_keys) | {"fold"}],
*[c for c in label_columns if c not in set(join_cols)],
]
seen: set[str] = set()
panel_feature_names = [
c
for c in ordered
if c not in ID_COLS
and c not in META_LEAK
and c != label_col
and not (c in seen or seen.add(c))
]
# CV splits — read buffer from setup.yaml (explicit, handles non-standard labels)
setup = yaml.safe_load((case_dir / "config" / "setup.yaml").read_text())
label_buffer = resolve_label_buffer(case_study_id, primary_label, setup)
# Whether that duration counts sessions or calendar time is the label's to declare,
# not each consumer's to guess: `35D` to option expiry is calendar, `21D` on a daily
# equity panel is 21 sessions, and the duration alone cannot tell them apart.
buffer_unit = resolve_label_buffer_unit(case_study_id, primary_label, setup)
if not label_buffer:
raise ValueError(
f"No explicit label buffer found for '{primary_label}' in "
f"case_studies/{case_study_id}/config/setup.yaml. "
f"Add buffer to labels.buffer (primary) or labels.variant_buffers (variants)."
)
# Fold identity belongs to the label contract. Deriving boundaries from the
# feature-joined frame lets warm-up nulls or feature availability shift the
# calendar and makes model selection disagree with canonical_window().
splits = generate_cv_splits(
unreduced_labels_lazy.select(date_col).unique().collect(),
case_study_id=case_study_id,
label_buffer=label_buffer,
outcome_horizon=resolve_label_horizon(case_study_id, primary_label, setup),
date_col=date_col,
buffer_unit=buffer_unit,
)
temporal_artifact_splits: list[dict[str, Any]] = []
if temporal_by_fold_pd is not None:
assert temporal is not None
validate_temporal_fold_coverage(
dataset,
temporal,
splits,
date_col=date_col,
)
# The artifact carries one fold set, built on the primary label's geometry, and
# the rows above were selected from it by ``fold`` id. So the values this load
# returns for fold F were fit on data through the *primary* label's train_end,
# and a label whose own validation opens earlier than that is scored on sessions
# its features already saw. Coverage does not see this: an artifact can cover
# every date a variant is scored on and still have been fit past the start of
# that window. Checked here rather than in a notebook because every model
# notebook reaches its data through this call, and only the label being loaded
# is checked, so the cost is one label timeline rather than all of them.
from case_studies.utils.cv_window import (
assert_variant_folds_are_out_of_sample,
configured_labels,
temporal_artifact_fold_boundaries,
)
configured = configured_labels(case_study_id)
artifact_label = configured[0] if configured else primary_label
temporal_artifact_splits = temporal_artifact_fold_boundaries(
case_study_id,
artifact_label,
temporal_path,
)
if configured and primary_label != configured[0]:
assert_variant_folds_are_out_of_sample(
case_study_id, configured[0], variants=[primary_label]
)
# WalkForwardConfig for library integration
# Normalize month-based buffers to days (pd.Timedelta rejects 'M' as ambiguous)
wf_horizon = label_buffer
if wf_horizon and wf_horizon.endswith("M") and wf_horizon[:-1].isdigit():
wf_horizon = f"{int(wf_horizon[:-1]) * 30}D"
try:
cv_config = make_wf_config(
case_study_id,
label_horizon=wf_horizon,
date_col=date_col,
buffer_unit=buffer_unit,
)
except Exception as exc:
warnings.warn(f"WalkForwardConfig creation failed for {case_study_id}: {exc}", stacklevel=2)
cv_config = None
# Detect label type (regression vs classification)
task_type, num_classes, class_values = detect_label_type(label_col, dataset[label_col])
# Classification labels: load the continuous-return label they were derived
# from so IC can be computed against returns rather than the binary target.
eval_label_col: str | None = None
eval_label_path: Path | None = None
if task_type == "classification":
eval_label_col = get_classification_eval_label(case_study_id, label_col)
eval_label_path = resolve_storage_path(
case_study_id,
load_label_spec(case_study_id, eval_label_col),
f"labels/{eval_label_col}.parquet",
)
if not eval_label_path.exists():
raise FileNotFoundError(
f"Eval label parquet missing for classification label {label_col!r}: "
f"expected {eval_label_path}. Generate it via 02_labels.py or update "
f"labels.classification_eval_label[{label_col}] in setup.yaml."
)
eval_labels = pl.read_parquet(eval_label_path)
# Normalize date column name if needed
if alt_date in eval_labels.columns and date_col not in eval_labels.columns:
eval_labels = eval_labels.rename({alt_date: date_col})
# Inner-join eval column on the same join_cols (drops rows missing eval)
eval_join_cols = sorted(set(eval_labels.columns) & set(dataset.columns) & ID_COLS)
eval_labels = eval_labels.select([*eval_join_cols, eval_label_col])
dataset = dataset.join(eval_labels, on=eval_join_cols, how="left")
# Refresh feature_names so the eval label column is not used as a feature
feature_names = [
c for c in dataset.columns if c not in ID_COLS and c not in {label_col, eval_label_col}
]
panel_feature_names = [c for c in panel_feature_names if c != eval_label_col]
input_artifacts = {
"financial": features_path,
"label": label_path,
}
if temporal_path.exists():
input_artifacts["model_based"] = temporal_path
if eval_label_path is not None:
input_artifacts["eval_label"] = eval_label_path
# Two depths, and the flag chooses between them rather than between checking
# and not. The row-count comparison reads parquet metadata and no column, so it
# runs always: write_artifact writes the parquet before its sidecar, and an
# interrupted regeneration therefore leaves new data beside the previous
# record, which almost always has a different number of rows. Comparing the
# content digest is what catches a value change that keeps the row count, and
# that means hashing every row - a multi-GB read plus a sort over 68M row
# hashes on us_equities_panel, paid per load. It is the same work
# scripts/verify_artifact_sidecars.py does once over everything after a
# regeneration, which is where a sweep should pay it.
verify_artifact_sidecars(
input_artifacts,
require_sidecar=verify_input_digests,
values=verify_input_digests,
)
return ModelingDataset(
dataset=dataset,
feature_names=feature_names,
panel_feature_names=panel_feature_names,
label_col=label_col,
date_col=date_col,
entity_cols=entity_cols,
join_cols=join_cols,
splits=splits,
label_buffer=label_buffer,
case_study_id=case_study_id,
cv_config=cv_config,
task_type=task_type,
num_classes=num_classes,
class_values=class_values,
temporal_by_fold=temporal_by_fold_pd,
temporal_keys=_temporal_keys,
temporal_feature_names=_temporal_feature_names,
temporal_artifact_splits=temporal_artifact_splits,
feature_dtype="float32" if storage_dtype == pl.Float32 else "float64",
eval_label_col=eval_label_col,
lineage_inputs={
"artifacts": input_artifacts,
"max_symbols": max_symbols,
"symbols": symbols,
},
)
def reduce_to_top_entities(
dataset: pl.DataFrame,
primary_entity: str,
max_symbols: int,
) -> pl.DataFrame:
"""Keep the ``max_symbols`` entitieShown in full with attribution under the source's licence. Licence: MIT
This summary was written by Stratmill's research agent from the original; it is not a copy of the source.