Skip to content

NGED Data API

NGED JSON Data

This package reads NGED's telemetry JSON files from S3 and parses them into the PowerTimeSeries and TimeSeriesMetadata schemas (see contracts). The metadata roster is the only file this package owns and writes; the parsed power observations are handed back to the caller, which appends them to the power_time_series Delta table (see Usage below). nged_data.storage's module docstring, below on this page, says which functions read that Delta table and which write the roster.

Public surface

Only upsert_metadata is re-exported from the package root (from nged_data import upsert_metadata); the other five live in nged_data.storage (from nged_data.storage import list_timeseries_json_files, etc.).

  • nged_data.storage.list_timeseries_json_files(store) — lists the timeseries JSON files on NGED's S3 bucket, parsing time_series_id, start_time, and end_time out of each file's path.
  • nged_data.storage.remove_small_files_from_listing(file_listing, size_threshold_bytes=520) — drops files too small to carry any readings, so download_and_parse_files never fetches and parses one only to discard the result.
  • nged_data.storage.download_and_parse_files(store, paths_df) — downloads and parses each listed file, returning a DownloadAndParseResult of metadata (TimeSeriesMetadata), power_time_series (PowerTimeSeries), and n_implausible_power_rows_dropped. Raises NoNewData if the listing was empty, or if every listed file's data field was null.
  • nged_data.storage.select_new_rows(time_series, delta_path, storage_options=None) — filters time_series down to rows genuinely missing from the power_time_series Delta table at delta_path. PowerTimeSeries rows are filtered by existence, an anti-join on (time_series_id, time), so a reading is kept regardless of arrival order. The file listing that list_timeseries_json_files returns is filtered by comparing each file's end_time against its series' on-disk high-water mark, loosened by a lookback margin so a late file is still downloaded.
  • nged_data.storage.time_series_coverage(delta_path, storage_options=None) — the earliest and latest observation time on disk for each time_series_id in the power_time_series Delta table.
  • nged_data.upsert_metadata(new_metadata, metadata_path, storage_options=None) — merges a TimeSeriesMetadata snapshot into the stored metadata Parquet file, keeping the newest values per time_series_id and rewriting the file only if the incoming metadata differs from what is stored.

nged_data.read_nged_json parses one downloaded JSON file into the two schemas. All three of its functions are private, and the two that parse a whole file are called only by download_and_parse_files, so the module appears on the API page below carrying just its ExtractedPowerTimeSeries result type.

Data quality

download_and_parse_files drops rows whose time is malformed — outside the plausible datetime range, null, or not aligned to the top or bottom of the hour — via PowerTimeSeries.drop_implausible_rows, and reports how many as n_implausible_power_rows_dropped. Degrading rather than raising on a malformed reading follows inherent stability: a malformed time originates upstream of our pipeline, at the meter or in the telemetry export, not in our own code. Ingestion therefore keeps the rest of the batch rather than aborting it. No other cleaning happens during ingestion.

Usage

This package is used by the power_time_series_and_metadata Dagster asset in src/nged_substation_forecast/defs/assets.py.

nged_data.read_nged_json

Extracts metadata and time series from NGED JSON data.

The parser expects two properties of each JSON file. Metadata fields sit at the top level. A data field holds an array of time series data points.

Attributes

log = logging.getLogger(__name__) module-attribute

Classes

ExtractedPowerTimeSeries

Bases: NamedTuple

Result of parsing PowerTimeSeries rows out of one NGED JSON file.

n_dropped counts rows dropped by PowerTimeSeries.drop_implausible_rows for a malformed time — see that method's docstring for why ingestion degrades rather than raising.

Source code in packages/nged_data/src/nged_data/read_nged_json.py
28
29
30
31
32
33
34
35
36
37
class ExtractedPowerTimeSeries(NamedTuple):
    """Result of parsing ``PowerTimeSeries`` rows out of one NGED JSON file.

    ``n_dropped`` counts rows dropped by ``PowerTimeSeries.drop_implausible_rows`` for a
    malformed ``time`` — see that method's docstring for why ingestion degrades rather than
    raising.
    """

    dataframe: pt.DataFrame[PowerTimeSeries]
    n_dropped: int
Attributes
dataframe instance-attribute
n_dropped instance-attribute

nged_data.storage

Listing, downloading, and parsing NGED's telemetry JSON from S3.

Each file yields TimeSeriesMetadata describing one series and PowerTimeSeries power observations from it. The metadata is upserted here, into a Parquet roster (see upsert_metadata). The power observations are not. This module never writes PowerTimeSeries rows to disk. It returns them to the caller, and the caller appends them to the power_time_series Delta table. time_series_coverage and select_new_rows read that same Delta table, to find what is new, but neither writes to it.

Attributes

log = logging.getLogger(__name__) module-attribute

Classes

NoNewData

Bases: Exception

Raised by download_and_parse_files when none of its listed files carried power data.

Raised when the file listing was empty, or when every listed file's data field was null. A file whose data field is present contributes a DataFrame, even when every row in that file is dropped as implausible. A file that yields zero usable rows therefore does not raise.

Source code in packages/nged_data/src/nged_data/storage.py
170
171
172
173
174
175
176
class NoNewData(Exception):
    """Raised by `download_and_parse_files` when none of its listed files carried power data.

    Raised when the file listing was empty, or when every listed file's ``data`` field was null.
    A file whose ``data`` field is present contributes a DataFrame, even when every row in that
    file is dropped as implausible. A file that yields zero usable rows therefore does not raise.
    """

DownloadAndParseResult

Bases: NamedTuple

Result of download_and_parse_files.

n_implausible_power_rows_dropped sums ExtractedPowerTimeSeries.n_dropped across every file in the batch — see PowerTimeSeries.drop_implausible_rows for what gets dropped and why.

Source code in packages/nged_data/src/nged_data/storage.py
179
180
181
182
183
184
185
186
187
188
189
class DownloadAndParseResult(NamedTuple):
    """Result of ``download_and_parse_files``.

    ``n_implausible_power_rows_dropped`` sums ``ExtractedPowerTimeSeries.n_dropped`` across every
    file in the batch — see ``PowerTimeSeries.drop_implausible_rows`` for what gets dropped and
    why.
    """

    metadata: pt.DataFrame[TimeSeriesMetadata]
    power_time_series: pt.DataFrame[PowerTimeSeries]
    n_implausible_power_rows_dropped: int
Attributes
metadata instance-attribute
power_time_series instance-attribute
n_implausible_power_rows_dropped instance-attribute

TimeSeriesCoverage

Bases: Model

Per-series observation-time span of the power_time_series Delta table.

first_time/last_time are the earliest/latest observation time for each time_series_id. This frame is a transient intermediate, never persisted. Three callers read the frame: the freshness asset check reads last_time to detect staleness, select_new_rows reads last_time to find genuinely-new rows, and cross-validation (CV) fold-eligibility (eligible_time_series_ids) reads both first_time and last_time. A CV fold is one train/test split of the history, and a series is eligible for a fold only if its data covers that split.

The freshness check reads this on-disk recency rather than the asset's materialisation timestamp. A materialisation-freshness policy would miss the failure the check exists to catch. When NGED's telemetry stalls, the ingest asset keeps materialising successfully on schedule, and writes nothing. The materialisation looks fresh; the newest observation on disk does not. Full reasoning: https://openclimatefix.github.io/nged-substation-forecast/architecture/production-deployment/#warn-on-stale-power-data-with-a-dagster-asset-check.

Source code in packages/nged_data/src/nged_data/storage.py
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
class TimeSeriesCoverage(pt.Model):
    """Per-series observation-time span of the ``power_time_series`` Delta table.

    ``first_time``/``last_time`` are the earliest/latest observation ``time`` for each
    ``time_series_id``. This frame is a transient intermediate, never persisted. Three callers
    read the frame: the freshness asset check reads ``last_time`` to detect staleness,
    ``select_new_rows`` reads ``last_time`` to find genuinely-new rows, and cross-validation (CV)
    fold-eligibility (``eligible_time_series_ids``) reads both ``first_time`` and ``last_time``.
    A CV fold is one train/test split of the history, and a series is eligible for a fold only if
    its data covers that split.

    The freshness check reads this on-disk recency rather than the asset's materialisation
    timestamp. A materialisation-freshness policy would miss the failure the check exists to
    catch. When NGED's telemetry stalls, the ingest asset keeps materialising successfully on
    schedule, and writes nothing. The materialisation looks fresh; the newest observation on disk
    does not. Full reasoning:
    <https://openclimatefix.github.io/nged-substation-forecast/architecture/production-deployment/#warn-on-stale-power-data-with-a-dagster-asset-check>.
    """

    time_series_id: int = _get_time_series_id_dtype(unique=True)
    first_time: int = pt.Field(dtype=PowerTimeSeries.dtypes["time"])
    last_time: int = pt.Field(dtype=PowerTimeSeries.dtypes["time"])
Attributes
time_series_id = _get_time_series_id_dtype(unique=True) class-attribute instance-attribute
first_time = pt.Field(dtype=PowerTimeSeries.dtypes['time']) class-attribute instance-attribute
last_time = pt.Field(dtype=PowerTimeSeries.dtypes['time']) class-attribute instance-attribute

UpsertMetadataStats

Bases: TypedDict

What the TimeSeriesMetadata upsert did, published as Dagster output metadata.

Source code in packages/nged_data/src/nged_data/storage.py
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
class UpsertMetadataStats(TypedDict, total=False):
    """What the ``TimeSeriesMetadata`` upsert did, published as Dagster output metadata."""

    metadata_n_new_TimeSeriesIDs: int
    metadata_n_updated_TimeSeriesIDs: int
    metadata_updated_TimeSeriesIDs: Sequence[int]
    metadata_upsert_failed: str
    """Set by the asset when the whole upsert raised, so the power write went ahead without it.

    Read this field's presence as "the roster is stale, retry next hour", not as a power-ingest
    failure — the power write is unaffected. See [Degraded input
    data](https://openclimatefix.github.io/nged-substation-forecast/live_service/operations/#degraded-input-data-nwp-feed-down-or-telemetry-stalled),
    under "Reading a failed roster upsert", for the operational read of this field and what a stale
    roster costs while the failure persists.
    """
Attributes
metadata_n_new_TimeSeriesIDs instance-attribute
metadata_n_updated_TimeSeriesIDs instance-attribute
metadata_updated_TimeSeriesIDs instance-attribute
metadata_upsert_failed instance-attribute

Set by the asset when the whole upsert raised, so the power write went ahead without it.

Read this field's presence as "the roster is stale, retry next hour", not as a power-ingest failure — the power write is unaffected. See Degraded input data, under "Reading a failed roster upsert", for the operational read of this field and what a stale roster costs while the failure persists.

Functions:

list_timeseries_json_files(store)

List all the timeseries JSON files in NGED's S3 bucket.

Each path encodes start_time, end_time, and time_series_id, and is assumed to be of the form:

timeseries/1774512000000_1774533600000/TimeSeries_23_20260326T080000Z_20260326T140000Z.json

_process_file_listing parses those three fields back out, and its aligned comment traces the regex against this same key.

A key that does not match yields null captures and fails _ProcessedFileListing.validate. That validation failure aborts the listing for the whole bucket rather than skipping the one object. Aborting the whole listing is deliberately stricter than the null-data handling in download_and_parse_files, which logs the offending file and carries on. The difference is where the fault lies. A malformed reading originates upstream of our pipeline, at the meter or in the telemetry export. A key we cannot parse means NGED's naming convention has changed. Every start_time, end_time, and time_series_id this function returns is then suspect — including the values parsed from the keys that still match.

Source code in packages/nged_data/src/nged_data/storage.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
def list_timeseries_json_files(
    store: obstore.store.S3Store,
) -> pt.DataFrame[_ProcessedFileListing]:
    """List all the timeseries JSON files in NGED's S3 bucket.

    Each path encodes `start_time`, `end_time`, and `time_series_id`, and is assumed to be of the
    form:

        timeseries/1774512000000_1774533600000/TimeSeries_23_20260326T080000Z_20260326T140000Z.json

    `_process_file_listing` parses those three fields back out, and its aligned comment traces the
    regex against this same key.

    A key that does not match yields null captures and fails `_ProcessedFileListing.validate`. That
    validation failure aborts the listing for the whole bucket rather than skipping the one object.
    Aborting the whole listing is deliberately stricter than the null-`data` handling in
    `download_and_parse_files`, which logs the offending file and carries on. The difference is
    where the fault lies. A malformed reading originates upstream of our pipeline, at the meter or
    in the telemetry export. A key we cannot parse means NGED's naming convention has changed. Every
    `start_time`, `end_time`, and `time_series_id` this function returns is then suspect — including
    the values parsed from the keys that still match.
    """
    raw_file_listing: list[_RawFileListItem] = []
    total_objects = 0
    for chunk in store.list(prefix="timeseries"):
        # `list()` returns the file listing in chunks of `chunk_size=50` items per chunk.
        for object_meta in chunk:
            total_objects += 1
            if object_meta["path"].endswith(".json"):
                raw_file_listing.append(
                    _RawFileListItem(
                        path=object_meta["path"],
                        filesize_bytes=object_meta["size"],
                    ),
                )
    log.info(f"JSON files on NGED's S3: {len(raw_file_listing)} out of {total_objects=}")
    return _process_file_listing(raw_file_listing)

remove_small_files_from_listing(file_listing, size_threshold_bytes=520)

Remove files too small to carry any readings.

The filter skips NGED JSON files that have no data field, so download_and_parse_files never has to fetch and parse them only to discard the result. It is an optimisation, not a correctness requirement: download_and_parse_files already tolerates a null data field.

size_threshold_bytes defaults to 520, derived from the real files on NGED's S3. Some NGED files carry a Well-Known Text (WKT) geometry string, which dominates their size; the rest are WKT-less. A WKT-less file with zero readings tops out at 488 bytes. A WKT-less file with a single reading starts at 556 bytes. The 520-byte default sits in the gap between 488 and 556. WKT-bearing (Primary substation) files run far larger: zero-reading examples measured between 4,405 and 20,148 bytes. The WKT-less floor is therefore the binding constraint on the threshold.

That 68-byte gap comes from V1's 32 series (phased rollout), The gap is narrow, so V2's ~2,500 series want re-measuring before this default is trusted there. Two changes would close the gap. A populated information field would push a zero-reading file above 520 bytes; TimeSeriesMetadata records that field as always null in the V1 trial area. A substation name shorter than any in V1 would pull a one-reading file below 520 bytes. Re-run the measurement rather than assume the gap survives.

Source code in packages/nged_data/src/nged_data/storage.py
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
def remove_small_files_from_listing(
    file_listing: pt.DataFrame[_ProcessedFileListing],
    size_threshold_bytes: int = 520,
) -> pt.DataFrame[_ProcessedFileListing]:
    """Remove files too small to carry any readings.

    The filter skips NGED JSON files that have no `data` field, so `download_and_parse_files`
    never has to fetch and parse them only to discard the result. It is an optimisation, not a
    correctness requirement: `download_and_parse_files` already tolerates a null `data` field.

    `size_threshold_bytes` defaults to 520, derived from the real files on NGED's S3. Some NGED
    files carry a Well-Known Text (WKT) geometry string, which dominates their size; the rest are
    WKT-less. A WKT-less file with zero readings tops out at 488 bytes. A WKT-less file with a
    single reading starts at 556 bytes. The 520-byte default sits in the gap between 488 and 556.
    WKT-bearing (Primary substation) files run far larger: zero-reading examples measured between
    4,405 and 20,148 bytes. The WKT-less floor is therefore the binding constraint on the
    threshold.

    That 68-byte gap comes from V1's 32 series ([phased
    rollout](https://openclimatefix.github.io/nged-substation-forecast/background/requirements/#phased-rollout)),
    The gap is narrow, so V2's ~2,500 series want re-measuring before this default is trusted there.
    Two changes would close the gap. A populated `information` field would push a zero-reading file
    above 520 bytes; `TimeSeriesMetadata` records that field as always null in the V1 trial area. A
    substation name shorter than any in V1 would pull a one-reading file below 520 bytes. Re-run the
    measurement rather than assume the gap survives.
    """
    n_files_before_filter = file_listing.height
    filtered = file_listing.filter(pl.col("filesize_bytes") > size_threshold_bytes)
    log.info(
        f"Files retained after the size filter: {filtered.height} out of {n_files_before_filter=}"
    )
    return filtered

download_and_parse_files(store, paths_df)

Download and parse each listed file, one end_time group at a time.

The listing is taken in ascending end_time order. Two files can cover overlapping periods for the same time_series_id. Processing in end_time order means the more recent file's readings overwrite the older file's duplicate rows, in the unique(..., keep="last") dedupe below.

Parameters:

Name Type Description Default
store S3Store

The NGED S3 bucket to download each file from.

required
paths_df DataFrame[_ProcessedFileListing]

The file listing to process — typically already filtered by remove_small_files_from_listing and select_new_rows.

required

Returns:

Type Description
DownloadAndParseResult

A DownloadAndParseResult bundling every file's parsed TimeSeriesMetadata and

DownloadAndParseResult

PowerTimeSeries rows, deduplicated across files — see DownloadAndParseResult's

DownloadAndParseResult

docstring for what each field holds.

Raises:

Type Description
NoNewData

if paths_df was empty, or if every listed file's data field was null. A null data field means NGED's meter reported nothing for the period that file covers. The guard counts the DataFrames collected, not the rows in them. A file whose data field is present therefore still counts, even when every one of its rows is dropped as implausible, and does not raise. The power_time_series_and_metadata asset catches NoNewData and reports an empty ingest, so an empty listing degrades the run rather than failing it.

Source code in packages/nged_data/src/nged_data/storage.py
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
243
244
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
def download_and_parse_files(
    store: obstore.store.S3Store, paths_df: pt.DataFrame[_ProcessedFileListing]
) -> DownloadAndParseResult:
    """Download and parse each listed file, one `end_time` group at a time.

    The listing is taken in ascending `end_time` order. Two files can cover overlapping periods for
    the same `time_series_id`. Processing in `end_time` order means the more recent file's readings
    overwrite the older file's duplicate rows, in the `unique(..., keep="last")` dedupe below.

    Args:
        store: The NGED S3 bucket to download each file from.
        paths_df: The file listing to process — typically already filtered by
            `remove_small_files_from_listing` and `select_new_rows`.

    Returns:
        A `DownloadAndParseResult` bundling every file's parsed `TimeSeriesMetadata` and
        `PowerTimeSeries` rows, deduplicated across files — see `DownloadAndParseResult`'s
        docstring for what each field holds.

    Raises:
        NoNewData: if `paths_df` was empty, or if every listed file's `data` field was null. A null
            `data` field means NGED's meter reported nothing for the period that file covers. The
            guard counts the DataFrames collected, not the rows in them. A file whose `data` field
            is present therefore still counts, even when every one of its rows is dropped as
            implausible, and does not raise. The `power_time_series_and_metadata` asset catches
            `NoNewData` and reports an empty ingest, so an empty listing degrades the run rather
            than failing it.
    """
    metadata_dfs = []
    power_time_series_dfs = []
    n_implausible_power_rows_dropped = 0
    for _end_time, df_for_end_time in paths_df.group_by("end_time", maintain_order=True):
        for path in df_for_end_time["path"]:
            # TODO: Use `store.get_async` to get all files for this group concurrently.
            result = store.get(path)
            json_bytes = bytes(result.bytes())
            df = pl.read_json(json_bytes)

            # Extract TimeSeriesMetadata from df:
            new_metadata_df = _extract_time_series_metadata(df)
            metadata_dfs.append(new_metadata_df)
            time_series_id: int = new_metadata_df["time_series_id"].item()

            # Extract PowerTimeSeries from df:
            try:
                extracted = _extract_power_time_series(df=df, time_series_id=time_series_id)
            except pl.exceptions.InvalidOperationError as e:
                if "invalid dtype: expected 'Struct', got 'Null' for 'data'" in str(e):
                    log.warning(
                        f"The 'data' field is 'null' in {path=}. This is expected behaviour if"
                        " NGED's meter reported no values for the period covered by the JSON file."
                    )
                else:
                    raise
            else:
                power_time_series_dfs.append(extracted.dataframe)
                n_implausible_power_rows_dropped += extracted.n_dropped

    log.info(
        f"{len(metadata_dfs)} new TimeSeriesMetadata DataFrames and {len(power_time_series_dfs)}"
        " new PowerTimeSeries dataframes extracted from NGED JSON data."
    )

    if len(metadata_dfs) == 0 or len(power_time_series_dfs) == 0:
        raise NoNewData

    # Concatenate and return:
    metadata_df = (
        pl.concat(metadata_dfs, how="diagonal")
        .unique(subset="time_series_id", keep="last")
        .sort("time_series_id")
    )
    time_series_df = (
        pl.concat(power_time_series_dfs)
        .unique(subset=["time_series_id", "time"], keep="last")
        .sort(by=PowerTimeSeries.columns_to_sort_by)
    )

    return DownloadAndParseResult(
        metadata=TimeSeriesMetadata.validate(metadata_df),
        power_time_series=PowerTimeSeries.validate(time_series_df),
        n_implausible_power_rows_dropped=n_implausible_power_rows_dropped,
    )

time_series_coverage(delta_path, storage_options=None)

Return the earliest/latest observation time on disk per time_series_id.

Returns an empty (but correctly typed) frame if the Delta table does not exist yet. min/max grouped by time_series_id are value aggregations, so they are safe from the Polars 32-bit row-count wraparound even on a very large table (see https://openclimatefix.github.io/nged-substation-forecast/architecture/code-style/#data-handling).

Cost: a full two-column scan-and-aggregate, O(rows in the table). Projection pushdown drops the power column. A group-wise min/max cannot be answered from Parquet row-group statistics, because no engine on our stack does aggregate-from-statistics, so every time/time_series_id value is read. Computing both bounds instead of one takes ~20% more wall-clock time and no extra memory, because the shared scan dominates.

The collect uses the streaming engine to keep peak memory bounded, because this scan runs hourly on a small control-plane VM. The scan runs twice in each of those hours. The power_data_is_fresh asset check runs it once. power_time_series_and_metadata runs it again, inside the select_new_rows call that asset makes on the file listing. The second select_new_rows call, on the parsed rows, uses _existing_power_time_series_keys instead — a scan restricted to the reporting series' own history, not the whole-table scan this function runs. The measurement used a synthetic V2 table: 2,500 series, half-hourly, partitioned by time_series_id, holding a year of history (43.8M rows). The streaming engine took ~0.21 s at ~190 MB peak. The in-memory engine peaked at ~1.3 GB for the same result, so streaming uses ~7x less memory.

Cost scales linearly with accumulated history. If the scan ever becomes a problem, both bounds can instead be read from the Delta add-action min.time/max.time file statistics. Delta's transaction log records one add-action per data file, carrying that file's per-column minimum and maximum. That read is metadata-only, O(files): ~0.02 s and <100 MB at the same scale. The statistics read is the same Delta-log-metadata trick used to count whole-table rows without scanning.

delta_path is a local path or remote URI for the power_time_series Delta table; storage_options carries the object-store credentials/endpoint for a remote delta_path.

Source code in packages/nged_data/src/nged_data/storage.py
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
def time_series_coverage(
    delta_path: str,
    storage_options: ObjectStoreOptions | None = None,
) -> pt.DataFrame[TimeSeriesCoverage]:
    """Return the earliest/latest observation ``time`` on disk per ``time_series_id``.

    Returns an empty (but correctly typed) frame if the Delta table does not exist yet.
    ``min``/``max`` grouped by ``time_series_id`` are value aggregations, so they are safe from
    the Polars 32-bit row-count wraparound even on a very large table (see
    <https://openclimatefix.github.io/nged-substation-forecast/architecture/code-style/#data-handling>).

    Cost: a full two-column scan-and-aggregate, O(rows in the table). Projection pushdown drops
    the ``power`` column. A group-wise ``min``/``max`` cannot be answered from Parquet row-group
    statistics, because no engine on our stack does aggregate-from-statistics, so every
    ``time``/``time_series_id`` value is read. Computing both bounds instead of one takes ~20%
    more wall-clock time and no extra memory, because the shared scan dominates.

    The ``collect`` uses the streaming engine to keep peak memory bounded, because this scan runs
    hourly on a small control-plane VM. The scan runs twice in each of those hours. The
    ``power_data_is_fresh`` asset check runs it once. ``power_time_series_and_metadata`` runs it
    again, inside the ``select_new_rows`` call that asset makes on the file listing. The second
    ``select_new_rows`` call, on the parsed rows, uses ``_existing_power_time_series_keys``
    instead — a scan restricted to the reporting series' own history, not the whole-table scan
    this function runs. The measurement used a synthetic V2 table: 2,500 series, half-hourly,
    partitioned by ``time_series_id``, holding a year of history (43.8M rows). The streaming
    engine took ~0.21 s at ~190 MB peak. The in-memory engine peaked at ~1.3 GB for the same
    result, so streaming uses ~7x less memory.

    Cost scales linearly with accumulated history. If the scan ever becomes a problem, both
    bounds can instead be read from the Delta add-action ``min.time``/``max.time`` file
    statistics. Delta's transaction log records one add-action per data file, carrying that
    file's per-column minimum and maximum. That read is metadata-only, O(files): ~0.02 s and <100
    MB at the same scale. The statistics read is the same Delta-log-metadata trick used to count
    whole-table rows without scanning.

    `delta_path` is a local path or remote URI for the ``power_time_series`` Delta table;
    `storage_options` carries the object-store credentials/endpoint for a remote `delta_path`.
    """
    if not delta_table_exists(delta_path, storage_options):
        log.info(f"{delta_path=} does not exist yet; returning an empty coverage frame.")
        empty = pl.DataFrame(
            schema={name: TimeSeriesCoverage.dtypes[name] for name in TimeSeriesCoverage.columns}
        )
        return pt.DataFrame(empty).set_model(TimeSeriesCoverage).validate()

    coverage = (
        pl.scan_delta(delta_path, storage_options=typeddict_to_dict(storage_options))
        .group_by("time_series_id")
        .agg(first_time=pl.min("time"), last_time=pl.max("time"))
        # Streaming engine: bounds peak memory (~7x lower than in-memory at V2 scale) so the
        # hourly full-table aggregate stays comfortable on a small control-plane VM. See docstring.
        .collect(engine="streaming")
    )
    log.info(
        f"Found on-disk coverage for {coverage.height} time_series_ids from {delta_path}."
        f" {coverage['last_time'].min()=}. {coverage['last_time'].max()=}"
    )
    return pt.DataFrame(coverage).set_model(TimeSeriesCoverage).validate()

select_new_rows(time_series, delta_path, storage_options=None)

select_new_rows(
    time_series: pt.DataFrame[PowerTimeSeries],
    delta_path: str,
    storage_options: ObjectStoreOptions | None = None,
) -> pt.DataFrame[PowerTimeSeries]
select_new_rows(
    time_series: pt.DataFrame[_ProcessedFileListing],
    delta_path: str,
    storage_options: ObjectStoreOptions | None = None,
) -> pt.DataFrame[_ProcessedFileListing]

Return rows in time_series genuinely missing from the Delta table.

time_series is either PowerTimeSeries rows or the _ProcessedFileListing a raw S3 listing parses into. The function tells the two apart by which of time/end_time is present. The two @overload declarations above tell a type checker which input type produces which output type. delta_path is a local path or remote URI for the power_time_series Delta table; storage_options carries the object-store credentials/endpoint for a remote delta_path.

When the Delta table does not exist yet, a call returns its input unchanged and scans nothing.

For PowerTimeSeries rows, the filter is a genuine existence check: an anti-join on (time_series_id, time) against _existing_power_time_series_keys. A late file is therefore ingested even when a later reading for the same series is already on disk. So is a file that fills a gap earlier in a series' history. See _existing_power_time_series_keys's docstring for the cost the existence check trades in return.

For the file listing, there is no per-row time to check existence against before download — only the file's start_time/end_time window from its S3 key. The filter therefore compares end_time against each series' on-disk last_time, which it takes from time_series_coverage and which acts as that series' watermark. _LATE_FILE_LOOKBACK loosens that comparison, so a file whose end_time falls a short while before the watermark is still downloaded. A file whose end_time falls more than _LATE_FILE_LOOKBACK before last_time is still dropped before download.

Cost: the file-listing branch runs time_series_coverage, one full two-column scan of power_time_series — see that function for the measured figures. The PowerTimeSeries branch instead runs _existing_power_time_series_keys, restricted to the time_series_ids in time_series — see that function's docstring for its own measured figures. power_time_series_and_metadata calls select_new_rows once on the file listing. The asset calls select_new_rows again on the parsed rows, but only if that listing turned up files worth downloading. An hour in which NGED published nothing new stops after the first call, because download_and_parse_files raises NoNewData in between and the asset returns.

Source code in packages/nged_data/src/nged_data/storage.py
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
def select_new_rows(
    time_series: pt.DataFrame[PowerTimeSeries] | pt.DataFrame[_ProcessedFileListing],
    delta_path: str,
    storage_options: ObjectStoreOptions | None = None,
) -> pt.DataFrame[PowerTimeSeries] | pt.DataFrame[_ProcessedFileListing]:
    """Return rows in `time_series` genuinely missing from the Delta table.

    `time_series` is either `PowerTimeSeries` rows or the `_ProcessedFileListing` a raw S3
    listing parses into. The function tells the two apart by which of `time`/`end_time` is
    present. The two `@overload` declarations above tell a type checker which input type produces
    which output type. `delta_path` is a local path or remote URI for the ``power_time_series``
    Delta table; `storage_options` carries the object-store credentials/endpoint for a remote
    `delta_path`.

    When the Delta table does not exist yet, a call returns its input unchanged and scans nothing.

    For `PowerTimeSeries` rows, the filter is a genuine existence check: an anti-join on
    `(time_series_id, time)` against `_existing_power_time_series_keys`. A late file is therefore
    ingested even when a later reading for the same series is already on disk. So is a file that
    fills a gap earlier in a series' history. See `_existing_power_time_series_keys`'s docstring
    for the cost the existence check trades in return.

    For the file listing, there is no per-row `time` to check existence against before download —
    only the file's `start_time`/`end_time` window from its S3 key. The filter therefore compares
    `end_time` against each series' on-disk `last_time`, which it takes from
    `time_series_coverage` and which acts as that series' watermark. `_LATE_FILE_LOOKBACK`
    loosens that comparison, so a file whose `end_time` falls a short while before the watermark
    is still downloaded. A file whose `end_time` falls more than `_LATE_FILE_LOOKBACK` before
    `last_time` is still dropped before download.

    Cost: the file-listing branch runs `time_series_coverage`, one full two-column scan of
    `power_time_series` — see that function for the measured figures. The `PowerTimeSeries`
    branch instead runs `_existing_power_time_series_keys`, restricted to the `time_series_id`s
    in `time_series` — see that function's docstring for its own measured figures.
    `power_time_series_and_metadata` calls `select_new_rows` once on the file listing. The asset
    calls `select_new_rows` again on the parsed rows, but only if that listing turned up files
    worth downloading. An hour in which NGED published nothing new stops after the first call,
    because `download_and_parse_files` raises `NoNewData` in between and the asset returns.
    """
    if not delta_table_exists(delta_path, storage_options):
        log.info(f"{delta_path=} does not exist yet.")
        return time_series

    if "time" in time_series.columns:
        reporting_ids = time_series["time_series_id"].unique().to_list()
        existing_keys = _existing_power_time_series_keys(delta_path, storage_options, reporting_ids)
        filtered_df = (
            time_series.lazy()
            .join(existing_keys, on=["time_series_id", "time"], how="anti")
            .sort(by=PowerTimeSeries.columns_to_sort_by)
            .collect()
        )
        return pt.DataFrame(filtered_df).set_model(PowerTimeSeries).validate()
    if "end_time" in time_series.columns:
        # Strip the Patito model from `coverage` so Polars' cross-subclass join check accepts it,
        # and keep only `last_time` (the most recent time on disk per series) for the filter below.
        coverage = time_series_coverage(delta_path, storage_options)
        plain_last_times = pl.LazyFrame._from_pyldf(coverage.lazy()._ldf).select(
            "time_series_id", "last_time"
        )
        filtered_df = (
            time_series.lazy()
            .join(plain_last_times, on="time_series_id", how="left")
            # A null last_time means this is a new time_series_id, so keep it unconditionally.
            .filter(
                pl.col("last_time").is_null()
                | (pl.col("end_time") > pl.col("last_time") - _LATE_FILE_LOOKBACK)
            )
            .drop("last_time")
            .sort(by="end_time")
            .collect()
        )
        return pt.DataFrame(filtered_df).set_model(_ProcessedFileListing).validate()
    raise ValueError(
        "Expected `time_series` to have either a `time` column or an `end_time` column,"
        f" not {time_series.columns=}"
    )

upsert_metadata(new_metadata, metadata_path, storage_options=None)

Upserts metadata to a Parquet file, keeping the newest version of each time series.

If the Parquet file does not exist, it saves the new_metadata. If it exists, it merges the new_metadata into it and rewrites the file only if the incoming metadata differs from what is stored. new_metadata is a snapshot of the metadata for the series that reported this run. The snapshot need not carry the same columns, or the same column order, as the stored roster. Rows are matched on time_series_id. A series that new_metadata covers is replaced wholesale, so a field the snapshot has stopped carrying is cleared for that series. A series that new_metadata omits keeps its last stored values indefinitely. The roster therefore holds every time series we have ever seen, not only the series in the latest snapshot.

This function is not safe under concurrent callers: it assumes it is called by one thread at a time, and takes no lock.

The rewrite is not atomic either. write_parquet overwrites the roster in place, with no write-to-temporary-file-and-rename. The roster therefore does not get the all-or-nothing commit that Delta gives the tables around it. See principle 10, every write is atomic and idempotent. A crash or an out-of-memory kill part-way through a local write leaves a partial file. pl.read_parquet below is what rejects that partial file on the next run, before TimeSeriesMetadata.validate ever sees it. The error reads ComputeError: parquet: File out of specification: The file must end with PAR1. validate is the guard for the other case: a roster that reads back cleanly but is off-contract, from an older writer or a hand-edit. Either way the asset records metadata_upsert_failed, and the roster stays broken until an operator acts. A corrupt file is not a missing file, so the create branch below never runs again by itself.

Deleting the file is not on its own a fix. power_time_series_and_metadata extracts metadata only from the files select_new_rows judged new. The next hourly run would therefore rebuild the roster from whichever series happened to publish that hour, rather than from every series NGED publishes. Nothing is permanently lost, and a rebuild is cheaper than re-reading the bucket. Every JSON file carries its own series' metadata in its top-level fields, and list_timeseries_json_files returns a time_series_id per key. The newest file per series is therefore enough — one download per time series rather than one per file NGED has ever published.

Parameters:

Name Type Description Default
new_metadata DataFrame[TimeSeriesMetadata]

The new metadata DataFrame.

required
metadata_path str

Local path or remote URI of the Parquet file where we store our version of the metadata.

required
storage_options ObjectStoreOptions | None

Object-store credentials/endpoint for a remote metadata_path; None/empty for a local path.

None

Returns:

Type Description
UpsertMetadataStats

An UpsertMetadataStats. metadata_n_new_TimeSeriesIDs counts the time_series_ids in

UpsertMetadataStats

new_metadata that are new to the roster. metadata_n_updated_TimeSeriesIDs counts the

UpsertMetadataStats

time_series_ids already in the roster that changed. metadata_updated_TimeSeriesIDs holds

UpsertMetadataStats

the sorted list of changed time_series_ids. Both counts are zero and the id list

UpsertMetadataStats

is omitted when the parquet file was up to date already. The id list is also omitted on

UpsertMetadataStats

a first-ever write, when every id in new_metadata counts as new rather than updated.

Source code in packages/nged_data/src/nged_data/storage.py
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
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
630
631
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
def upsert_metadata(
    new_metadata: pt.DataFrame[TimeSeriesMetadata],
    metadata_path: str,
    storage_options: ObjectStoreOptions | None = None,
) -> UpsertMetadataStats:
    """Upserts metadata to a Parquet file, keeping the newest version of each time series.

    If the Parquet file does not exist, it saves the new_metadata. If it exists, it merges the
    new_metadata into it and rewrites the file only if the incoming metadata differs from what is
    stored. ``new_metadata`` is a snapshot of the metadata for the series that reported this run.
    The snapshot need not carry the same columns, or the same column order, as the stored roster.
    Rows are matched on ``time_series_id``. A series that ``new_metadata`` covers is replaced
    wholesale, so a field the snapshot has stopped carrying is **cleared** for that series. A series
    that ``new_metadata`` omits keeps its last stored values indefinitely. The roster therefore
    holds every time series we have ever seen, not only the series in the latest snapshot.

    This function is not safe under concurrent callers: it assumes it is called by one thread at a
    time, and takes no lock.

    The rewrite is not atomic either. `write_parquet` overwrites the roster in place, with no
    write-to-temporary-file-and-rename. The roster therefore does not get the all-or-nothing
    commit that Delta gives the tables around it. 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).
    A crash or an out-of-memory kill part-way through a local write leaves a partial file.
    `pl.read_parquet` below is what rejects that partial file on the next run, before
    `TimeSeriesMetadata.validate` ever sees it. The error reads `ComputeError: parquet: File out of
    specification: The file must end with PAR1`. `validate` is the guard for the other case: a
    roster that reads back cleanly but is off-contract, from an older writer or a hand-edit. Either
    way the asset records `metadata_upsert_failed`, and the roster stays broken until an operator
    acts. A corrupt file is not a missing file, so the create branch below never runs again by
    itself.

    Deleting the file is not on its own a fix. `power_time_series_and_metadata` extracts metadata
    only from the files `select_new_rows` judged new. The next hourly run would therefore rebuild
    the roster from whichever series happened to publish that hour, rather than from every series
    NGED publishes. Nothing is permanently lost, and a rebuild is cheaper than re-reading the
    bucket. Every JSON file carries its own series' metadata in its top-level fields, and
    `list_timeseries_json_files` returns a `time_series_id` per key. The newest file per series is
    therefore enough — one download per time series rather than one per file NGED has ever
    published.

    Args:
        new_metadata: The new metadata DataFrame.
        metadata_path: Local path or remote URI of the Parquet file where we store our version
            of the metadata.
        storage_options: Object-store credentials/endpoint for a remote `metadata_path`;
            ``None``/empty for a local path.

    Returns:
        An `UpsertMetadataStats`. `metadata_n_new_TimeSeriesIDs` counts the `time_series_id`s in
        `new_metadata` that are new to the roster. `metadata_n_updated_TimeSeriesIDs` counts the
        `time_series_id`s already in the roster that changed. `metadata_updated_TimeSeriesIDs` holds
        the sorted list of changed `time_series_id`s. Both counts are zero and the id list
        is omitted when the parquet file was up to date already. The id list is also omitted on
        a first-ever write, when every id in `new_metadata` counts as new rather than updated.
    """
    COMPRESSION: Final[str] = "zstd"

    # The annotation is not enforced at runtime and this is the package's only public entry point,
    # so check the caller's snapshot rather than trust it.
    new_metadata = TimeSeriesMetadata.validate(new_metadata.sort("time_series_id"))

    if not object_exists(metadata_path, storage_options):
        log.info(f"Metadata file not found at {metadata_path}. Creating new file.")
        # write_parquet doesn't create missing parent directories, so a first-ever run against a
        # fresh local data root would fail here. This create branch runs before any Delta write that
        # would otherwise create the dir. Create the parent for a local metadata_path.
        if_local_path_then_make_parent_dir(metadata_path)
        new_metadata.write_parquet(
            metadata_path,
            compression=COMPRESSION,
            storage_options=typeddict_to_dict(storage_options),
        )
        return UpsertMetadataStats(
            metadata_n_new_TimeSeriesIDs=new_metadata.height,
            metadata_n_updated_TimeSeriesIDs=0,
        )

    existing_metadata = pl.read_parquet(
        metadata_path, storage_options=typeddict_to_dict(storage_options)
    )
    # The stored roster is outside this code's control: it can come from an older writer, a
    # hand-edit, or a truncated upload. An off-contract file must therefore not be merged blind into
    # the roster we write back. As with any raise from this function, the asset contains it rather
    # than failing: it records `metadata_upsert_failed` and lets the power write proceed (see
    # `defs/assets.py`).
    TimeSeriesMetadata.validate(existing_metadata)

    # `how="diagonal"` because the snapshot and the stored roster can differ in both width and
    # column order: four TimeSeriesMetadata fields are `allow_missing`. Aligning the two frames into
    # one also makes the `hash_rows` diff below insensitive to the stored column order. Hashing the
    # two frames separately would not be.
    combined = pl.concat([new_metadata, existing_metadata], how="diagonal")
    new_rows = combined.head(new_metadata.height)
    stored_rows = combined.slice(new_metadata.height)

    # Compare metadata. `metadata_diff` contains all rows in `new_metadata` that do not have an
    # exact match in `existing_metadata`. Adapted from https://stackoverflow.com/a/79888719
    metadata_diff = new_rows.filter(~new_rows.hash_rows().is_in(stored_rows.hash_rows().implode()))
    # The first frame carrying the union of both inputs' columns: the concat adds to the snapshot's
    # rows any `allow_missing` field only the stored roster had. All four fields are nullable, so
    # this validate is a shape check on a frame neither validation above saw, not a guard against a
    # known fault. Of the four validate calls in this function this is the weakest, and the first
    # to reconsider if the validate calls get trimmed.
    TimeSeriesMetadata.validate(metadata_diff)

    if metadata_diff.is_empty():
        log.info("TimeSeriesMetadata is up to date.")
        return UpsertMetadataStats(
            metadata_n_new_TimeSeriesIDs=0,
            metadata_n_updated_TimeSeriesIDs=0,
        )

    log.info(
        f"New TimeSeriesMetadata available for {metadata_diff.height} timeseries_ids."
        f" Updating {metadata_path}."
    )

    # Merge metadata. Put new_metadata first so that unique(keep="first") keeps the new version
    merged_metadata = combined.unique(subset="time_series_id", keep="first").sort("time_series_id")

    # The last gate before the stored roster is overwritten. `unique` draws rows from both sides of
    # the concat, so no validation above has seen this row set.
    TimeSeriesMetadata.validate(merged_metadata)

    merged_metadata.write_parquet(
        metadata_path, compression=COMPRESSION, storage_options=typeddict_to_dict(storage_options)
    )

    # Compute stats
    new_ids = set(new_metadata["time_series_id"]) - set(existing_metadata["time_series_id"])
    updated_ids = list(
        set(metadata_diff["time_series_id"]).intersection(existing_metadata["time_series_id"])
    )
    return UpsertMetadataStats(
        metadata_n_new_TimeSeriesIDs=len(new_ids),
        metadata_n_updated_TimeSeriesIDs=len(updated_ids),
        metadata_updated_TimeSeriesIDs=sorted(updated_ids),
    )