Skip to content

Delta Store API

delta_store

Physical storage policy for the project's Delta tables.

Why this package exists

contracts owns each table's logical shape and meaning. This package owns its physical layout: parquet writer properties (codec + per-column encodings), compression-friendly sort orders, and significand-precision rounding, plus the write helpers that apply them. Dagster assets stay thin by writing through this package rather than calling write_deltalake with ad-hoc settings — and it becomes impossible to land rows in a table without its storage format applied.

Every lever applied by power_forecasts and nwp below is measured against real data rather than assumed — see Storage formats: measured, not assumed for the comparison between those two tables, and design principle 12 for why measuring rather than assuming matters project-wide. This package covers six tables. The other four carry no writer-properties tuning, because no measurement has yet been made to justify any tuning.

The flagship tuned table is the internal power_forecasts table, whose storage format shrank a 403.6M-row development copy of that table from 6.33 GB to 0.73 GB. The levers are ZSTD, DELTA_BINARY_PACKED timestamps, BYTE_STREAM_SPLIT floats, member-adjacent sorting (the ~51 ensemble members of one forecast target on adjacent rows), and rounding power_fcst to a 13-bit significand. POWER_FORECASTS_WRITER_PROPERTIES, below on this page, gives the lever-by-lever breakdown, measured on the table's least compressible single file rather than on the full table.

Contents

  • precision.round_to_significand_bits() — rounds a Float32 expression to a chosen number of significand bits in pure Polars arithmetic (Veltkamp splitting), zeroing the low mantissa bits so BYTE_STREAM_SPLIT + zstd can compress them away. The trick and its preconditions are rigorously documented on the function.
  • power_forecasts — the power_forecasts table's writer properties, sort order, precision policy, and write_power_forecasts().
  • nwp — the nwp table's writer properties, sort order, precision policy, and write_nwp(); its writer properties are deliberately different from power_forecasts's, because the same encodings measured worse on numerical weather prediction (NWP) data.
  • power_time_series — write_power_time_series(), an append-only write to the power_time_series table.
  • eligible_time_series — write_eligible_time_series(), a per-fold_id-partition overwrite to the eligible_time_series table.
  • effective_capacity — write_effective_capacity(), a whole-table overwrite to the effective_capacity table.
  • forecast_metrics — write_forecast_metrics(), a per-(experiment_name, fold_id)-partition overwrite to the forecast_metrics table, including the Enum→String cast delta-rs needs before writing.

delta_store.precision

Reduce the significand precision of Float32 columns so parquet compression can work.

In a large forecast or weather table, nearly every full-precision Float32 value is distinct, so the low significand bits are incompressible noise — they defeat every general-purpose codec. Rounding each value to a small number of significand bits zeroes those low bits. A general-purpose codec then has repetition to find, at the cost of a strictly bounded relative error. Which codec then wins is table-dependent, and is measured per table. delta_store.power_forecasts adds a BYTE_STREAM_SPLIT encoding on top of zstd; see that module for measured numbers. delta_store.nwp measured BYTE_STREAM_SPLIT as worse on NWP data, and uses zstd alone. The rounding is what both tables have in common.

The rounding here is pure Polars arithmetic, with no bit-twiddling and no numpy round-trip. The rounding works because of a classical floating-point identity, Veltkamp splitting, documented in detail on round_to_significand_bits.

Attributes

FLOAT32_SIGNIFICAND_BITS = 24 module-attribute

IEEE 754 binary32 significand precision p: 23 explicit fraction bits + 1 implicit bit.

Functions:

round_to_significand_bits(expr, *, keep_bits)

Round a Float32 expression to keep_bits significand bits (round-to-nearest).

The result is exactly representable with a keep_bits-bit significand. Write p for the format's significand width, which FLOAT32_SIGNIFICAND_BITS above sets to 24. Every finite output therefore has its low 24 - keep_bits explicit fraction bits set to zero, which is what gives the compression codec repetition to find. The relative error is bounded by the unit roundoff of a keep_bits-bit format:

|result - x| <= 2**-keep_bits * |x|

(e.g. keep_bits=13 -> max relative error 2^-13 ~= 1.2e-4). Rounding is to nearest, so unlike truncation it introduces no systematic bias toward zero.

How it works — Veltkamp splitting. With s = 24 - keep_bits and the constant C = 2**s + 1, the expression computes, entirely in Float32 round-to-nearest arithmetic (RN)::

c = RN(x * C)  # = RN(x*2^s + x)
result = RN(c - RN(c - x))

Veltkamp's theorem (Veltkamp 1968; Dekker (1971), "A floating-point technique for extending the available precision", Numerische Mathematik 18; Muller et al. (2018), Handbook of Floating-Point Arithmetic, 2nd ed., §4.4, Algorithm 4.9 "Split") states that for 2 <= s <= p - 2 and no overflow, both subtractions are exact and result is x rounded to nearest onto p - s = keep_bits significand bits. The intuition: x*C stacks a copy of x shifted s exponent positions above itself. Rounding that sum to p bits, which forms c, is precisely what discards the low s bits of the original x, and discards them with correct rounding. The two exact subtractions then peel the shifted copy back off, leaving x with its low s significand bits rounded away.

Preconditions — each one is load-bearing:

  • expr must be Float32. The splitting constant is pinned to Float32, but Polars promotes mixed arithmetic. If a Float64 expression is passed, the whole computation runs at p = 53 and silently keeps 53 - s bits instead of keep_bits. A .cast(pl.Float32) is deliberately not applied here. Silently changing the dtype of a Float64 column would be its own trap, and callers own their dtypes.
  • Arithmetic must be evaluated operation-by-operation in IEEE round-to-nearest. Polars (Rust) guarantees operation-by-operation evaluation: floats are never reassociated, and a*b + c is never contracted into an FMA. Reassociation or contraction into an FMA would break the exactness of the subtractions.
  • Non-finite and overflow cases are guarded. For |x| > f32_max / C the product c overflows to inf, and the naive result would be inf - inf = NaN. For x NaN or ±inf, c is likewise non-finite. Wherever c is non-finite, the input value is passed through unchanged. A huge-but-finite value is therefore stored at full precision rather than corrupted, and NaN/±inf survive verbatim.
  • Subnormal x (|x| < 2**-126) degrades gracefully. The relative-error bound loosens, because precision is already at the absolute floor of the format. The result is still a faithful nearby value.

Parameters:

Name Type Description Default
expr Expr

A Float32 expression (see the dtype precondition above).

required
keep_bits int

Significand bits to keep, in [2, 22] (the theorem's 2 <= s <= p - 2). Note this counts significand bits: the implicit leading 1 plus the explicit fraction bits. keep_bits=13 therefore keeps 12 explicit fraction bits.

required

Returns:

Type Description
Expr

A Float32 expression: expr rounded to keep_bits significand bits, with

Expr

non-finite and overflowing inputs passed through unchanged.

Source code in packages/delta_store/src/delta_store/precision.py
25
26
27
28
29
30
31
32
33
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
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
def round_to_significand_bits(expr: pl.Expr, *, keep_bits: int) -> pl.Expr:
    """Round a ``Float32`` expression to ``keep_bits`` significand bits (round-to-nearest).

    The result is exactly representable with a ``keep_bits``-bit significand. Write ``p`` for the
    format's significand width, which `FLOAT32_SIGNIFICAND_BITS` above sets to 24. Every finite
    output therefore has its low ``24 - keep_bits`` explicit fraction bits set to zero, which is
    what gives the compression codec repetition to find. The relative error is bounded by the unit
    roundoff of a ``keep_bits``-bit format:

        |result - x| <= 2**-keep_bits * |x|

    (e.g. ``keep_bits=13`` -> max relative error 2^-13 ~= 1.2e-4). Rounding is to nearest, so
    unlike truncation it introduces no systematic bias toward zero.

    **How it works — Veltkamp splitting.** With ``s = 24 - keep_bits`` and the constant
    ``C = 2**s + 1``, the expression computes, entirely in ``Float32`` round-to-nearest
    arithmetic (``RN``)::

        c = RN(x * C)  # = RN(x*2^s + x)
        result = RN(c - RN(c - x))

    Veltkamp's theorem (Veltkamp 1968; [Dekker (1971)](https://doi.org/10.1007/BF01397083), "A
    floating-point technique for extending the available precision", *Numerische Mathematik* 18;
    [Muller et al. (2018)](https://doi.org/10.1007/978-3-319-76526-6), *Handbook of Floating-Point
    Arithmetic*, 2nd ed., §4.4, Algorithm 4.9 "Split") states that for
    ``2 <= s <= p - 2`` and no overflow, both subtractions are **exact** and ``result`` is ``x``
    rounded to nearest onto ``p - s = keep_bits`` significand bits. The intuition: ``x*C`` stacks a
    copy of ``x`` shifted ``s`` exponent positions above itself. Rounding that sum to ``p`` bits,
    which forms ``c``, is precisely what discards the low ``s`` bits of the original ``x``, and
    discards them with correct rounding. The two exact subtractions then peel the shifted copy back
    off, leaving ``x`` with its low ``s`` significand bits rounded away.

    **Preconditions — each one is load-bearing:**

    - ``expr`` must be ``Float32``. The splitting constant is pinned to ``Float32``, but Polars
      promotes mixed arithmetic. If a ``Float64`` expression is passed, the whole computation runs
      at ``p = 53`` and silently keeps ``53 - s`` bits instead of ``keep_bits``. A
      ``.cast(pl.Float32)`` is deliberately *not* applied here. Silently changing the dtype of a
      ``Float64`` column would be its own trap, and callers own their dtypes.
    - Arithmetic must be evaluated operation-by-operation in IEEE round-to-nearest. Polars (Rust)
      guarantees operation-by-operation evaluation: floats are never reassociated, and ``a*b + c``
      is never contracted into an FMA. Reassociation or contraction into an FMA would break the
      exactness of the subtractions.
    - Non-finite and overflow cases are guarded. For ``|x| > f32_max / C`` the product ``c``
      overflows to ``inf``, and the naive result would be ``inf - inf = NaN``. For ``x`` NaN or
      ``±inf``, ``c`` is likewise non-finite. Wherever ``c`` is non-finite, the input value is
      passed through unchanged. A huge-but-finite value is therefore stored at full precision
      rather than corrupted, and ``NaN``/``±inf`` survive verbatim.
    - Subnormal ``x`` (``|x| < 2**-126``) degrades gracefully. The relative-error bound loosens,
      because precision is already at the absolute floor of the format. The result is still a
      faithful nearby value.

    Args:
        expr: A ``Float32`` expression (see the dtype precondition above).
        keep_bits: Significand bits to keep, in ``[2, 22]`` (the theorem's ``2 <= s <= p - 2``).
            Note this counts *significand* bits: the implicit leading 1 plus the explicit fraction
            bits. ``keep_bits=13`` therefore keeps 12 explicit fraction bits.

    Returns:
        A ``Float32`` expression: ``expr`` rounded to ``keep_bits`` significand bits, with
        non-finite and overflowing inputs passed through unchanged.
    """
    shift = FLOAT32_SIGNIFICAND_BITS - keep_bits
    if not 2 <= shift <= FLOAT32_SIGNIFICAND_BITS - 2:
        raise ValueError(
            f"keep_bits must be in [2, {FLOAT32_SIGNIFICAND_BITS - 2}], got {keep_bits}"
        )
    splitter = pl.lit(float(2**shift + 1), dtype=pl.Float32)
    c = expr * splitter
    rounded = c - (c - expr)
    return pl.when(c.is_finite()).then(rounded).otherwise(expr)

delta_store.power_forecasts

Storage policy for the internal power_forecasts Delta table.

Owns everything about how PowerForecast rows are laid out on disk: the parquet writer properties (codec + per-column encodings), the compression-friendly row order, and the power_fcst precision reduction. Callers write through write_power_forecasts so it is impossible to land rows in the table without this format applied.

Measured impact: the full 403.6M-row development table shrank from 6.33 GB to 0.73 GB when rewritten into this format. The 6.33 GB figure is the delta-rs defaults: SNAPPY, dictionary encoding, unsorted, full precision. See the POWER_FORECASTS_WRITER_PROPERTIES docstring for the per-lever breakdown.

Attributes

POWER_FCST_SIGNIFICAND_BITS = 13 module-attribute

Significand bits kept when storing power_fcst (1 implicit + 12 explicit fraction bits).

Rounding to nearest at 13 significand bits caps the relative error at 2⁻¹³ ≈ 1.2×10⁻⁴, orders of magnitude below forecast error. The rounding also zeroes the 11 low fraction bits. Those 11 bits are otherwise pure entropy that defeats every compression codec, because nearly every full-precision power_fcst value is distinct. Mirrors the significand-rounding scheme delta_store.nwp applies to NWP data, but with different writer properties. See that module's docstring for why the choice doesn't transfer between tables. See POWER_FORECASTS_WRITER_PROPERTIES for the measured size impact.

POWER_FORECASTS_SORT_COLS = PowerForecast.PRIMARY_KEY module-attribute

Within-file row order for power_forecasts writes.

power_fcst becomes locally smooth, and the timestamp columns become stepped sequences, when the ~51 ensemble members of one (series, init time, valid time) target sit on adjacent rows. Locally smooth values and stepped sequences are exactly what the BYTE_STREAM_SPLIT and DELTA_BINARY_PACKED encodings in POWER_FORECASTS_WRITER_PROPERTIES need to compress well. Leading with time_series_id also lets parquet row-group statistics prune scans that filter on one series.

Defined as PowerForecast.PRIMARY_KEY rather than repeating its columns, because the set that identifies a row uniquely is exactly the set whose adjacency makes a row compress. The order is this module's own choice, and test_sort_cols_lead_with_time_series_id pins that order. A reordered primary key is still a correct key but a worse layout, so this tuple is not free to follow a reordered primary key.

POWER_FORECASTS_WRITER_PROPERTIES = WriterProperties(compression='ZSTD', compression_level=3, column_properties={'valid_time': _TIMESTAMP_COLUMN_PROPERTIES, 'power_fcst_init_time': _TIMESTAMP_COLUMN_PROPERTIES, 'nwp_init_time': _TIMESTAMP_COLUMN_PROPERTIES, 'power_fcst': ColumnProperties(encoding='BYTE_STREAM_SPLIT', dictionary_enabled=False)}) module-attribute

Parquet writer settings for the power_forecasts Delta table.

power_fcst is ~76% of the bytes. The delta-rs defaults (SNAPPY + dictionary encoding everywhere) leave that column essentially uncompressed. Measured on the table's largest single file (18.3M rows, 105 MB): ZSTD-3 alone → 80% of the original size; adding these column encodings plus the POWER_FORECASTS_SORT_COLS row order → 50%; adding the POWER_FCST_SIGNIFICAND_BITS precision reduction → 29%. That file is the table's least compressible slice. Across the full table the format achieved 11.5%.

Classes

Functions:

write_power_forecasts(forecasts, table_uri, *, replace_partition=None, replace_predicate_extra=None, storage_options=None)

Write PowerForecast rows to the power_forecasts Delta table in its storage format.

Applies the three storage levers before writing: rounds power_fcst to POWER_FCST_SIGNIFICAND_BITS significand bits, sorts rows by POWER_FORECASTS_SORT_COLS, and writes with POWER_FORECASTS_WRITER_PROPERTIES.

The table is partitioned by (experiment_name, fold_id). A multi-chunk materialisation passes replace_partition on its first chunk, and None (append) on the rest. Passing replace_partition on the first chunk overwrites that partition, so prior rows are always replaced, even when the first chunk is empty.

Every row of forecasts must satisfy the predicate the overwrite path builds. delta-rs checks each row against that predicate. If any row falls outside the predicate, delta-rs rejects the whole write — DeltaError, nothing committed — rather than writing the rows that do match. So a frame spanning two (experiment_name, fold_id) pairs fails under a single replace_partition, and a frame spanning two power_fcst_init_time values fails under a narrowing replace_predicate_extra. Confirmed empirically against deltalake 1.6.3: a frame carrying one row outside the predicate raised Invalid data found: 1 rows failed validation check. The table was left at both its previous row count and its previous Delta version. write_nwp documents the same delta-rs behaviour and is no safer. write_nwp's predicate is built from the frame's first row (nwp.item(0, ...)), so a frame spanning two (nwp_model_id, init_time) partitions is rejected there too. Both functions leave it to the caller to pass a frame the predicate covers, and neither checks before writing.

The append path takes no predicate, so the append path performs no row-by-row check. Two identical appends leave two copies, confirmed empirically, because nothing deduplicates on PowerForecast.PRIMARY_KEY at write time. PowerForecast's uniqueness constraint is checked by validate(), on the way in, not by the table.

What stops a retry double-counting is the caller's chunking discipline, not anything in this function. The cross-validation asset cv_power_forecasts restarts its chunk loop from the first chunk on every materialisation, and that first chunk always passes replace_partition. The first chunk therefore clears whatever an interrupted run left in the (experiment_name, fold_id) partition, before the later chunks append. A caller that appended without an overwriting first chunk would silently double the partition's rows instead — see principle 10, every write is atomic and idempotent.

Parameters:

Name Type Description Default
forecasts DataFrame[PowerForecast]

Validated forecast rows. experiment_name and fold_id are String, which is their PowerForecast dtype. String is exactly what delta-rs needs for Hive-style partition directories, so no cast is required.

required
table_uri str | Path

Path or URI of the power_forecasts Delta table.

required
replace_partition tuple[str, str] | None

(experiment_name, fold_id) to overwrite, or None to append. Passed explicitly rather than derived from forecasts so an empty first chunk still clears the partition.

None
replace_predicate_extra str | None

An additional AND-ed SQL predicate clause narrowing the overwrite below the (experiment_name, fold_id) partition. For example, "power_fcst_init_time = '2026-07-04T06:00:00+00:00'" lets live_forecasts replace one 6-hourly slot's rows without wiping the rest of the "live" fold's partition. The production service forecasts every 6 hours, and stores every one of those runs under the reserved fold_id of "live" rather than under a cross-validation fold. delta-rs' replaceWhere supports predicates on non-partition columns (confirmed empirically: a datetime.isoformat() literal round-trips correctly against a Timestamp column). Only meaningful alongside replace_partition.

None
storage_options ObjectStoreOptions | None

delta-rs object-store options (credentials/endpoint) for a remote table_uri; None/empty for a local path.

None
Source code in packages/delta_store/src/delta_store/power_forecasts.py
 80
 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
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
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
def write_power_forecasts(
    forecasts: pt.DataFrame[PowerForecast],
    table_uri: str | Path,
    *,
    replace_partition: tuple[str, str] | None = None,
    replace_predicate_extra: str | None = None,
    storage_options: ObjectStoreOptions | None = None,
) -> None:
    """Write ``PowerForecast`` rows to the ``power_forecasts`` Delta table in its storage format.

    Applies the three storage levers before writing: rounds ``power_fcst`` to
    ``POWER_FCST_SIGNIFICAND_BITS`` significand bits, sorts rows by
    ``POWER_FORECASTS_SORT_COLS``, and writes with ``POWER_FORECASTS_WRITER_PROPERTIES``.

    The table is partitioned by ``(experiment_name, fold_id)``. A multi-chunk materialisation passes
    ``replace_partition`` on its first chunk, and ``None`` (append) on the rest. Passing
    ``replace_partition`` on the first chunk overwrites that partition, so prior rows are always
    replaced, even when the first chunk is empty.

    **Every row of ``forecasts`` must satisfy the predicate the overwrite path builds.** delta-rs
    checks each row against that predicate. If any row falls outside the predicate, delta-rs rejects
    the whole write — ``DeltaError``, nothing committed — rather than writing the rows that do
    match. So a frame spanning two ``(experiment_name, fold_id)`` pairs fails under a single
    ``replace_partition``, and a frame spanning two ``power_fcst_init_time`` values fails under a
    narrowing ``replace_predicate_extra``. Confirmed empirically against ``deltalake`` 1.6.3: a
    frame carrying one row outside the predicate raised ``Invalid data found: 1 rows failed
    validation check``. The table was left at both its previous row count and its previous Delta
    version. ``write_nwp`` documents the same delta-rs behaviour and is no safer. ``write_nwp``'s
    predicate is built from the frame's *first row* (``nwp.item(0, ...)``), so a frame spanning two
    ``(nwp_model_id, init_time)`` partitions is rejected there too. Both functions leave it to the
    caller to pass a frame the predicate covers, and neither checks before writing.

    The append path takes no predicate, so the append path performs no row-by-row check. Two
    identical appends leave two copies, confirmed empirically, because nothing deduplicates on
    ``PowerForecast.PRIMARY_KEY`` at write time. ``PowerForecast``'s uniqueness constraint is
    checked by ``validate()``, on the way in, not by the table.

    **What stops a retry double-counting is the caller's chunking discipline, not anything in this
    function.** The cross-validation asset ``cv_power_forecasts`` restarts its chunk loop from the
    first chunk on every materialisation, and that first chunk always passes ``replace_partition``.
    The first chunk therefore clears whatever an interrupted run left in the ``(experiment_name,
    fold_id)`` partition, before the later chunks append. A caller that appended without an
    overwriting first chunk would silently double the partition's rows instead — see [principle 10,
    every write is atomic and
    idempotent](https://openclimatefix.github.io/nged-substation-forecast/design-philosophy/design-principles/#10-every-write-is-atomic-and-idempotent-and-every-failure-is-confined-to-one-partition).

    Args:
        forecasts: Validated forecast rows. ``experiment_name`` and ``fold_id`` are ``String``,
            which is their ``PowerForecast`` dtype. ``String`` is exactly what delta-rs needs for
            Hive-style partition directories, so no cast is required.
        table_uri: Path or URI of the ``power_forecasts`` Delta table.
        replace_partition: ``(experiment_name, fold_id)`` to overwrite, or ``None`` to append.
            Passed explicitly rather than derived from ``forecasts`` so an *empty* first chunk
            still clears the partition.
        replace_predicate_extra: An additional ``AND``-ed SQL predicate clause narrowing the
            overwrite below the ``(experiment_name, fold_id)`` partition. For example,
            ``"power_fcst_init_time = '2026-07-04T06:00:00+00:00'"`` lets ``live_forecasts`` replace
            one 6-hourly slot's rows without wiping the rest of the ``"live"`` fold's
            partition. The production service forecasts every 6 hours, and stores every one of
            those runs under the reserved ``fold_id`` of ``"live"`` rather than under a
            cross-validation fold. delta-rs' ``replaceWhere`` supports predicates on non-partition
            columns (confirmed empirically: a `datetime.isoformat()` literal round-trips correctly
            against a ``Timestamp`` column). Only meaningful alongside ``replace_partition``.
        storage_options: delta-rs object-store options (credentials/endpoint) for a remote
            ``table_uri``; ``None``/empty for a local path.
    """
    prepared = (
        forecasts.with_columns(
            power_fcst=round_to_significand_bits(
                pl.col("power_fcst"), keep_bits=POWER_FCST_SIGNIFICAND_BITS
            )
        )
        .sort(*POWER_FORECASTS_SORT_COLS)
        .to_arrow()
    )
    if replace_partition is not None:
        experiment_name, fold_id = replace_partition
        predicate = f"experiment_name = '{experiment_name}' AND fold_id = '{fold_id}'"
        if replace_predicate_extra is not None:
            predicate = f"{predicate} AND {replace_predicate_extra}"
        write_deltalake(
            table_or_uri=table_uri,
            data=prepared,
            mode="overwrite",
            predicate=predicate,
            partition_by=["experiment_name", "fold_id"],
            writer_properties=POWER_FORECASTS_WRITER_PROPERTIES,
            storage_options=typeddict_to_dict(storage_options),
        )
    else:
        write_deltalake(
            table_or_uri=table_uri,
            data=prepared,
            mode="append",
            partition_by=["experiment_name", "fold_id"],
            writer_properties=POWER_FORECASTS_WRITER_PROPERTIES,
            storage_options=typeddict_to_dict(storage_options),
        )

delta_store.nwp

Storage policy for the nwp Delta table.

Owns everything about how Nwp rows are laid out on disk: the parquet writer properties, the compression-friendly row order, and the significand-precision reduction of the continuous weather variables. Callers write through write_nwp so it is impossible to land rows in the table without this format applied.

Stores plain Float32 + delta_store.precision.round_to_significand_bits — the technique used by delta_store.power_forecasts, but with different writer properties. Measured on real numerical weather prediction (NWP) data (9 partitions spread across two years of history): BYTE_STREAM_SPLIT made every continuous column larger, not smaller — the opposite of the power_forecasts result. Working hypothesis: significand rounding collapses NWP values into a small set of repeats, because many H3 cells and many ensemble members round to the same value. Parquet's default dictionary+RLE encoding captures those repeats directly. BYTE_STREAM_SPLIT instead scatters that repetition across four separate byte planes, and so loses more than it gains. power_forecasts's target values are near-continuous ML output, and have no such repetition. BYTE_STREAM_SPLIT wins there instead, so the two tables need different writer properties. See https://openclimatefix.github.io/nged-substation-forecast/architecture/performance/#storage-formats-measured-not-assumed for the measured GB/yr numbers, and https://openclimatefix.github.io/nged-substation-forecast/api/dynamical_data/ for the member-early-sort read-speed benchmark.

Attributes

NWP_SIGNIFICAND_BITS = 13 module-attribute

Significand bits kept for every continuous NWP variable (1 implicit + 12 explicit fraction bits) — the same budget as delta_store.power_forecasts.POWER_FCST_SIGNIFICAND_BITS. Caps the relative error at 2⁻¹³ ≈ 1.2×10⁻⁴. Measured max absolute error on real data: ≤0.004 °C for temperature, ≤8 Pa (0.08 hPa) for mean-sea-level pressure — both well inside tolerance (temperature ≤0.25 K, MSL pressure ≤1 hPa).

NWP_SORT_COLS = ('init_time', 'ensemble_member', 'valid_time', 'h3_index') module-attribute

Within-file row order for nwp writes — member before valid_time. That is the opposite priority from power_forecasts, which sorts member-adjacent for a different reason. In power_forecasts member adjacency is about compressing near-duplicate ensemble values; in nwp it is about row-group pruning. Sorting ensemble_member early puts each member's rows in one contiguous block, which _member_aligned_row_group_size then turns into one row group per member. A single-member predicate then matches a single row group's min/max range and skips the rest of the partition instead of decoding those rows. That pruning holds for any member, not only the control member. The sort alone is not enough: row groups that straddle member boundaries advertise the whole span between their extremes.

Measured on the stored table, a single-member read decodes 1.96% of a partition's rows — one row group in 51 — and it does so for every member. A census sampled the stored table's partitions. Every partition the census sampled held 51 row groups, each spanning a single member, and those 51 row groups together covered members 0 to 50. It takes 30 ms and 400 MB of peak resident memory to read 29 daily partitions, nine H3 cells, and the control member alone. The same read of a table sorted valid_time-first takes 170 ms and 2,200 MB. The member-early sort holds 3.7% more stored bytes than the valid_time-first sort. One measurement took a real partition and removed half its rows at random members, so that the member-to-row-group alignment degrades rather than holding exactly. Under that degraded alignment the worst member still reads 5.88% of the partition's rows, against the 1.96% floor above. The method and the full figures live beside the storage measurements in https://openclimatefix.github.io/nged-substation-forecast/api/dynamical_data/.

Two conditions have to hold for the predicate to reach the Parquet scan at all. The predicate must survive Nwp.scan_delta's cast, which requires that cast to be a no-op (see the Nwp.ensemble_member field). And the row groups have to stay member-aligned, which is what _member_aligned_row_group_size and NWP_TARGET_FILE_SIZE_BYTES exist to guarantee.

NWP_TARGET_FILE_SIZE_BYTES = 2000000000 module-attribute

Target size for each Parquet file delta-rs writes, sized to keep one partition in one file.

The ensemble forecast of the European Centre for Medium-Range Weather Forecasts is abbreviated ECMWF ENS. One daily partition of ECMWF ENS is ~158 MB, so the target leaves more than a tenfold headroom.

The file-size target is an optimisation, not a correctness requirement. A partition that outgrows the target still writes correctly and still prunes well, because the single Arrow chunk and the member-aligned row-group size below do the real work. Measured on a real partition, a single-member read touches 1.96% of rows when the partition lands in one file and 3.92% when delta-rs splits it in two.

NWP_ROW_GROUP_SIZE_LIMITS = (1024, 1048576) module-attribute

Floor and ceiling clamped around the member-aligned row-group size.

The floor stops a frame with very few rows per member producing row groups too small to compress or too numerous to track. The floor has not bound on any real ECMWF ENS run measured so far, where one member occupies ~142,000 rows. The ceiling is Parquet's conventional maximum, which also bounds the streaming engine's peak memory per morsel.

Classes

Functions:

write_nwp(nwp, table_uri, storage_options=None)

Write one NWP run into the nwp Delta table in its storage format.

Rounds every continuous weather variable to NWP_SIGNIFICAND_BITS significand bits, sorts rows by NWP_SORT_COLS, and writes one row group per ensemble member (see _member_aligned_row_group_size) into a single file per partition. The table is partitioned by (nwp_model_id, init_time), matching Nwp.scan_delta's partition-pruning assumptions; the first write creates the table.

The write replaces the (nwp_model_id, init_time) partition named by the frame's first row, so re-materialising an ecmwf_ens partition leaves one copy of the run. delta-rs checks every row against that predicate. If any row falls outside the predicate, delta-rs rejects the whole write and leaves the table untouched. That was confirmed empirically against deltalake 1.6.3, locally and on S3, on a partition column despite its percent-encoded Hive directory name. Two materialisations of the same partition at once contend and the loser raises CommitFailedError; disjoint partitions do not.

schema_mode="overwrite" is passed on every write, not just to migrate an old table. Passing it on every write is safe. That safety claim holds only for a widening contract change, and a widening change is the only kind this function has been used for so far. Confirmed empirically against deltalake 1.6.3: widening updates only the table's logical schema — the _delta_log metadata — to match the incoming frame's dtypes. Widening leaves every other partition's physical Parquet bytes untouched. A later read at the new logical dtype is therefore correct and lossless, even for a partition still physically stored at an older, narrower dtype.

A validated pt.DataFrame[Nwp] input is what makes schema_mode="overwrite" safe to leave on. The nwp argument carries the full column set at the current contract's dtypes. The write can only ever change the table's schema through a deliberate widening of an Nwp dtype, and the write cannot silently drop a column.

A narrowing contract change is a different, worse failure mode, also confirmed empirically: the write that narrows a column succeeds silently, because schema_mode="overwrite" accepts that write at write time. The table is then left with a logical schema that its own, previously-written partitions can no longer safely satisfy. Nothing fails until a later read of the whole table, by anyone, not necessarily the writer that broke the table. That read raises SchemaError: incoming dtype cannot safely cast to target dtype. The error is confusing and far removed from its cause, especially if further partitions land on the bad schema before anyone reads across all of them.

Leaving schema_mode="overwrite" on permanently is an accepted, disclosed risk rather than a defect. Nothing in this contract's own dtype-widening history has ever narrowed a column, and there is no guard against a future change that does. Without schema_mode="overwrite", delta-rs instead auto-widens the incoming (narrower) data up to the table's existing (wider) schema at write time, and the table stays readable throughout.

Parameters:

Name Type Description Default
nwp DataFrame[Nwp]

Validated, non-empty NWP rows for a single (nwp_model_id, init_time) partition.

required
table_uri str | Path

Path or URI of the nwp Delta table.

required
storage_options ObjectStoreOptions | None

delta-rs object-store options (credentials/endpoint) for a remote table_uri; None/empty for a local path.

None
Source code in packages/delta_store/src/delta_store/nwp.py
142
143
144
145
146
147
148
149
150
151
152
153
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
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
def write_nwp(
    nwp: pt.DataFrame[Nwp],
    table_uri: str | Path,
    storage_options: ObjectStoreOptions | None = None,
) -> None:
    """Write one NWP run into the ``nwp`` Delta table in its storage format.

    Rounds every continuous weather variable to ``NWP_SIGNIFICAND_BITS`` significand bits, sorts
    rows by ``NWP_SORT_COLS``, and writes one row group per ensemble member (see
    `_member_aligned_row_group_size`) into a single file per partition. The table is
    partitioned by ``(nwp_model_id, init_time)``, matching ``Nwp.scan_delta``'s
    partition-pruning assumptions; the first write creates the table.

    The write **replaces** the ``(nwp_model_id, init_time)`` partition named by the frame's first
    row, so re-materialising an ``ecmwf_ens`` partition leaves one copy of the run.
    delta-rs checks every row against that predicate. If any row falls outside the predicate,
    delta-rs rejects the whole write and leaves the table untouched. That was confirmed empirically
    against ``deltalake`` 1.6.3, locally and on S3, on a partition column despite its
    percent-encoded Hive directory name. Two materialisations of the *same* partition at once
    contend and the loser raises ``CommitFailedError``; disjoint partitions do not.

    ``schema_mode="overwrite"`` is passed on every write, not just to migrate an old table.
    Passing it on every write is safe. That safety claim holds only for a **widening** contract
    change, and a widening change is the only kind this function has been used for so far. Confirmed
    empirically against ``deltalake`` 1.6.3: widening updates only the table's *logical* schema —
    the ``_delta_log`` metadata — to match the incoming frame's dtypes. Widening leaves every other
    partition's physical Parquet bytes untouched. A later read at the new logical dtype is therefore
    correct and lossless, even for a partition still physically stored at an older, narrower dtype.

    **A validated ``pt.DataFrame[Nwp]`` input is what makes ``schema_mode="overwrite"`` safe to
    leave on.** The ``nwp`` argument carries the full column set at the *current* contract's dtypes.
    The write can only ever change the table's schema through a deliberate widening of an ``Nwp``
    dtype, and the write cannot silently drop a column.

    A **narrowing** contract change is a different, worse failure mode, also confirmed
    empirically: the write that narrows a column succeeds silently, because
    ``schema_mode="overwrite"`` accepts that write at write time. The table is then left with a
    logical schema that its own, previously-written partitions can no longer safely satisfy. Nothing
    fails until a *later* read of the whole table, by anyone, not necessarily the writer that broke
    the table. That read raises ``SchemaError: incoming dtype cannot safely cast to target dtype``.
    The error is confusing and far removed from its cause, especially if further partitions land on
    the bad schema before anyone reads across all of them.

    **Leaving ``schema_mode="overwrite"`` on permanently is an accepted, disclosed risk rather
    than a defect.** Nothing in this contract's own dtype-widening history has ever narrowed a
    column, and there is no guard against a future change that does. Without
    ``schema_mode="overwrite"``, delta-rs instead auto-widens the incoming (narrower) data up to
    the table's existing (wider) schema at write time, and the table stays readable throughout.

    Args:
        nwp: Validated, non-empty NWP rows for a single ``(nwp_model_id, init_time)`` partition.
        table_uri: Path or URI of the ``nwp`` Delta table.
        storage_options: delta-rs object-store options (credentials/endpoint) for a remote
            ``table_uri``; ``None``/empty for a local path.
    """
    continuous_vars = sorted(Nwp.continuous_var_names())
    rounded = nwp.with_columns(
        **{
            var: round_to_significand_bits(pl.col(var), keep_bits=NWP_SIGNIFICAND_BITS)
            for var in continuous_vars
        }
    ).sort(*NWP_SORT_COLS)

    # Combine the frame into one Arrow chunk rather than the 32 chunks a measured partition arrived
    # in. delta-rs consumes a multi-chunk table out of order when partitioning the write. That
    # out-of-order read scatters the member-sorted rows across row groups and widens every row
    # group's ensemble_member min/max range. Measured on a real partition, a single-member read went
    # from 1.96% of rows to 33% purely from the chunking. Combining allocates one extra copy of the
    # frame.
    prepared = rounded.to_arrow().combine_chunks()

    write_deltalake(
        table_or_uri=table_uri,
        data=prepared,
        mode="overwrite",
        schema_mode="overwrite",
        predicate=(
            f"nwp_model_id = '{nwp.item(0, 'nwp_model_id')}' "
            f"AND init_time = '{nwp.item(0, 'init_time').isoformat()}'"
        ),
        partition_by=["nwp_model_id", "init_time"],
        writer_properties=_writer_properties(
            max_row_group_size=_member_aligned_row_group_size(nwp)
        ),
        target_file_size=NWP_TARGET_FILE_SIZE_BYTES,
        storage_options=typeddict_to_dict(storage_options),
    )

delta_store.power_time_series

Storage policy for the power_time_series Delta table.

See this package's __init__ docstring for which tables carry writer-properties tuning.

Classes

Functions:

write_power_time_series(power_ts, table_uri, *, storage_options=None)

Append PowerTimeSeries rows to the power_time_series Delta table.

The table is partitioned by time_series_id; the first write creates it. Appends only — the caller is responsible for de-duplicating rows already on disk before calling this (see power_time_series_and_metadata's use of select_new_rows).

Parameters:

Name Type Description Default
power_ts DataFrame[PowerTimeSeries]

Validated, de-duplicated new rows to append.

required
table_uri str | Path

Path or URI of the power_time_series Delta table.

required
storage_options ObjectStoreOptions | None

delta-rs object-store options (credentials/endpoint) for a remote table_uri; None/empty for a local path.

None
Source code in packages/delta_store/src/delta_store/power_time_series.py
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
def write_power_time_series(
    power_ts: pt.DataFrame[PowerTimeSeries],
    table_uri: str | Path,
    *,
    storage_options: ObjectStoreOptions | None = None,
) -> None:
    """Append ``PowerTimeSeries`` rows to the ``power_time_series`` Delta table.

    The table is partitioned by ``time_series_id``; the first write creates it. Appends only —
    the caller is responsible for de-duplicating rows already on disk before calling this (see
    ``power_time_series_and_metadata``'s use of ``select_new_rows``).

    Args:
        power_ts: Validated, de-duplicated new rows to append.
        table_uri: Path or URI of the ``power_time_series`` Delta table.
        storage_options: delta-rs object-store options (credentials/endpoint) for a remote
            ``table_uri``; ``None``/empty for a local path.
    """
    write_deltalake(
        table_or_uri=table_uri,
        data=power_ts.to_arrow(),
        mode="append",
        partition_by=["time_series_id"],
        storage_options=typeddict_to_dict(storage_options),
    )

delta_store.eligible_time_series

Storage policy for the eligible_time_series Delta table.

See this package's __init__ docstring for which tables carry writer-properties tuning.

Classes

Functions:

write_eligible_time_series(eligible, table_uri, *, fold_id, storage_options=None)

Write one fold's EligibleTimeSeries rows to the eligible_time_series Delta table.

The table is partitioned by fold_id; the write replaces that partition, so re-materialising a fold leaves one copy of its eligible population. fold_id is an explicit parameter rather than read off eligible. A fold with zero eligible series must still clear its own partition, and an empty frame has no row to read a fold id from.

Parameters:

Name Type Description Default
eligible DataFrame[EligibleTimeSeries]

Validated eligible-series rows for the fold named by fold_id. May be empty.

required
table_uri str | Path

Path or URI of the eligible_time_series Delta table.

required
fold_id str

The fold this write replaces, used in the overwrite predicate.

required
storage_options ObjectStoreOptions | None

delta-rs object-store options (credentials/endpoint) for a remote table_uri; None/empty for a local path.

None
Source code in packages/delta_store/src/delta_store/eligible_time_series.py
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
def write_eligible_time_series(
    eligible: pt.DataFrame[EligibleTimeSeries],
    table_uri: str | Path,
    *,
    fold_id: str,
    storage_options: ObjectStoreOptions | None = None,
) -> None:
    """Write one fold's ``EligibleTimeSeries`` rows to the ``eligible_time_series`` Delta table.

    The table is partitioned by ``fold_id``; the write **replaces** that partition, so
    re-materialising a fold leaves one copy of its eligible population. ``fold_id`` is an explicit
    parameter rather than read off ``eligible``. A fold with zero eligible series must still clear
    its own partition, and an empty frame has no row to read a fold id from.

    Args:
        eligible: Validated eligible-series rows for the fold named by ``fold_id``. May be empty.
        table_uri: Path or URI of the ``eligible_time_series`` Delta table.
        fold_id: The fold this write replaces, used in the overwrite predicate.
        storage_options: delta-rs object-store options (credentials/endpoint) for a remote
            ``table_uri``; ``None``/empty for a local path.
    """
    write_deltalake(
        table_or_uri=table_uri,
        data=eligible.to_arrow(),
        mode="overwrite",
        predicate=f"fold_id = '{fold_id}'",
        partition_by=["fold_id"],
        storage_options=typeddict_to_dict(storage_options),
    )

delta_store.effective_capacity

Storage policy for the effective_capacity Delta table.

See this package's __init__ docstring for which tables carry writer-properties tuning.

Classes

Functions:

write_effective_capacity(capacity, table_uri, *, storage_options=None)

Overwrite the effective_capacity Delta table with capacity.

The table is small — one row per series. The whole table is replaced on every write. There is no partitioning or predicate to scope the overwrite.

Parameters:

Name Type Description Default
capacity DataFrame[EffectiveCapacity]

Validated capacity rows, one per time_series_id.

required
table_uri str | Path

Path or URI of the effective_capacity Delta table.

required
storage_options ObjectStoreOptions | None

delta-rs object-store options (credentials/endpoint) for a remote table_uri; None/empty for a local path.

None
Source code in packages/delta_store/src/delta_store/effective_capacity.py
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
def write_effective_capacity(
    capacity: pt.DataFrame[EffectiveCapacity],
    table_uri: str | Path,
    *,
    storage_options: ObjectStoreOptions | None = None,
) -> None:
    """Overwrite the ``effective_capacity`` Delta table with ``capacity``.

    The table is small — one row per series. The whole table is replaced on every write. There is no
    partitioning or predicate to scope the overwrite.

    Args:
        capacity: Validated capacity rows, one per ``time_series_id``.
        table_uri: Path or URI of the ``effective_capacity`` Delta table.
        storage_options: delta-rs object-store options (credentials/endpoint) for a remote
            ``table_uri``; ``None``/empty for a local path.
    """
    write_deltalake(
        table_or_uri=table_uri,
        data=capacity.to_arrow(),
        mode="overwrite",
        storage_options=typeddict_to_dict(storage_options),
    )

delta_store.forecast_metrics

Storage policy for the forecast_metrics Delta table.

See this package's __init__ docstring for which tables carry writer-properties tuning. The Enum→String cast below is on-disk-format knowledge, so it lives here rather than in evaluation logic.

Classes

Functions:

write_forecast_metrics(metrics, table_uri, *, experiment_name, fold_id, storage_options=None)

Write Metrics rows to the forecast_metrics Delta table.

Casts the Enum columns (e.g. metric_name, horizon_slice) to String before writing. Without the cast, delta-rs writes the column as a dictionary-typed parquet column but records it as Utf8 in the Delta log. The write itself succeeds. A later pl.read_delta then raises SchemaError: data type mismatch for column ...: incoming: Enum([...]) != target: String (verified against deltalake 1.6.3 / polars 1.44.2). Casting before the write keeps the file's physical type and the log's logical type in step. Performs an idempotent overwrite of the (experiment_name, fold_id) partition so re-materialising the asset replaces rows rather than duplicating them.

Parameters:

Name Type Description Default
metrics DataFrame[Metrics]

Fully populated Metrics rows, with all provenance columns set by enrich_metrics_rows().

required
table_uri str | Path

Path or URI of the forecast_metrics Delta table.

required
experiment_name str

Experiment name; used in the Delta overwrite predicate to scope the replacement to this (experiment_name, fold_id) partition.

required
fold_id str

Fold identifier; used alongside experiment_name in the predicate.

required
storage_options ObjectStoreOptions | None

delta-rs object-store options (credentials/endpoint) for a remote table_uri; None/empty for a local path.

None
Source code in packages/delta_store/src/delta_store/forecast_metrics.py
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
def write_forecast_metrics(
    metrics: pt.DataFrame[Metrics],
    table_uri: str | Path,
    *,
    experiment_name: str,
    fold_id: str,
    storage_options: ObjectStoreOptions | None = None,
) -> None:
    """Write ``Metrics`` rows to the ``forecast_metrics`` Delta table.

    Casts the ``Enum`` columns (e.g. ``metric_name``, ``horizon_slice``) to ``String`` before
    writing. Without the cast, delta-rs writes the column as a dictionary-typed parquet column but
    records it as ``Utf8`` in the Delta log. The write itself succeeds. A later ``pl.read_delta``
    then raises ``SchemaError: data type mismatch for column ...: incoming: Enum([...]) != target:
    String`` (verified against deltalake 1.6.3 / polars 1.44.2). Casting before the write keeps the
    file's physical type and the log's logical type in step. Performs an idempotent overwrite of the
    ``(experiment_name, fold_id)`` partition so re-materialising the asset replaces rows rather than
    duplicating them.

    Args:
        metrics: Fully populated ``Metrics`` rows, with all provenance columns set by
            ``enrich_metrics_rows()``.
        table_uri: Path or URI of the ``forecast_metrics`` Delta table.
        experiment_name: Experiment name; used in the Delta overwrite predicate to scope the
            replacement to this ``(experiment_name, fold_id)`` partition.
        fold_id: Fold identifier; used alongside ``experiment_name`` in the predicate.
        storage_options: delta-rs object-store options (credentials/endpoint) for a remote
            ``table_uri``; ``None``/empty for a local path.
    """
    enum_cols = [c for c, dtype in metrics.schema.items() if isinstance(dtype, pl.Enum)]
    delta_data = metrics.with_columns(pl.col(c).cast(pl.String) for c in enum_cols).to_arrow()
    write_deltalake(
        table_or_uri=table_uri,
        data=delta_data,
        mode="overwrite",
        predicate=f"experiment_name = '{experiment_name}' AND fold_id = '{fold_id}'",
        partition_by=["experiment_name", "fold_id"],
        storage_options=typeddict_to_dict(storage_options),
    )