ML Core API
ml_core
The model-agnostic half of the forecasting stack: the interface every model implements, the feature pipeline that feeds it, the metrics that score it, and the helpers that train, promote, and serve it.
The vocabulary the rest of this page assumes
The rest of this page leans on vocabulary from four fields at once — Dagster orchestration, the Polars and Patito data stack, MLflow experiment tracking, and operational weather forecasting — so the nouns come first. Each noun is defined only to the depth this package needs, and the pages linked from the sections below go deeper.
- A Dagster asset is a named, schedulable unit of computation producing one stored output. Dagster is the orchestrator the pipeline runs in. An asset check runs after an asset and reports a warning without blocking. A partition is one slice of an asset — one day, or one cross-validation fold — materialised and re-runnable on its own, and re-materialising a partition means running that slice again.
- A frame is a Polars dataframe, the type every function here passes around. Patito is the library that declares a schema class for a frame, attaches that class to the frame, and validates the frame against the class. Delta Lake is the on-disk table format the pipeline reads and writes.
- An MLflow run is one recorded training job, holding parameters, metrics, tags, and stored
artifact files under a run id. MLflow is the experiment tracker, and the leaderboard is
MLflow's ranked comparison of those runs' scores. Promoting a model means choosing one run's
saved model to be the model the live service forecasts with, and
promotion.jsonrecords that choice. - The roster is the
TimeSeriesMetadatatable of every time series this project forecasts, one row pertime_series_id, carrying that series' location and its static attributes. H3 is a hexagonal grid covering the globe at numbered resolutions. Gridded weather is stored one row per resolution-5 cell, and each roster row records the resolution-5 cell its series sits in, so the weather table and the roster join on that cell. - A numerical weather prediction (NWP) run is one execution of a weather model, initialised at
nwp_init_timeand predicting many futurevalid_times. Each run carries several ensemble members — alternative weather futures for the same times — of which member 0 is the control member, the single unperturbed one. The analysis proxy answers a target time already in the past from the control member of the freshest run, standing in for the weather that actually happened. - A cross-validation (CV) fold is one train-and-validate split of the history, fixed by a pair
of date ranges and identified by a
fold_id. A slot is one scheduled time at which the live service issues a forecast. A lag feature is a value observed a fixed number of hours before the time being forecast, and the forecast lead time is the gap between the moment a forecast is issued and the time that forecast is for.
Why this package exists
ml_core owns everything about forecasting that is not specific to one model family. A concrete
model supplies five members — train, predict, save, load, and the trained_time_series_ids
property — and inherits the rest from this package: the feature pipeline that builds its input, the
MLflow wiring that gives its run an identity, the archive format its weights are shipped in, the
checks that decide whether a saved model can still be served, and the scoring that puts it on the
leaderboard. XGBoostForecaster in xgboost_forecaster is the only subclass outside the tests
today, so the second model family is the real test of that split. Writing that second family should
mean writing five members, not a second pipeline.
The Dagster assets delegate here rather than implementing the forecasting logic themselves. The
cross-validation, metrics, and live-inference assets in src/nged_substation_forecast/defs/ are
thin shells over the functions below. That thinness keeps every one of those functions unit-testable
without a Dagster instance, an MLflow server, or an object store. The assets hold the orchestration
— partitions, schedules, retries, and asset checks — and nothing else.
Every neighbouring package owns a piece of the data; ml_core owns what is done with that data.
contracts owns what every frame means, and ml_core consumes those Patito schemas without
declaring any of its own. delta_store owns how a table is physically written, and ml_core writes
no Delta table at all. The files ml_core does write are a saved model's frozen roster copy, the
gzipped tar archive that model ships in, and the promotion.json recording which run was promoted.
nged_data and dynamical_data own the ingest of observed power and of gridded weather, and
ml_core starts from whatever those two packages landed. geo owns H3 indexing, and
weather_utils owns the analysis-proxy query the weather-lag join is built on. What is left —
turning power, weather, and metadata into a model-ready frame, and turning a model's output into a
leaderboard row — is this package.
Two invariants span the whole package
Confusing the moment a forecast is issued with the moment the weather model ran is the failure
this package is built to prevent. power_fcst_init_time is the first of those two moments and
nwp_init_time is the second. The pipeline carries both moments as separate columns from end to
end. The two columns differ by the publication delay — the hours between an NWP run's init_time and
the moment that run reaches our own disk. A feature is therefore legitimate only if it was knowable
at power_fcst_init_time, never if it was merely knowable at nwp_init_time. A power lag shorter
than or equal to the forecast lead time would read an observation from at or after
power_fcst_init_time, which nobody had yet when the forecast was issued. Every such lag is set to
null, in _nullify_leaky_lags. Weather lags reaching back before power_fcst_init_time are
answered from the control member of the freshest NWP run covering that time, rather than from the
row's own run and its own ensemble member. Answering a past target time that way is the
dual-strategy join in _apply_weather_lag.
The train==predict population invariant keeps the leaderboard comparable. A model scores exactly
the time_series_id population it trained on, whatever the eligibility rules would admit today, and
BaseForecaster.trained_time_series_ids is that frozen record. Power coverage changes over time, so
a series can newly qualify or newly drop out between a training run and a scoring run. Letting the
scored population drift with the coverage would silently compare two experiments over two different
populations. The same invariant stops the live service forecasting a series the promoted model never
saw.
Research and production share one execution path
What differs between a cross-validation run and a live forecast is failure policy, not code. The
cross-validation, training, and metrics paths fail fast, because a quietly degraded training run
poisons every comparison built on it. The live path degrades instead: an absent input routes into
the always-output branch, the degradation is logged and reported to Sentry, and the forecast still
comes out with wider uncertainty. _engineer_features carries both policies in one function,
choosing between them on whether the caller supplied a power_fcst_init_time. A live forecast is
issued at one known moment, so the live caller passes that moment as power_fcst_init_time. Bulk
training and backtesting leave the argument unset and let every row derive its own moment from the
NWP run it sits on. Whether the argument is present is therefore exactly whether the call is a live
forecast. The reasoning for
that asymmetry is on Inherent
stability
and design principle
3.
Contents
features— theFeatureEngineerstrategy interface, its defaultTabularFeatureEngineerimplementation, and the declarative tabular pipeline behind that implementation. A forecaster wanting a different view of the data (a spatial crop for a convolutional neural network, say) supplies a differentFeatureEngineerrather than editing the shared pipeline. The package docstring below lists the sub-modules and says what each sub-module owns.base_forecaster—BaseForecaster, the abstract interface every model implements, andBaseForecasterConfig, the serialisable config carrying a trained model's experiment identity. Also holds the MLflow round-trip that both research and production read a saved model through, which ships a model as one replaceable archive rather than as a directory of files. The constant_MLFLOW_MODEL_ARTIFACTdocuments why that distinction matters.metrics— the scoring pipeline: effective capacity, the ensemble-aware metrics (fair continuous ranked probability score, Fortin-corrected spread, pinball loss, prediction-interval coverage probability, and interval width), the horizon-slice bands, and the MLflow aggregate keys. Effective capacity is the per-series denominator the normalised errors divide by, taken as the 99th percentile of observed|power|. A horizon-slice band is a group of forecast lead times the scores are reported over separately, so that a day-ahead score never hides behind an intraday one. The equations and the argument for each choice are on Evaluation metrics; the docstrings below say what this implementation does with those equations.cv_helpers— pure, input-only helpers for the cross-validation assets: fold date arithmetic, the eligibility rule, partition-key parsing, and config flattening. Fold design itself is on Cross-validation folds.mlflow_runs— idempotent resolution of an MLflow experiment, its parent run, and its per-fold child runs, always by tag rather than through a live handle, so separate processes and Dagster retries converge on the same run.production_helpers— the live-inference helpers: picking the NWP run a slot may legitimately use, building the dense forecast spine the power-centric join needs, and refusing to promote a saved model this code can no longer serve. The join is power-centric, meaning it keeps one output row per power row. No power observation exists for a time still in the future, so a live slot has to construct its future rows itself. The spine is that construction: a complete half-hourly grid of target times, one row per series, with the observed power joined in where any exists.repro— the git commit and the Delta-table versions stamped onto an MLflow run, so a result can be traced back to the exact code and data that produced it. Every function here is non-raising, because provenance must never fail the run it describes.
ml_core.features
Feature engineering package.
Four terms carry the rest of this docstring. NWP is numerical weather prediction: gridded
weather-model output, published as runs, each run initialised at an nwp_init_time and
predicting many future valid_times. The package runs in two join modes — bulk training over
every NWP run in the input, and single-run inference against one named run. A lag feature is a
value observed a fixed number of hours before the time being forecast, and a lag is leaky when
that observation would not yet exist at the moment the forecast is issued. The dual-strategy
join answers a weather lag from the row's own NWP run for a target time still in the future, and
from the freshest earlier run for a target time already past.
Public interface:
FeatureEngineer— abstract base; implement to swap the full feature pipeline.TabularFeatureEngineer— default implementation: containing-cell NWP spatial join followed by the declarative tabular pipeline.NWP_PUBLICATION_DELAY_HOURS— hours after an NWP run'sinit_timebefore we treat that run as usable. Public because the constant is the default of three public signatures:FeatureEngineer.engineer,TabularFeatureEngineer.engineer, andml_core.production_helpers.select_nwp_init_time. A caller reasoning about any of those three signatures needs a name to reach the constant by.
That public interface lives in the two public sub-modules below, plus the private _nwp. This
package re-exports NWP_PUBLICATION_DELAY_HOURS from _nwp. The two public sub-modules are:
feature_engineer— theFeatureEngineerabstract base, plusDEFAULT_LOCAL_TIMEZONE, the IANA zone the local-time features are computed in. The zone lives on the interface rather than on the default implementation. AFeatureEngineerforecasting another region can therefore override the zone through the sameengineer()call every production caller already makes.tabular_feature_engineer—TabularFeatureEngineerand the_engineer_featuresorchestrator the class delegates the tabular pipeline to. The module name carries no leading underscore because the class is public; every function the module defines is private.
Private sub-modules (importable by tests, not part of the public API):
_parsed_features— typed feature descriptors and theParsedFeaturesparser._nwp— NWP upsampling, processing, and the two power/NWP join modes._lags— power lags, weather lags (the dual-strategy join), and leaky-lag nullification.
Attributes
NWP_PUBLICATION_DELAY_HOURS = 9
module-attribute
Hours after an NWP run's init_time before we treat that run as usable.
The delay models when a run reaches our disk, not when Dynamical publish that run. Dynamical
publish each 00Z run between 08:05 and 08:20 UTC, and ecmwf_ens_schedule downloads it at 08:30
UTC. A 00Z run is therefore ours from roughly 08:30, which is 8.5 hours after that run's
init_time. Nine is the nearest whole hour at or after 8.5.
The feature pipeline uses the delay to derive power_fcst_init_time from nwp_init_time in
bulk mode, and to derive nwp_init_time when a single-run caller omits nwp_init_time.
select_nwp_init_time uses the delay to reconstruct availability for "replay" backfills.
select_analysis_proxy needs no delay: in single-run mode _engineer_features caps the
proxy at the NWP run that call already selected.
Of select_nwp_init_time's two modes, only "replay" needs the delay. A live forecast run
joins whatever NWP runs are genuinely on disk, so reality already constrains the NWP table to the
runs that were genuinely published. A replay of a past init time would otherwise join runs that only
landed afterwards — lookahead bias rather than mere inaccuracy. The asymmetry in full:
https://openclimatefix.github.io/nged-substation-forecast/architecture/production-deployment/#resolve-nwp-availability-asymmetrically-live-vs-replay
Two bounds constrain the value, given one 00Z run a day and forecast slots at 00/06/12/18 UTC. The
06:00 slot must not see that morning's run, which has not landed yet, so the value must exceed 6
hours. The 12:00 slot must see that morning's run, so the value must not exceed 12 hours. Both
bounds move if ecmwf_ens_schedule's start time changes.
__all__ = ['NWP_PUBLICATION_DELAY_HOURS', 'FeatureEngineer', 'TabularFeatureEngineer']
module-attribute
Classes
FeatureEngineer
Bases: ABC
Turns raw power + NWP + metadata into a model-ready AllFeatures frame.
Source code in packages/ml_core/src/ml_core/features/feature_engineer.py
34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 | |
Methods:
engineer(*, selected_features, power_time_series, time_series_metadata, nwp, power_fcst_init_time=None, nwp_init_time=None, nwp_publication_delay_hours=NWP_PUBLICATION_DELAY_HOURS, local_timezone=DEFAULT_LOCAL_TIMEZONE)
abstractmethod
Engineer features for training/inference.
Two operating modes, matching tabular_feature_engineer._engineer_features (see
that function for the full detail this docstring summarises):
- Bulk mode (
power_fcst_init_time=None, the default): NWP-centric, vectorised over every NWP run in the input, one forecast per(nwp_init_time, valid_time)pair. Used for training and multi-run backtesting. - Single-run mode (
power_fcst_init_timegiven): power-centric. Joins exactly one NWP run, named bynwp_init_time, or derived frompower_fcst_init_timeminusnwp_publication_delay_hourswhennwp_init_timeis omitted. Stamps a constantpower_fcst_init_timeon every row. Used for production inference and replay backfilling.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
selected_features
|
set[str]
|
The feature names to produce. |
required |
power_time_series
|
LazyFrame[PowerTimeSeries]
|
Observed power, one row per |
required |
time_series_metadata
|
DataFrame[TimeSeriesMetadata]
|
Per-time-series metadata. Carries |
required |
nwp
|
LazyFrame[Nwp]
|
Gridded NWP in physical units, keyed by |
required |
power_fcst_init_time
|
datetime | None
|
|
None
|
nwp_init_time
|
datetime | None
|
The NWP run to join in single-run mode. Must be |
None
|
nwp_publication_delay_hours
|
int
|
Delay used to derive whichever of
|
NWP_PUBLICATION_DELAY_HOURS
|
local_timezone
|
str
|
IANA zone the local-time features are computed in. |
DEFAULT_LOCAL_TIMEZONE
|
Returns:
| Type | Description |
|---|---|
LazyFrame[AllFeatures]
|
A lazy |
Source code in packages/ml_core/src/ml_core/features/feature_engineer.py
37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 | |
TabularFeatureEngineer
Bases: FeatureEngineer
Containing res-5 NWP cell per time series, then the declarative tabular pipeline.
Source code in packages/ml_core/src/ml_core/features/tabular_feature_engineer.py
83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 | |
Methods:
engineer(*, selected_features, power_time_series, time_series_metadata, nwp, power_fcst_init_time=None, nwp_init_time=None, nwp_publication_delay_hours=NWP_PUBLICATION_DELAY_HOURS, local_timezone=DEFAULT_LOCAL_TIMEZONE)
Map each NWP cell to the time series inside it, then run the tabular feature pipeline.
See FeatureEngineer.engineer for the argument and operating-mode contract, and
_engineer_features in this module for the pipeline itself — including
local_timezone, the IANA zone the local-time features are computed in.
Source code in packages/ml_core/src/ml_core/features/tabular_feature_engineer.py
86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 | |
ml_core.base_forecaster
BaseForecaster: the interface every forecasting model implements.
Also holds BaseForecasterConfig, which carries a trained model's experiment identity.
Attributes
TRAINED_METADATA_FILENAME = 'time_series_metadata.parquet'
module-attribute
Name of the file holding the TimeSeriesMetadata rows a saved model carries.
The file holds one row per series in the model's trained population. Production inference reads each
series' location from that file, never from the TimeSeriesMetadata roster. An unreadable or
thinned roster therefore cannot fail a live slot, and cannot silently drop a series from a slot.
Being the model's own frozen copy of what it trained against also keeps a series' H3 cell and static
feature values identical between training and serving. See
https://openclimatefix.github.io/nged-substation-forecast/design-philosophy/inherent-stability/#the-rules.
Logging the frozen metadata copy as a second MLflow artifact rather than putting the copy in the
model directory would reopen the merge problem _MLFLOW_MODEL_ARTIFACT documents.
Classes
BaseForecasterConfig
Bases: BaseModel
Universal configuration for all forecasting models.
Subclasses add model-specific hyperparameters. Having a shared base ensures that every
forecaster carries its own feature list and optional MLflow experiment id in one serialisable
object. That single object simplifies save/load and the conf/model/*.yaml config wiring.
An MLflow tag is a label that can be rewritten later, while an MLflow param is written once
and then fixed. The tag fields (weather_source, training_strategy) are stamped onto MLflow
runs for leaderboard grouping. The tag fields live here so that constructing the config from
a model YAML's model_params validates them at load time. Every conf/model/*.yaml file
carries both tag fields under model_params.
experiment_name is the per-experiment key, set to the MLflow experiment name at
registration and stored in the saved config. Every PowerForecast row carries
experiment_name, distinct from the model-family MODEL_NAME. random_seed is
threaded into each model's training so that re-training a fold reproduces the same model —
keeping retries and the leaderboard stable.
Model identity (name and version) lives on the BaseForecaster class itself as
MODEL_NAME and MODEL_VERSION. Both constants are properties of the implementation,
not of the experiment config.
Serialisation must be canonical. A config is compared and stored as its serialised form:
register_experiment stamps model_dump_json() onto the MLflow experiment as the
config tag and compares a re-registration against that tag. register_experiment also
logs flatten_config(...) as write-once MLflow params. Python's per-process string hash
randomisation means a set iterates in a different order in every process. A set dumped
straight to a list would therefore make two dumps of the same config differ. A
re-registration would then look like a config change, and its param write would be rejected.
selected_features is therefore serialised sorted. A subclass that adds a set-valued (or
otherwise unordered) field must do the same.
Unknown keys are rejected, not ignored (extra="forbid"). A key no field declares
raises ValidationError, so a misspelled hyperparameter in a run's config_overrides
fails at registration. A stored config carrying a key the current code no longer declares
is refused rather than silently losing that key. The recovery is to re-train, never to
hand-edit. Why this is worth failing over:
https://openclimatefix.github.io/nged-substation-forecast/ml_experimentation/model-configuration/#tweaking-a-config-for-an-experiment.
Source code in packages/ml_core/src/ml_core/base_forecaster.py
189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 | |
Attributes
model_config = ConfigDict(extra='forbid')
class-attribute
instance-attribute
selected_features
instance-attribute
ml_flow_experiment_id = None
class-attribute
instance-attribute
experiment_name = ''
class-attribute
instance-attribute
weather_source = ''
class-attribute
instance-attribute
training_strategy = ''
class-attribute
instance-attribute
random_seed = 0
class-attribute
instance-attribute
BaseForecaster
Bases: ABC
Defines the universal interface for all energy forecasting ML models.
Every forecasting model subclasses this abstract base to allow shared Dagster assets and evaluation code to remain completely agnostic to the underlying model implementation.
Subclasses must define MODEL_NAME, MODEL_VERSION, and CONFIG_CLASS as class-level
constants. MODEL_NAME and MODEL_VERSION are stamped onto every PowerForecast row
at predict time. MODEL_NAME is also stamped on the experiment's parent MLflow run, as
that run's model_family tag. The MLflow experiment name comes from neither constant: the
name is the experiment_name the launcher hands to register_experiment. Bumping
MODEL_VERSION requires a code change (intentional), not a config edit.
Lazy evaluation contract: train and predict both accept a pt.LazyFrame[AllFeatures].
Callers must not collect before passing data in. Collecting early wastes memory and prevents
Polars from optimising the full query plan. A subclass materialises the data at the model
boundary (typically a single .collect(), streamed). Keeping that collect bounded is the
caller's responsibility. The caller prunes the inputs: the NWP control member (member 0,
the one unperturbed member of the weather model's ensemble), the relevant H3 cells, and the
window's init_time partitions. Where all of the ensemble's members are needed, the caller
processes one init_time chunk at a time. Filtering the engineered output cannot prune the
upstream join, nor the upstream upsample that puts the weather onto the half-hourly power
grid. See "Bounding feature-engineering memory: prune the inputs, not the output" in
https://openclimatefix.github.io/nged-substation-forecast/architecture/performance/#bounding-feature-engineering-memory-prune-the-inputs-not-the-output.
Persistence has two layers. Subclasses implement save/load for their own on-disk
format and need know nothing about MLflow. The concrete
save_to_mlflow/load_from_mlflow methods are shared by all subclasses. Those methods
wrap that disk format with MLflow's artifact store, so the same trained model can be shared
across machines by round-tripping through a run.
Source code in packages/ml_core/src/ml_core/base_forecaster.py
245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 | |
Attributes
MODEL_NAME
class-attribute
MODEL_VERSION
class-attribute
CONFIG_CLASS
class-attribute
The config class that load rebuilds from a saved model_params mapping.
A subclass's load must build its config through this attribute rather than naming its config
class again. ml_core.production_helpers._check_meta_is_servable then always validates a
saved config against the same class load uses.
feature_engineer = TabularFeatureEngineer()
class-attribute
The feature pipeline this forecaster's data is engineered through.
the forecaster references a feature engineer rather than
implementing feature engineering. A forecaster can therefore swap the whole pipeline by
overriding this attribute with a different FeatureEngineer. The default produces the
tabular AllFeatures
frame that train/predict consume.
model_params = model_params
instance-attribute
trained_time_series_ids
abstractmethod
property
The sorted time_series_id population this model will serve a predict for.
The population is the model's own frozen record of who the model trained on — not a
statement about how the model is internally structured. A model may hold one sub-model
per series, a single model spanning many series, or anything in between. The contract is
only about which time_series_ids the model will score.
Why the frozen population is load-bearing. The train==predict population
invariant: the population a model is scored on must equal the population it was trained
on. The invariant holds even if the live eligibility set has drifted since training.
Power coverage changes over time, so a series may newly qualify or newly drop out.
Consumers (cv_power_forecasts, live_forecasts) filter their inputs to this set,
so a model scores exactly the population it learned — never whatever eligibility says
today. The frozen population is what keeps the leaderboard apples-to-apples, and what
stops production forecasting a series the model never saw.
Subclasses persist and reconstruct this set through their own save/load (e.g. in
meta.json), so the population survives a round-trip through MLflow.
Methods:
__init__(model_params)
Store the config that carries this model's hyper-parameters and experiment identity.
Source code in packages/ml_core/src/ml_core/base_forecaster.py
297 298 299 | |
save(path)
abstractmethod
Save this trained model to a directory, replacing every file already there.
Two requirements on every implementation:
- The saved directory must contain a
meta.jsonwith amodel_classfield. That field holds the fully-qualified{module}.{qualname}of the concrete subclass, for example"xgboost_forecaster.forecaster.XGBoostForecaster". Production inference (ml_core.production_helpers.load_forecaster_from_dir) uses that field to reconstruct the correct class from a plain model directory — no caller-supplied class name, no class registry, and no MLflow run (issue #221). pathmust be cleared first, so that saving over a directory holding a larger model's files leaves none of them behind. A dropped time series' weights would survive a re-train, if a save merged into the directory instead of replacing its contents (issue #197).XGBoostForecaster.saveclears withshutil.rmtree(path, ignore_errors=True).
The clearing requirement makes path the model's to own while it saves, so any file a
caller left there is gone afterwards. Depositing a file after a save is fine.
Depositing after a save is how production_helpers.fetch_model_artifacts puts
promotion.json beside the model.
Source code in packages/ml_core/src/ml_core/base_forecaster.py
325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 | |
load(path)
abstractmethod
classmethod
Reconstruct a trained instance from a previously saved directory.
Source code in packages/ml_core/src/ml_core/base_forecaster.py
349 350 351 352 | |
save_to_mlflow(run_id, *, time_series_metadata)
Upload this trained model to the given MLflow run, as one replaceable archive.
Writes the model to a temporary directory via save (the subclass's own format), then
adds the frozen metadata copy (write_trained_metadata). Packs that directory into a
single model.tar.gz, and logs that one file to the run's artifact root. Logging one
archive rather than a directory of files is what makes a re-upload replace the previous
model instead of merging with it — see _MLFLOW_MODEL_ARTIFACT. The caller is responsible
for setting the tracking URI (mlflow.set_tracking_uri) beforehand.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
run_id
|
str
|
The MLflow run to attach the artifact to. |
required |
time_series_metadata
|
DataFrame[TimeSeriesMetadata]
|
The roster rows this model was engineered against. Required,
because a model uploaded without them cannot be promoted — see
|
required |
Source code in packages/ml_core/src/ml_core/base_forecaster.py
354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 | |
load_from_mlflow(run_id)
classmethod
Download a trained model's archive from an MLflow run, unpack it and load it.
Downloads into a temporary directory and loads from there. There is deliberately no
local-disk cache. A CV fold run is reused across re-materialisations:
get_or_create_fold_run resolves the same run for every re-run of a fold's partition, and
each re-run overwrites that run's single model archive. One run_id therefore names
different model bytes at different times, so a cache keyed by run_id would not be unique
for its contents. Production inference makes no MLflow call at all either. Full rationale:
https://openclimatefix.github.io/nged-substation-forecast/architecture/ml-orchestration/#why-there-is-no-local-cache.
The caller is responsible for setting the tracking URI (mlflow.set_tracking_uri)
beforehand.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
run_id
|
str
|
The MLflow run the model was saved under. |
required |
Returns:
| Type | Description |
|---|---|
Self
|
The reconstructed, trained forecaster. |
Source code in packages/ml_core/src/ml_core/base_forecaster.py
392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 | |
train(data, time_series_ids)
abstractmethod
Fit the model on data for the given time_series_id population.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
data
|
LazyFrame[AllFeatures]
|
The engineered features (lazy; the caller must not pre-collect). |
required |
time_series_ids
|
list[int]
|
The population the model must train on — the caller's eligible set.
The model
decides how to map that population onto boosters. |
required |
Source code in packages/ml_core/src/ml_core/base_forecaster.py
422 423 424 425 426 427 428 429 430 431 432 433 434 435 | |
predict(data, *, fold_id='live')
abstractmethod
Return power forecasts for all rows in the given AllFeatures data.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
data
|
LazyFrame[AllFeatures]
|
The engineered features to forecast from. |
required |
fold_id
|
str
|
The value stamped onto every row's |
'live'
|
Returns:
| Type | Description |
|---|---|
DataFrame[PowerForecast]
|
One row per |
DataFrame[PowerForecast]
|
present in |
DataFrame[PowerForecast]
|
megavolt-amperes (MVA) according to that series' unit in |
DataFrame[PowerForecast]
|
row also holds the model-family identity ( |
DataFrame[PowerForecast]
|
|
DataFrame[PowerForecast]
|
( |
DataFrame[PowerForecast]
|
|
DataFrame[PowerForecast]
|
|
Source code in packages/ml_core/src/ml_core/base_forecaster.py
437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 | |
Functions:
write_trained_metadata(model_dir, time_series_metadata)
Write a model's frozen TimeSeriesMetadata copy into its saved directory.
Call this after a subclass's save, which clears the directory first.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
model_dir
|
Path
|
The directory a subclass's |
required |
time_series_metadata
|
DataFrame[TimeSeriesMetadata]
|
The roster rows the model was engineered against.
|
required |
Source code in packages/ml_core/src/ml_core/base_forecaster.py
90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 | |
load_trained_metadata(model_dir)
Read back the TimeSeriesMetadata rows a saved model carries.
Patito offers two ways to put a schema on a frame: validate checks every row against the
schema and raises on a violation, while set_model attaches the schema without checking.
This function uses set_model, matching how the cross-validation, training, and metrics
assets read the roster itself. load_trained_metadata reads back what was written, so
re-checking those rows would only reject rosters the rest of the system already accepts.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
model_dir
|
Path
|
A directory written by |
required |
Returns:
| Type | Description |
|---|---|
DataFrame[TimeSeriesMetadata]
|
One row per series in |
DataFrame[TimeSeriesMetadata]
|
That population is the population the model will serve a |
DataFrame[TimeSeriesMetadata]
|
population the model was engineered over. |
Raises:
| Type | Description |
|---|---|
FileNotFoundError
|
The directory holds no such file, so the model was saved by code
predating |
Source code in packages/ml_core/src/ml_core/base_forecaster.py
110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 | |
ml_core.metrics
Scoring forecasts against the power that was actually observed.
A forecast here is an ensemble: for one time series and one forecast run, the model emits many
member predictions of the same future half-hour, and those members are what the probabilistic
scores below measure. A cross-validation (CV) fold is one train-and-validate split of the
history, identified by a fold_id. MLflow is the experiment tracker the scores are logged to,
and its leaderboard is the ranked comparison of one run's scores against another's.
Four public functions, in the order the cross-validation Dagster assets call them.
compute_effective_capacity derives the per-series denominator that normalised mean absolute
error (NMAE) divides by. compute_metrics joins predictions to observed power and returns the
tall Metrics frame. enrich_metrics_rows stamps the evaluation window and scope onto that
frame once the calling asset knows the window and the scope. The enriched frame is what the asset
writes to Delta. build_mlflow_aggregate_metrics takes the un-enriched frame that
compute_metrics returned. That function reduces the frame to the flat key/value dictionary
the MLflow leaderboard displays.
Every function here is pure: no Dagster, no MLflow, and no IO. Each function is therefore unit-testable on an in-memory frame. The asset that calls the function owns every read and write.
The equations, and the argument for choosing each metric over the alternatives, are on https://openclimatefix.github.io/nged-substation-forecast/techniques/evaluation-metrics/. The docstrings below say what this implementation does with those equations. Those docstrings also record the corners where a degenerate input needed a defined answer rather than a NaN. The two degenerate inputs are a single-member ensemble and a forecast with zero error.
Attributes
Classes
NoOverlappingActualsError
Bases: ValueError
Raised by compute_metrics when no forecast row joins to any observed actual.
A distinct subclass, so that a caller scoring a fold in per-series batches can treat "this
batch's series have no overlapping actuals" as skippable. Skipping mirrors what happens when
the whole fold is scored in one call: those series silently vanish from the inner join. Every
other ValueError — a negative lead time, a missing capacity row — still propagates.
Source code in packages/ml_core/src/ml_core/metrics.py
52 53 54 55 56 57 58 59 | |
Functions:
compute_effective_capacity(power_lf)
Compute version 0.1 of the effective capacity, per time series.
Version 0.1 of the capacity definition is the full-history P99 of |power|.
One row per time_series_id: effective_capacity_mw is the 99th percentile of
abs(power) over all non-null observations. A series is dropped when its P99 is null or
non-positive (e.g. all-null or all-zero power), since EffectiveCapacity requires
effective_capacity_mw > 0.
time is set to that series' latest timestep carrying a non-null power (time.max()
over the same filtered rows the percentile is taken over). The v0.1 capacity is a single
scalar per series, so time is really an "as of" marker rather than a timestep the value
varies over. The marker stamps the estimate as current to the end of the observed history.
Version 0.7 of the capacity definition makes capacity genuinely time-varying (one row per
(time_series_id, time)), and only then does time carry per-row meaning. Version 0.1
stays one scalar row per series, rather than the value repeated at every half-hour.
Densifying a constant adds rows without information. The metrics join is by
time_series_id alone until capacity varies.
Kept as a pure helper (no Dagster, no IO) so the P99 logic is unit-testable in isolation.
Source code in packages/ml_core/src/ml_core/metrics.py
62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 | |
compute_metrics(cv_forecasts, actuals, metadata, capacity)
Compute evaluation metrics from CV predictions and observed power.
Full metric definitions, equations, and design rationale: https://openclimatefix.github.io/nged-substation-forecast/techniques/evaluation-metrics/
For each (time_series_id, fold_id, power_fcst_model_name) group:
- Joins predictions to observed
poweron(time_series_id, valid_time). - Assigns each row a
horizon_slicefrom its lead time (valid_time − power_fcst_init_time) — see_horizon_slice_exprfor the bands. - Collapses the ensemble members within each forecast run (per
power_fcst_init_time) into per-timestamp quantities: the deterministic ensemble mean, the fair CRPS, the Fortin-corrected ensemble variance, and the empiricalDELIVERY_QUANTILES. Each run covering avalid_timeis scored independently, exactly as a production consumer would experience that run. Runs at different lead times are never pooled. - Aggregates per
horizon_slice, plus the"all"aggregate over every lead time. The delivery quantiles are the quantiles the forecast is published at, listed byDELIVERY_QUANTILES. The ensemble mean gives MAE, NMAE, RMSE, and mean bias error (MBE). The member-aware quantities give CRPS and the spread-skill ratio. Each delivery quantile gives a pinball loss, and the pinball losses also get a mean. Each symmetric quantile band gives a prediction-interval coverage probability (PICP) and an interval width. - Joins
time_series_typefrommetadataonto each row. - Returns one row per
(time_series_id, fold_id, power_fcst_model_name, horizon_slice, metric_name, metric_param)in the tallMetricsformat.
NMAE is normalised by the pre-computed effective_capacity_mw, joined per time_series_id
from capacity. That denominator is capacity-like, and is computed over the full history so
that the denominator stays stable across folds. This time_series_id-only join is correct
while capacity is one scalar row per series, the shape at version 0.1 of the capacity
definition. Version 0.7 makes capacity time-varying, at one row per (time_series_id, time).
This join must then become a temporal as-of join on (time_series_id, valid_time).
Single-member "ensembles" — a deterministic baseline forecaster, for example — are scored unconditionally. Their fair CRPS equals their MAE, and their spread-skill ratio is 0. Their quantile bands are degenerate: all quantiles coincide, so PICP ≈ 0 and interval width = 0. Those values are honest descriptions of a deterministic forecast, not errors.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
cv_forecasts
|
DataFrame[PowerForecast]
|
CV predictions to evaluate. Currently only CV fold rows are handled. Support
for |
required |
actuals
|
LazyFrame[PowerTimeSeries]
|
Observed half-hourly power (lazy — only the joined subset is collected).
Deduplicated on |
required |
metadata
|
DataFrame[TimeSeriesMetadata]
|
Substation metadata used to join |
required |
capacity
|
DataFrame[EffectiveCapacity]
|
Pre-computed per-series effective capacity; |
required |
Returns:
| Type | Description |
|---|---|
DataFrame[Metrics]
|
A validated tall |
Raises:
| Type | Description |
|---|---|
NoOverlappingActualsError
|
If no rows survive the inner join (forecasts cover a period with no observed data). |
ValueError
|
If any forecast row has a negative lead time, meaning |
Source code in packages/ml_core/src/ml_core/metrics.py
306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 | |
build_mlflow_aggregate_metrics(metrics_df)
Return a flat {metric_key: value} dict for mlflow.log_metrics.
Computes mean metric_value across all series ("all" aggregate), per
time_series_type, and per horizon_slice. The key token is {metric_name} for
metric_param="all" metrics. For parametric metrics the key token is
{metric_name}_{metric_param}, as in pinball_loss_p10 and picp_p10_p90. Parametric
metrics are restricted to the _MLFLOW_LOGGED_PARAMETRIC headline subset. Key formats:
"{token}__all"— overall aggregate (horizon_slice="all")."{token}__{type_slug}"— per-type aggregates (horizon_slice="all")."{token}__all__{horizon_slice}"— overall aggregate per lead-time band (e.g."nmae__all__day_ahead"). Per-type sliced aggregates are deliberately not logged. The per-type mean for each lead-time band stays queryable in theforecast_metricsDelta table, as do the pinball loss at all 13 delivery quantiles and the PICP and interval width of all 6 bands.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
metrics_df
|
DataFrame
|
Per-series |
required |
Returns:
| Type | Description |
|---|---|
dict[str, float]
|
Flat dict of MLflow metric key → mean float value. |
Source code in packages/ml_core/src/ml_core/metrics.py
572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 | |
enrich_metrics_rows(per_series_metrics, experiment_name, evaluation_scope, window_start, window_end, window_label, computed_at, mlflow_run_id)
Add scope and evaluation-window provenance columns to a per-series Metrics frame.
Called by the metrics Dagster asset after compute_metrics() returns, once the
window bounds and MLflow run ID are known. Kept here so the enrichment logic is
unit-testable without Dagster.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
per_series_metrics
|
DataFrame[Metrics]
|
Frame produced by |
required |
experiment_name
|
str
|
Experiment that produced these forecasts. |
required |
evaluation_scope
|
EvalScopeType
|
|
required |
window_start
|
datetime
|
Inclusive start of the evaluated |
required |
window_end
|
datetime
|
Inclusive end of the evaluated |
required |
window_label
|
str
|
Human-readable label ( |
required |
computed_at
|
datetime
|
UTC timestamp when this metric batch was computed. |
required |
mlflow_run_id
|
str | None
|
MLflow fold run ID; |
required |
Returns:
| Type | Description |
|---|---|
DataFrame[Metrics]
|
A validated |
Source code in packages/ml_core/src/ml_core/metrics.py
632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 | |
ml_core.cv_helpers
Pure, data-only helpers for the cross-validation Dagster assets.
Every function here is deliberately free of I/O (no Delta, MLflow, or Dagster imports) so it can be unit-tested in isolation. The CV asset bodies stay thin by delegating their logic here.
A fold is one train-and-validate split of the history. Each fold is a CvFoldConfig carrying a
fold_id plus two pairs of calendar dates: train_start/train_end bounding its training
window, and val_start/val_end bounding its validation window. Cross-validation scores one
model over several such folds.
Attributes
CV_PARTITION_KEY_SEPARATOR = '__'
module-attribute
Separator between experiment name and fold id in a CV partition key.
A double-underscore reduces collision risk with experiment names that contain single underscores.
Classes
Functions:
date_to_utc_datetime(d, *, end_of_day=False)
Return a tz-aware UTC datetime at the start (or inclusive end) of the given date.
Used to turn a fold's [start, end] calendar dates into the [start 00:00:00, end
23:59:59] UTC window. The training and validation data loads filter on that window, and
eligibility uses the same conversion for val_end. That window is closed rather than
half-open: both ends are inclusive, so a filter written against the window needs no < where
the caller meant <=.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
d
|
date
|
The calendar date. |
required |
end_of_day
|
bool
|
If True, return |
False
|
Returns:
| Type | Description |
|---|---|
datetime
|
|
datetime
|
Always the same calendar date as |
Source code in packages/ml_core/src/ml_core/cv_helpers.py
27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 | |
eligible_time_series_ids(coverage, fold, min_training_months)
Return the sorted time_series_ids eligible for a fold, from data coverage alone.
A time series is eligible when it has at least min_training_months of observations
before the fold's val_start and observations through the fold's val_end.
Eligibility is a function of the data only: eligibility depends on neither the model nor the experiment config. Every experiment therefore evaluates a fold on the identical population. That identical population is what makes leaderboard comparisons fair.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
coverage
|
DataFrame
|
One row per time series with the columns |
required |
fold
|
CvFoldConfig
|
The CV fold whose eligibility is being computed. |
required |
min_training_months
|
int
|
Minimum months of pre- |
required |
Returns:
| Type | Description |
|---|---|
list[int]
|
Sorted list of eligible |
Source code in packages/ml_core/src/ml_core/cv_helpers.py
63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 | |
parse_cv_partition_key(partition_key)
Return (experiment_name, fold_id) from a CV partition key.
Partition key format: "{experiment_name}__{fold_id}". The separator is a
double-underscore to reduce collision risk with experiment names that contain single
underscores. Splitting from the right keeps __ inside the experiment name intact.
Source code in packages/ml_core/src/ml_core/cv_helpers.py
96 97 98 99 100 101 102 103 104 | |
flatten_config(config, prefix='')
Flatten a (possibly nested) config into dotted-key string params for MLflow.
MLflow params are scalar strings, so nested dicts are flattened with dotted keys and every
leaf value is stringified. A pydantic BaseModel is first dumped to a plain dict.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
config
|
BaseModel | dict[str, Any]
|
A pydantic model or (possibly nested) dict to flatten. |
required |
prefix
|
str
|
Internal recursion prefix; callers should not set it. |
''
|
Returns:
| Type | Description |
|---|---|
dict[str, str]
|
A flat |
Source code in packages/ml_core/src/ml_core/cv_helpers.py
107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 | |
ml_core.mlflow_runs
Idempotent MLflow run-resolution helpers for every cross-validation and registration path.
Single concern: resolving MLflow runs by tag. MLflow's own objects are the vocabulary here. An experiment holds runs. One run is one recorded job, holding params, metrics, tags, and artifact files under a run id. A tag is a rewritable label on a run or on an experiment, and a run can be reopened later by its id. A Dagster asset is one named, schedulable unit of computation producing one stored output, and a partition is one slice of an asset that is materialised and re-run on its own.
The cross-validation (CV) assets are trained_cv_model, cv_power_forecasts, and
metrics. Those assets and register_experiment_job run in separate processes, and on
retries at separate times. A live Run handle therefore cannot be passed between them.
Instead, every code path discovers or creates the run it needs by tag and resumes it by ID.
Each helper returns an ID string, never an open run handle. The caller wraps the returned ID
in with mlflow.start_run(run_id=...) to log and close the run within its own process. Because
lookup is by tag, every helper is safe to call from any process, and idempotent under Dagster
retries. A re-run resumes the same run rather than duplicating the run.
The tracking URI is set by the caller (mlflow.set_tracking_uri(...)). These helpers simply
use mlflow / MlflowClient, which honour whatever URI is in effect. Honouring the ambient
URI is what lets the helpers run against a file-based MLflow in tests with no server.
Attributes
CV_PARENT_RUN_NAME = 'cv_summary'
module-attribute
Run name of an experiment's parent run, which carries the aggregate leaderboard metrics.
Classes
PromotableRun
dataclass
One MLflow fold run (cv_role=fold) — a valid promoted_model promotion candidate.
Source code in packages/ml_core/src/ml_core/mlflow_runs.py
160 161 162 163 164 165 166 167 | |
Attributes
run_id
instance-attribute
experiment_name
instance-attribute
fold_id
instance-attribute
start_time
instance-attribute
Methods:
__init__(run_id, experiment_name, fold_id, start_time)
Functions:
load_experiment_forecaster(experiment_name)
Reconstruct the forecaster class + resolved config from an experiment's MLflow tags.
register_experiment is the function register_experiment_job runs. That function stamps
the config JSON and the forecaster class's fully-qualified import path (forecaster_target)
on the experiment. The config JSON carries no class identity of its own. That JSON is therefore
deserialised into the BaseForecasterConfig subclass reached through the forecaster's
CONFIG_CLASS. CONFIG_CLASS is the same class the forecaster's load uses. The caller
is responsible for setting the tracking URI (mlflow.set_tracking_uri) beforehand.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
experiment_name
|
str
|
The MLflow experiment name. The CV assets are partitioned per experiment and fold, and each partition key is the experiment name followed by the fold id, so the experiment name is also that key's prefix. |
required |
Returns:
| Type | Description |
|---|---|
type[BaseForecaster]
|
A |
BaseForecasterConfig
|
|
Source code in packages/ml_core/src/ml_core/mlflow_runs.py
39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 | |
get_or_create_experiment(experiment_name)
Return the MLflow experiment id for experiment_name, creating it if absent.
get_or_create_experiment is the self-healing fallback path. The canonical creator is
register_experiment_job, which also stamps the experiment's config/description tags.
Created here untagged. Were the experiment created with tags here instead, a stray asset call
running alongside register_experiment_job could leave an experiment that exists but carries
none of the job's tags — half-registered.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
experiment_name
|
str
|
Human-readable, unique experiment name. |
required |
Returns:
| Type | Description |
|---|---|
str
|
The MLflow experiment id. |
Source code in packages/ml_core/src/ml_core/mlflow_runs.py
69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 | |
get_or_create_parent_run(experiment_id)
Return the experiment's parent run id (tag cv_role=parent), creating it if absent.
The parent run holds the experiment's aggregate (leaderboard) metrics and the flattened config params. Resolved by tag so any process finds the same run.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
experiment_id
|
str
|
The MLflow experiment id. |
required |
Returns:
| Type | Description |
|---|---|
str
|
The parent run id. |
Source code in packages/ml_core/src/ml_core/mlflow_runs.py
90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 | |
get_or_create_fold_run(experiment_id, parent_run_id, fold_id)
Return the fold's child run id (tags cv_role=fold, fold_id=...), creating it if absent.
The fold run holds that fold's per-fold tags and metrics, and never MLflow params. A fold run is
reused across every re-materialisation of its partition, and MLflow params are write-once.
Nothing that can legitimately change between materialisations may therefore be logged as a
param; see trained_cv_model in defs/cv_assets.py. The fold run is created nested
under the experiment's parent run, so the MLflow UI groups folds beneath that parent run's name,
cv_summary. Resolved by (cv_role, fold_id) so any process — and any Dagster retry of the
fold — finds the same run.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
experiment_id
|
str
|
The MLflow experiment id. |
required |
parent_run_id
|
str
|
The experiment's parent run id (from |
required |
fold_id
|
str
|
The fold identifier (e.g. |
required |
Returns:
| Type | Description |
|---|---|
str
|
The fold run id. |
Source code in packages/ml_core/src/ml_core/mlflow_runs.py
118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 | |
list_promotable_runs()
List up to 1000 fold runs (cv_role=fold) per MLflow experiment, newest first.
A read-only convenience for the promotable_model_runs asset
(defs/production_assets.py), which logs the returned list as a metadata table in the
Dagster UI. The experiment listing takes MLflow's default page size of 1000. The
per-experiment fold-run search passes the same 1000 explicitly. The listing therefore covers
the first 1000 active experiments, and returns at most 1000 fold runs from each experiment. A
promoted_model promotion candidate's run id can then be copy-pasted into that asset's
launchpad rather than retyped from memory. The launchpad is the form in Dagster's user
interface where an asset's run-time configuration is entered before that asset is
materialised. The champion is still picked by eye off the MLflow leaderboard;
list_promotable_runs only lists the candidates. The caller is responsible for setting the
tracking URI (mlflow.set_tracking_uri) beforehand.
Source code in packages/ml_core/src/ml_core/mlflow_runs.py
170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 | |
ml_core.production_helpers
IO-light helpers for production (live) inference.
The live service forecasts on a fixed schedule, and one scheduled forecast is a slot. The
weather input a slot forecasts from is a numerical weather prediction (NWP) run, initialised at
nwp_init_time and carrying several ensemble members, of which member 0 is the unperturbed
control member. A scan below is a lazy Polars query over a Delta table, not a materialised
frame. Every function here is unit-testable in isolation. The two data-shaping helpers
(select_nwp_init_time, build_live_power_frame) take power_fcst_init_time as an
explicit parameter rather than calling datetime.now() internally. A test can therefore pass
any fixed time and get a deterministic result. The two disk/MLflow helpers are
load_forecaster_from_dir and fetch_model_artifacts. Both do the IO. Both also check that
this code can still build a config for the saved model, and can still parse that model's
features. Of the five helpers this module exports, weather_lags_lack_their_control_member is
the only helper that executes a read against a data table rather than against a saved model. The
read is a bounded probe against the slot's NWP scan, and the scan is an argument, so a test can
pass an in-memory frame. The live_forecasts and promoted_model Dagster assets
(src/nged_substation_forecast/defs/production_assets.py) stay thin shells over these helpers.
Attributes
AvailabilityModeType = Literal['live', 'replay']
module-attribute
Which NWP-availability rule select_nwp_init_time applies.
"live": the scheduled path. No modelled publication delay — the Delta table only contains runs that have genuinely been published, so the cutoff ispower_fcst_init_timeitself."replay": re-running a past slot over the data the table holds today. The cutoff ispower_fcst_init_time - nwp_publication_delay_hours, reconstructing what was actually available at that historicalpower_fcst_init_time. Subtracting the delay is what makes the cutoff earlier than the live cutoff. Without the subtraction we would leak runs that only landed afterwards.
Classes
Functions:
select_nwp_init_time(available_init_times, *, power_fcst_init_time, availability_mode, nwp_publication_delay_hours=NWP_PUBLICATION_DELAY_HOURS)
Return the freshest NWP init_time available at power_fcst_init_time.
Which runs count as available depends on availability_mode.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
available_init_times
|
Sequence[datetime]
|
The |
required |
power_fcst_init_time
|
datetime
|
The scheduled forecast time (the partition's window end). |
required |
availability_mode
|
AvailabilityModeType
|
|
required |
nwp_publication_delay_hours
|
int
|
Only used in |
NWP_PUBLICATION_DELAY_HOURS
|
Returns:
| Type | Description |
|---|---|
datetime
|
The freshest |
Raises:
| Type | Description |
|---|---|
ValueError
|
If no available |
Source code in packages/ml_core/src/ml_core/production_helpers.py
58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 | |
weather_lags_lack_their_control_member(nwp, *, selected_features)
Whether this slot asks for weather lags that its NWP run cannot fully supply.
A weather model's analysis is its best estimate of the weather that actually happened. This
pipeline has no analysis field, so it approximates an analysis with the control member of the
freshest run — the analysis proxy. A weather lag reaching back before power_fcst_init_time
is answered by that analysis proxy, which reads the control member (ensemble_member == 0)
alone. A run may carry no control-member rows at all, which happens when an ECMWF ENS download
is partial or malformed. Such a run therefore nulls each weather lag over the first
lag_hours of the horizon, where the lag still points into the past. The rest of the horizon
is answered by the same-run join, which reads whichever ensemble members the run does carry.
_engineer_features already degrades rather than failing there and logs a warning naming the
run. This function exists so live_forecasts can also report the degradation on the Sentry
channel. Rule
4
requires that Sentry report alongside the log, not as a substitute for the log. Reading the logs
is not a monitoring strategy: the operator reads the alert.
Returns False when the model selects no weather lag. A model with no weather lags never
reads the control member, because the same-run weather join reads whichever ensemble members the
run does carry.
H3 is a hexagonal grid over the globe, and NWP is stored one row per H3 cell, so each time
series is matched to the cell it sits in. live_forecasts probes the NWP scan before that
H3 spatial join, while _engineer_features probes the scan after. The earlier probe point
makes the frame here a superset of the frame the pipeline sees. The alert can therefore in
principle miss a degradation the pipeline hits, but can never fire on a degradation the pipeline
does not hit. In practice load_engineering_inputs has already pruned the scan to the model's
own frozen H3 cells — frozen meaning the copy the model saved at training time, in
TRAINED_METADATA_FILENAME — so the two frames hold the same cells today.
The probe reads at most one row, so a healthy NWP run answers the probe from a single row group. Proving an NWP run has no control member requires scanning the whole run, because absence cannot be shown early.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
nwp
|
LazyFrame[Nwp]
|
The slot's NWP rows, already narrowed to the selected run. |
required |
selected_features
|
set[str]
|
The promoted model's feature names. |
required |
Returns:
| Type | Description |
|---|---|
bool
|
|
Source code in packages/ml_core/src/ml_core/production_helpers.py
99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 | |
build_live_power_frame(observed_power, time_series_ids, *, power_fcst_init_time, history, horizon)
Build a dense half-hourly (time_series_id, time) spine for live inference.
Needed because ml_core.features._nwp._join_nwp_single_run is power-centric — with no future
power rows a live run would emit zero forecast rows. Left-joins observed power onto a spine
covering (power_fcst_init_time - history, power_fcst_init_time + horizon] for every
requested time_series_id. Rows beyond the last observation are therefore present with
power = null. The spine is also harmless for replay. Future observations already exist in
replay, and _nullify_leaky_lags prevents lag leakage regardless.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
observed_power
|
LazyFrame[PowerTimeSeries]
|
Lazy observed power, one row per |
required |
time_series_ids
|
list[int]
|
The series to build a spine for (typically
|
required |
power_fcst_init_time
|
datetime
|
The forecast init time. The spine's window is anchored on this. |
required |
history
|
timedelta
|
How far before |
required |
horizon
|
timedelta
|
How far after |
required |
Returns:
| Type | Description |
|---|---|
LazyFrame[PowerTimeSeries]
|
A lazy |
LazyFrame[PowerTimeSeries]
|
half-hourly grid, observed values joined in, future/missing rows null. |
Source code in packages/ml_core/src/ml_core/production_helpers.py
154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 | |
load_forecaster_from_dir(path)
Load the production model from a plain disk directory (no MLflow at inference time).
Reads meta.json and resolves model_class via contracts.config_schemas.import_class
(the same mechanism ml_core.mlflow_runs.load_experiment_forecaster uses), then calls the
concrete subclass's load(path).
The forecaster returned is a model this code can actually serve, not merely a model this code could deserialise. A config this code cannot rebuild is rejected here, as is a feature vocabulary this code cannot parse. Rejecting here beats failing partway through a live tick's feature engineering.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
path
|
Path
|
Directory previously populated by |
required |
Returns:
| Type | Description |
|---|---|
BaseForecaster
|
The reconstructed, trained forecaster. |
Raises:
| Type | Description |
|---|---|
FileNotFoundError
|
|
ValueError
|
This code cannot serve the saved model — see |
Source code in packages/ml_core/src/ml_core/production_helpers.py
297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 | |
fetch_model_artifacts(run_id, dest)
Download and unpack an MLflow run's saved model into dest, replacing it atomically.
Downloads and unpacks into a temporary directory first, so a failed or interrupted download
never touches dest. Only a fully-downloaded model is moved into place, via rmtree +
move. dest is local disk by convention — Settings.production_model_path derives from
local_artifacts_path, though nothing enforces that. Unlike a Delta table, dest is a
directory of many files with no commit protocol over it, and a part-written directory would be
served. The run holds the model as a single archive artifact
(ml_core.base_forecaster._MLFLOW_MODEL_ARTIFACT). dest therefore gets exactly the files
the last save_to_mlflow wrote, and can never inherit a stale file from an earlier, larger
model.
The downloaded model's saved config is checked against the running code before the swap,
reading the staged meta.json rather than loading the model. A model this code cannot serve
is therefore refused while the previous champion stays in dest and keeps serving. Reading
the JSON is deliberate: it applies the same validation the subclass's load would apply,
without pulling every booster into memory first.
Also writes a promotion.json ({"mlflow_run_id", "promoted_at"}) into dest for
provenance. That extra file is harmless, because a BaseForecaster.load implementation reads
its own population from its saved record, never from a directory listing. XGBoostForecaster,
for example, reads the population from meta.json's trained_time_series_ids.
The caller is responsible for setting the tracking URI (mlflow.set_tracking_uri)
beforehand.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
run_id
|
str
|
The MLflow run the model was saved under (via |
required |
dest
|
Path
|
Directory to populate — typically |
required |
Raises:
| Type | Description |
|---|---|
MlflowException
|
|
ValueError
|
The run holds no |
Source code in packages/ml_core/src/ml_core/production_helpers.py
336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 | |
ml_core.repro
Reproducibility provenance: the git SHA and Delta-table versions behind an MLflow run.
Answers "exactly which code and which data produced this?" for any MLflow run. The git SHA pins
the code. The SHA is stamped explicitly because MLflow's mlflow.source.git.commit
auto-detection needs gitpython installed and the working directory inside the repo. Neither
condition holds in a production container. Each Delta table's version() pins the data: Delta
Lake time travel makes data versioning one integer per table. A run can therefore later be
replayed with pl.scan_delta(path, version=N) after git checkout {sha}.
Every function here is deliberately non-raising: the git SHA, the dirty flag, and each Delta
table's version are a record about a run rather than an input to that run. The surrounding
training or forecasting run must never fail because of a missing .git directory (containers)
or an absent Delta table. Each absence degrades to the sentinels "unknown" / "absent"
instead.
provenance_tags stage-prefixes its keys (register_, train_, predict_,
metrics_) because four separate writers stamp provenance onto the same MLflow runs. Three are
Dagster assets writing one fold run: trained_cv_model, cv_power_forecasts, and
metrics. Each of the three can be on a different code revision and at a different set of
Delta table versions. The fourth is the register_experiment op inside
register_experiment_job, which stamps the experiment's parent run; the metrics asset
stamps that parent run too. Un-prefixed keys would clobber one another, and the prefix preserves
every writer's provenance snapshot side by side.
Attributes
logger = logging.getLogger(__name__)
module-attribute
MlflowTags = dict[str, str]
module-attribute
A {tag_key: tag_value} mapping ready to hand to mlflow.set_tags.
StageType = Literal['register', 'train', 'predict', 'metrics']
module-attribute
The four stages that stamp provenance; each value becomes a tag-key prefix (see
provenance_tags). "train", "predict", and "metrics" are Dagster assets, and
"register" is the register_experiment op. Add a new StageType value when a new asset or
op starts stamping.
TableNameType = Literal['power_time_series', 'nwp_data', 'eligible_time_series', 'power_forecasts', 'effective_capacity']
module-attribute
Logical names of the Delta tables whose versions get stamped — the keys of a delta_paths
mapping. Add a new TableNameType value when a stage starts reading (and stamping) another
table.
UNKNOWN = 'unknown'
module-attribute
Sentinel git SHA / dirty flag returned when no git repository is reachable (e.g. a container).
ABSENT = 'absent'
module-attribute
Sentinel Delta version returned when a table does not exist or cannot be read.
Classes
Functions:
get_git_info(cwd=None)
Return {"git_sha", "git_dirty"} for the current checkout, never raising.
git_dirty is "true" when the working tree has uncommitted changes, "false" when
clean. When the SHA cannot be read (no .git, no git binary, a timeout) both values are
UNKNOWN. When the SHA is read but the dirty check fails, the SHA is kept and only
git_dirty degrades to UNKNOWN, because a good SHA is never discarded.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
cwd
|
Path | None
|
Directory the |
None
|
Returns:
| Type | Description |
|---|---|
MlflowTags
|
|
MlflowTags
|
commit hash, and |
MlflowTags
|
|
Source code in packages/ml_core/src/ml_core/repro.py
81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 | |
get_delta_versions(paths, storage_options=None)
Return {f"delta_version__{name}": str(version)} for each named Delta table.
Never raises.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
paths
|
dict[TableNameType, str]
|
|
required |
storage_options
|
ObjectStoreOptions | None
|
object-store options for a remote URI; |
None
|
Returns:
| Type | Description |
|---|---|
MlflowTags
|
One entry per input path. A table that does not exist (or cannot be read) maps to |
MlflowTags
|
|
Source code in packages/ml_core/src/ml_core/repro.py
124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 | |
provenance_tags(stage, delta_paths=None, storage_options=None)
Build stage-prefixed MLflow tags stamping code (git) and data (Delta version) provenance.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
stage
|
StageType
|
Prefix identifying the writing asset or op. Keeps the tags of the assets and ops that share one MLflow run from clobbering. |
required |
delta_paths
|
dict[TableNameType, str] | None
|
|
None
|
storage_options
|
ObjectStoreOptions | None
|
object-store options for remote table URIs; |
None
|
Returns:
| Type | Description |
|---|---|
MlflowTags
|
e.g. for |
MlflowTags
|
"train_delta_version__power_time_series", ...}``. |
Source code in packages/ml_core/src/ml_core/repro.py
157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 | |