fix(precompute): in-place z-score + drop feature_vectors early

Previous workflow train-4qwtc hit memory-pressure thrash (~56Gi
cgroup.current sitting at the 56Gi pod limit, kernel reclaim
hammering page cache) right after OFI completed on the 9-quarter
17.8M-bar dataset. Two refactors reduce peak by ~12GB:

(1) walk_forward.rs: new `normalize_batch_in_place(&mut features)`
    that rewrites the slice in place. The previous `normalize_batch`
    `.collect()`s a new Vec — at this dataset size that's a
    transient ~6GB peak while both pre- and post-normalised arrays
    are alive.

(2) precompute_features.rs:
    - call `normalize_batch_in_place` instead of the rebinding form.
    - explicit `drop(feature_vectors)` after copying the slice into
      `features` — `feature_vectors` would otherwise stay alive
      until end-of-main shadowing the ~6GB allocation through
      every downstream step.

Combined with the prior `t.into_iter()` refactor (a27cb40a9), the
peak transient drops from ~56GB to ~44GB — well under the 56Gi pod
limit on the existing ci-compile-cpu pool (POP2-HC-32C-64G).

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
This commit is contained in:
jgrusewski
2026-05-16 10:32:11 +02:00
parent 110d3b4125
commit 623ebfcf71
2 changed files with 25 additions and 5 deletions

View File

@@ -448,6 +448,10 @@ async fn main() -> Result<()> {
let feature_vectors = extract_features_from_bars(&all_bars)
.context("Feature extraction failed")?;
info!("Extracted {} feature vectors in {:.1}s", feature_vectors.len(), t1.elapsed().as_secs_f64());
// feature_vectors will be CONSUMED by the `to_vec()` slice copy below
// and then dropped explicitly at the end of step 3. Without the drop
// it'd live alongside the 6GB normalised features Vec through every
// downstream step.
// ── Step 3: Build targets [preproc_close, preproc_next, raw_close, raw_next, raw_open, mid_price_open]
// [0:1] = log-return-normalized close/next (network input; matches contract
@@ -468,7 +472,10 @@ async fn main() -> Result<()> {
// shrinks by `LOOKAHEAD_HORIZON_MAX` bars (we need
// `all_bars[i + WARMUP + 30]` for the 30-bar log-return).
let n = feature_vectors.len().saturating_sub(LOOKAHEAD_HORIZON_MAX);
let features: Vec<[f64; 42]> = feature_vectors[..n].to_vec();
let mut features: Vec<[f64; 42]> = feature_vectors[..n].to_vec();
// feature_vectors now duplicates `features` data; drop the original
// so the ~6GB Vec doesn't shadow us through the OFI + alpha steps.
drop(feature_vectors);
let mut targets: Vec<[f64; 6]> = Vec::with_capacity(n);
for i in 0..n {
let raw_curr = all_bars[i + WARMUP].close;
@@ -732,11 +739,13 @@ async fn main() -> Result<()> {
// ── Normalize features (z-score) ──────────────────────────────────────────
// Raw features contain OHLCV prices (up to 25,000). Normalize once at
// precompute time so every consumer
// (training, hyperopt, inference) gets consistent normalized features.
// precompute time so every consumer (training, hyperopt, inference)
// gets consistent normalized features. Use the in-place variant so
// we don't transiently hold pre- + post-normalised copies (~6GB
// saving over `.normalize_batch(&features)` at this dataset size).
let norm_stats = ml::walk_forward::NormStats::from_features(&features);
let features = norm_stats.normalize_batch(&features);
info!("Features z-score normalized ({} bars × 42 dims)", features.len());
norm_stats.normalize_batch_in_place(&mut features);
info!("Features z-score normalized in place ({} bars × 42 dims)", features.len());
// Defence-in-depth gate: never write a poisoned fxcache to disk. If anything
// upstream produced an out-of-bounds value, fail loudly here rather than

View File

@@ -797,6 +797,17 @@ impl NormStats {
features.iter().map(|f| self.normalize(f)).collect()
}
/// In-place variant of `normalize_batch`: rewrites `features` rather
/// than allocating a new Vec. Saves a transient ~6GB peak when
/// `features` is the precompute_features pass over 17.8M bars
/// (which previously held both the pre- and post-normalised arrays
/// alive simultaneously while `.collect()` built the new one).
pub fn normalize_batch_in_place(&self, features: &mut [[f64; FEATURE_DIM]]) {
for f in features.iter_mut() {
*f = self.normalize(f);
}
}
/// Inverse of `normalize`: recover raw `feature` from a previously-
/// normalised value via `raw = norm * std + mean`. Caveat: the forward
/// path applies `clamp(±NORMALIZED_FEATURE_BOUND)` post-divide, so any