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, parsingtime_series_id,start_time, andend_timeout 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, sodownload_and_parse_filesnever 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 aDownloadAndParseResultofmetadata(TimeSeriesMetadata),power_time_series(PowerTimeSeries), andn_implausible_power_rows_dropped. RaisesNoNewDataif the listing was empty, or if every listed file'sdatafield was null.nged_data.storage.select_new_rows(time_series, delta_path, storage_options=None)— filterstime_seriesdown to rows genuinely missing from thepower_time_seriesDelta table atdelta_path.PowerTimeSeriesrows are filtered by existence, an anti-join on(time_series_id, time), so a reading is kept regardless of arrival order. The file listing thatlist_timeseries_json_filesreturns is filtered by comparing each file'send_timeagainst 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 observationtimeon disk for eachtime_series_idin thepower_time_seriesDelta table.nged_data.upsert_metadata(new_metadata, metadata_path, storage_options=None)— merges aTimeSeriesMetadatasnapshot into the stored metadata Parquet file, keeping the newest values pertime_series_idand 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 | |
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 | |
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 | |
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 | |
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 | |
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 | |
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 | |
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
|
required |
Returns:
| Type | Description |
|---|---|
DownloadAndParseResult
|
A |
DownloadAndParseResult
|
|
DownloadAndParseResult
|
docstring for what each field holds. |
Raises:
| Type | Description |
|---|---|
NoNewData
|
if |
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 | |
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 | |
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 | |
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 |
None
|
Returns:
| Type | Description |
|---|---|
UpsertMetadataStats
|
An |
UpsertMetadataStats
|
|
UpsertMetadataStats
|
|
UpsertMetadataStats
|
the sorted list of changed |
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 |
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 | |