Skip to content

Performance and Scale

The full pipeline — training and a complete 51-member backtest — runs on an ordinary laptop. Everything on this page exists to keep that true, because it is what keeps experimentation fast (H2, a hundred experiments per person in a peak month) and the running service cheap (H4, it runs for pocket money).

The page holds the measured performance engineering behind three of the design principles — the whole system must be exercisable on one laptop, push the work down to the query engine; materialise once, as late as possible, and measure; do not assume.

Storage formats: measured, not assumed

Both big Delta tables use the same precision lever: values are rounded to a 13-bit significand in pure Polars arithmetic — Veltkamp splitting, rigorously documented on delta_store.precision.round_to_significand_bits (max relative error 2⁻¹³ ≈ 1.2×10⁻⁴) — which zeroes the low mantissa bits that are otherwise incompressible entropy. Beyond that shared lever, the two tables need different writer properties, and each choice below was measured rather than assumed:

  • NWP (delta_store.nwp) stores physical-unit Float32 with plain ZSTD-3 and parquet's default dictionary+RLE encodings — BYTE_STREAM_SPLIT measures worse here, because the significand rounding collapses many cells/members onto repeated values that a dictionary exploits better. Rows are sorted member-early so single-member reads can skip row groups. The result is ~58 GB per year for the full ECMWF ENS dataset for Great Britain, and a single day's NWP data takes about 1 minute to download and convert. The writer-property comparison behind the format choice is in PR #271; the ~58 GB per year is the measured average of the 899 daily runs on disk today, which is higher than the figure that PR extrapolated.

  • power_forecasts (delta_store.power_forecasts) sorts ensemble members of the same target adjacent, uses DELTA_BINARY_PACKED timestamps and BYTE_STREAM_SPLIT for power_fcst (near-continuous ML output has no repeats for a dictionary to exploit), shrinking the 403.6M-row development table from 6.33 GB to 0.73 GB (8.7×). Details in PR #268.

Lazy evaluation strategy

Rule: everything stays lazy until the model boundary.

The entire feature-engineering pipeline operates on Polars LazyFrames. No code between reading from storage and calling BaseForecaster.train() or predict() is allowed to call .collect(). The pipeline is:

Nwp.scan_delta(path)                     # lazy scan, already Float32 physical units
  → FeatureEngineer.engineer()           # spatial join + feature pipeline;
                                         #   returns pt.LazyFrame[AllFeatures]; no collect
  → BaseForecaster.train/predict()       # subclass materialises the data here, at the boundary

Why it matters:

  • Memory: Polars only materialises the data at the last moment, at the model boundary. Intermediate representations (individual join outputs, lag columns before filtering) are never written to RAM.
  • Query optimisation: Polars can see the full plan end-to-end and push down filters (e.g. date ranges, time_series_id subsets) into the Delta Lake scan before any data crosses the wire.
  • Clean boundaries: The materialisation site is the model boundary — the one place where a third-party library (XGBoost, PyTorch, …) needs a concrete in-memory array. Keeping it there makes the boundary explicit and easy to find.

The one exception is a limit(1).collect() guard in _engineer_features (inside ml_core), reached only when a weather lag feature is requested. It checks the raw NWP scan for the control member. In bulk mode (training/backtesting) a missing control member fails the pipeline loudly, because a silently-degraded training run would poison every comparison built on it. In single-run mode (production inference, replay) a missing control member is the outside world misbehaving, not a contract violation, so the guard instead logs a warning naming the run and falls through to the ordinary join-miss path. The output stays full height, with every affected weather-lag column null rather than the row dropped. It probes the raw frame rather than the upsampled one deliberately: ensemble_member is one of the upsample's group-by keys, so the two are equivalent, but SLICE cannot push through the upsample's window functions — probing the upsampled frame would execute the whole upsample of the control member before answering.

Contract for BaseForecaster subclasses: the lazy plan is materialised at the model boundary, as late as possible, and never by callers — never ask a caller to collect before passing data in. A subclass typically does a single .collect(); keeping that bounded is the caller's job, by pruning the inputs before feature engineering (not by filtering the engineered output — see below).

Bounding feature-engineering memory: prune the inputs, not the output

The feature-engineering plan is dominated by the NWP scan. The NWP Delta is large — for the V1 trial it is ~142 GB: 899 daily init_time partitions, each up to ~7.24M rows = 1671 H3 cells × 51 ensemble members × 85 native steps (control member alone is 142k rows/partition). The 30-min upsample and the multi-run bulk join inflate that further. So the whole memory question is: how little of that NWP do we touch?

You cannot prune the scan by filtering the engineered output. data.filter(time_series_id == x) runs the cell-attach join, the 30-min upsample (group_by + explode + interpolate) and the bulk join first, then drops rows. And NWP is keyed by h3_index, not time_series_id, so a time_series_id predicate can never reach the NWP scan at all. Pruning must be applied to the raw inputs, in load_engineering_inputs, directly on the Nwp.scan_delta scan.

What actually prunes the NWP scan — verified with LazyFrame.explain():

Predicate on the raw NWP scan Effect
init_time ∈ [start − 16d, end] Partition prune — NWP is partitioned by (nwp_model_id, init_time), so only those partition directories are opened. A valid_time filter alone does not prune partitions (it scanned ~887 files).
ensemble_member ∈ {…} on the raw scan Only the requested members are decoded — this requires the predicate to reach the Parquet scan unchanged, which in turn requires Nwp's declared ensemble_member dtype to match what delta-rs actually stores (see Nwp.ensemble_member). Once it does, delta_store.nwp writes one row group per ensemble member, so a single-member predicate matches one row group's min/max range and skips the rest of the partition — for any member, not only the control member. Measured on the stored table, a single-member read decodes 1.96% of a partition — one row group in 51. Every partition the census sampled held 51 row groups, each spanning a single member, the 51 together covering members 0 to 50, so the share is the same for every member.
h3_index ∈ {cells} Restricts to the cells the requested series sit in, on the same condition as ensemble_member above (the declared dtype must match what's on disk). h3_index is not a sort-early column (see below), so this is row-level filtering during decode rather than row-group skipping.
time_series_id == x on the output Prunes the power scan + metadata join, but not the NWP scan (no time_series_id there), and only after the upsample. Doesn't help.
slice(n) / head(n) (row count) on the output Nothing — a row slice can't push through a join; the whole join runs first.

Pruning the partitions is necessary but not sufficient — you must also stream the collect. init_time prunes whole partition files, but ensemble_member and h3_index are not partition keys, so they act within each file. The ensemble_member predicate is cheap because of the on-disk layout: rows are ordered init → member → valid → h3 (delta_store.nwp.NWP_SORT_COLS) and the row-group size is set to the rows one member occupies, so each row group holds a single member and a single-member read skips every other row group via min/max stats, provided the predicate reaches the scan unchanged (see the table above). Sizing the row groups is what makes single-member pruning work for every member: row groups that straddle member boundaries advertise the whole span between their extremes, and only members at the ends of the range then prune. Reading 29 daily partitions, 9 H3 cells, and the control member alone takes 30 ms and 400 MB of peak resident memory, against 170 ms and 2,200 MB for the same read of a valid_time-first sorted table. Both tables were written freshly through write_nwp and differ only in the sort order, and each figure is the warm-cache median of five timed repetitions. The member-early sort uses 3.7% more storage. h3_index is not a sort-early column, though: cell filtering is still decode-then-filter within the surviving row groups, so an in-memory collect would still materialise every surviving row before dropping 1662/1671 of the cells. The streaming engine (collect(engine="streaming")) applies the predicates per morsel, so peak memory is set by the morsel/row-group size, not the window length: measured pre-migration, a full 15-month control-member collect stayed at ~2.5 GB. XGBoostForecaster.train/predict therefore collect with engine="streaming".

Many-to-one h3_index ↔ time_series_id: one NWP cell covers several series (the 32 V1 series live in just 9 cells, one holding 12), so h3_index pruning is keyed on the (few) cells the requested series occupy, and the feature engineer's spatial join replicates each cell's weather across its series.

Resulting design (load_engineering_inputs applies all three NWP predicates — init_time, ensemble_member, and h3_index = the requested series' cells — to the raw scan, and every collect streams):

control member (train) full 51-member ensemble (predict)
All eligible series at once ~5 GB → train collects once, groups in memory ~25 GB ✘ (OOMs)
One init_time chunk at a time — ~9 GB → cv_power_forecasts chunks by init_time, appends to Delta

Prediction is bounded by chunking on init_time (_PREDICT_INIT_CHUNK, 14 days), not by cell. init_time is one of the table's two partition columns and the axis that fans the output out across runs, so a chunk's forecast frame stays small while each partition is read exactly once and all series/cells/members are processed together (a per-cell loop instead OOMs on the busiest cell — 10 series × 51 members × the 10-month window ≈ 116M rows ≈ 25 GB). Measured end-to-end on the mid_2025_to_mid_2026 fold: training peaks ~5 GB and the full 51-member validation prediction (~321M forecast rows) peaks ~9 GB — both well under a 24 GB laptop.

Scoring the metrics: batch the series, and stream every scan

Scoring a whole fold at once OOM-kills a 29 GB machine, so _score_forecast_group scores a batch of time series at a time. compute_metrics is independent per series, which is what makes the batching sound. Measured on the mid_2025_to_mid_2026 fold when it held 28 series of about 13M forecast rows each, a batch of 4 series completes in roughly 50 seconds at about 18 GB peak process resident set size. That fold has gained series since, so the figures below understate today's fold and the case for batching is only stronger. That peak is dominated by the streaming Delta scan rather than by the data: each batch re-scans the partition, while the materialised batch frame itself is a few GB. A batch of 2 measured only about 2 GB lower, so a smaller batch saves little memory — the scan overhead is roughly constant in the fold size, which is what keeps the approach workable as folds grow to V2 scale.

Every scan in the scoring path streams, because the eager equivalent materialises full-length rows before it reduces them. Discovering which groups a fold holds projects just the two partition columns, experiment_name and fold_id, and takes their distinct values. Collected eagerly, those two columns are still materialised at full row length before the unique, which OOM-killed that fold at its then-size of 364M rows; the streaming collect peaks at 0.3 GB. The same rule applies to the time_series_id and ensemble_member values cv_power_forecasts accumulates from each chunk: taking the unique values gives a roughly 30-element Python list per chunk instead of a roughly 14M-element one. Measured on a 14M-row Int32 column, the full-length list cost 2.78 seconds and 560 MB peak against 0.05 seconds and no measurable allocation — and that 560 MB landed inside the loop whose whole job is holding the frame at 2 to 3 GB.

The other hard ceiling: Polars' 32-bit row index

RAM is not the only bound a Polars query must respect. Default Polars builds use a 32-bit row index (IdxSize), so any single materialised frame, row count, or row index is capped at 2³² ≈ 4.29 billion rows — and crossing the cap raises no error. Row counts wrap modulo 2³²: pl.scan_delta over the ~5.9-billion-row NWP dev table reports 1,652,180,189 rows via .select(pl.len()) (= 5,947,147,485 mod 2³², exactly), and group_by(...).agg(pl.len()) wraps identically for any single group past the cap, streaming engine included. This was investigated and empirically pinned down in issue #293: the wraparound reproduces on plain Parquet scans (it is not a Delta Lake bug), and the polars-u64-idx build (64-bit index, lags mainline releases) returns correct counts on the same queries.

What is and isn't affected — each verified empirically on 5-billion-row scans:

Operation on a >2³²-row scan Outcome
pl.len() / group_by(...).agg(pl.len()) where a count exceeds 2³² Silently wraps mod 2³²
Materialising a single frame of ≥2³² rows Unsupported (in practice RAM is exhausted first, which at least fails loudly)
Filtered / partition-pruned query whose result is < 2³² rows Correct
Value aggregations (sum, min/max, quantiles) over > 2³² rows Correct — values aren't indices

The rules that follow:

  • Never row-count a table that can exceed 2³² rows with Polars. Whole-table counts and data-completeness checks must come from the Delta transaction log — DeltaTable(path).count(), or summing num_records over get_add_actions(flatten=True) — which is metadata-only, exact, and reads no data files.
  • The lazy evaluation strategy above already keeps every production collect far below the cap (input pruning + init_time chunking), so pipeline reads are unaffected. The cap is a second, independent reason that rule exists: an unbounded collect of a V2-scale table wouldn't just OOM, its row accounting would be wrong.
  • Tables past the cap today: NWP (~5.9B rows). power_forecasts will pass it at V2 scale (~1 trillion rows — see forecast delivery). The metrics asset discovers (experiment_name, fold_id) groups via a streaming unique() scan, then scores each group in batches of four time_series_id values at a time, so peak memory is one batch, never a whole fold or the whole matched population — the same chunking that keeps the row-index cap out of reach.