Lab 3 · From NEMWEB to a model-ready dataset¶
Intermediate Production track
Every chapter so far ran on synthetic data. A model that is going to run every five
minutes needs real data instead, and it needs it on time, complete and exactly
once. This lab builds that pipe for AEMO's public dispatch data. It downloads files from
NEMWEB, checks them, stores them in layers, watches how fresh they are, and gives models a
single way in: load_dataset("nem/dispatch_price", ...).
The code is in energy_or.nem, the tests in tests/test_nem.py, and a Colab notebook
(notebooks/lab03_nem_data_pipeline.ipynb) runs every step offline on three real files
that ship with the repository.
Data source and licence
Source: Australian Energy Market Operator (AEMO), NEMWEB DispatchIS reports. AEMO gives general permission to use its public material for any purpose, with accurate attribution to AEMO as author (AEMO copyright permissions, retrieved 5 October 2026). Every number on this page comes from a real run on 5 October 2026. Files edited to stage incidents are labelled as edited.
1 · Where the data lives¶
NEMWEB is a set of plain web folders. For dispatch, two of them matter:
| Folder | Holds | Kept | Observed on 5 Oct 2026 |
|---|---|---|---|
Reports/Current/DispatchIS_Reports/ |
one zip per five-minute interval | about 48 hours | 576 files, about 21 KB each |
Reports/Archive/DispatchIS_Reports/ |
one zip per trading day, holding that day's 288 five-minute zips | about a year (older data: the monthly MMSDM archive) | 375 days, 24 Sep 2025 to 3 Oct 2026; about 6 MB per day |
The file name carries the data's identity:
PUBLIC_DISPATCHIS_202610051330_0000000541277619.zip
└──────────┘ └──────────────┘
interval ENDING 13:30 AEST AEMO's file sequence number
Three conventions bite everyone at least once:
- Intervals are labelled by their end. "13:30" is the interval 13:25–13:30. A trading
day runs from the interval ending 00:05 to the one ending 00:00 the next day, so
trading_day(interval_end) = (interval_end − 5 min).date(). - Market time is AEST all year (UTC+10, no daylight saving), even when Sydney and
Melbourne clocks say AEDT. Columns carry the zone in their name:
interval_end_aest. - The file arrives before its interval ends. Dispatch is run at the start of the interval. Across the 742 files loaded on 5 October, AEMO wrote each one between 4.4 and 4.9 minutes before the end of its interval.
2 · One file, seven tables¶
An AEMO report is not one table but several, stacked in one CSV in the MMS format:
C,NEMP.WORLD,DISPATCHIS,AEMO,PUBLIC,2026/10/05,13:20:14,0000000541277005,... ← header
I,DISPATCH,PRICE,5,SETTLEMENTDATE,RUNNO,REGIONID,...,RRP,... ← table, version 5, columns
D,DISPATCH,PRICE,5,"2026/10/05 13:25:00",1,NSW1,...,90.65,... ← data
...
C,"END OF REPORT",1193 ← line count
The file for the interval ending 13:25 holds seven tables: case solution (1 row), local price (94), price (5), regional summary (5), interconnector results (7), constraints (1,067) and interconnection (5). The constraint rows hold each network constraint's right-hand side and marginal value, the raw material for Case J's constraint forecasting.
parse_mms_csv reads every table and checks the final line count. That count is the only
integrity check inside the file. A download cut short in transit loses it, and the parser
then raises MMSFormatError instead of loading part of an interval as if it were complete.
3 · Four layers, one rule each¶
flowchart LR
N[NEMWEB<br/>Current + Archive] -->|sync| R[raw<br/>zips, byte for byte]
R --> B[bronze<br/>every table, text]
B --> S[silver<br/>typed, keyed, units]
S -->|publish| G[gold<br/>Parquet by month]
G -->|load_dataset| M[models]
R -.-> F[(meta.files<br/>manifest)]
B -.-> F
| Layer | Holds | Rule |
|---|---|---|
| raw | the downloaded zips | never modified; stored under their SHA-256, so a revised file sits beside the original rather than over it |
| bronze | every MMS table, every value as text, plus _source_file and _row |
one DuckDB table per (group, table, version); reloading a file replaces that file's rows only |
| silver | dispatch_price, dispatch_region: typed columns with units in their names |
primary key on the natural key (interval end, region, run, intervention); a newer LASTCHANGED wins |
| gold | Parquet partitioned by year and month | what models read; rebuilt from silver by publish() |
The manifest meta.files records every file seen: its name, hash, size, URL, the archive
it came in, when it was retrieved, and whether it loaded or was quarantined. That
manifest is the lineage. Any gold row names its source_file, the manifest names that
file's raw bytes and URL, and the hash proves they have not changed.
Why is bronze a DuckDB table rather than one Parquet file per five-minute interval? Because that would be 288 tiny files a day. Small files are a classic data-lake performance problem, so the lake writes Parquet only at the gold layer, in monthly partitions.
4 · Idempotency: run it twice, change nothing¶
A scheduled job will fail halfway one day: a timeout, a full disk, a reboot. The simplest reliable recovery is to run the same job again. That is safe only if running it again cannot double-count anything. Every write in this pipeline is designed for that:
- Before downloading,
syncasks the manifest which intervals already loaded, and fetches only the rest. - The same bytes again are a
duplicate: the manifest'stimes_seengoes up and nothing else changes. - Different bytes under a known name are a
reloadedrevision. Bronze deletes that file's rows before inserting, and silver upserts:
INSERT INTO silver.dispatch_price ... SELECT ... FROM bronze.... WHERE _source_file = ?
ON CONFLICT (interval_end_aest, region_id, run_no, intervention) DO UPDATE SET ...
WHERE excluded.last_changed_aest >= silver.dispatch_price.last_changed_aest
The WHERE clause on the update is what makes order irrelevant. A newer revision
replaces an older one, and an older revision arriving late (a retry, an archive copy) is
ignored. Without it, the last file to arrive would win, and that is not always the
newest.
| Test | What it proves |
|---|---|
test_ingesting_twice_changes_nothing |
Every bronze and silver row is identical after a second load |
test_newer_revision_wins_and_a_late_older_one_is_ignored |
The upsert respects LASTCHANGED, not arrival order |
test_truncated_file_is_quarantined_then_recovered |
A broken file is quarantined, flagged by the health check, and cleared by a clean copy |
test_schema_drift_stays_in_bronze |
An unknown schema version is kept in bronze and kept out of silver, with a warning |
test_intervention_prices_from_pricing_run_quantities_from_physical |
Gold takes prices from the pricing run and quantities from the physical run |
test_sync_uses_current_then_the_archive_and_is_idempotent |
Archive fallback works; a second sync makes no network calls at all |
test_fetcher_retries_transient_errors_but_not_404 |
Backoff on 503/timeouts; a 404 is information, not a glitch |
test_freshness_lag_gaps_and_incomplete_intervals |
Lag, missing intervals and missing regions are each detected |
test_load_dataset_windows_are_half_open |
(start, end] on the interval end, as AEMO labels it |
Intervention pricing
When AEMO intervenes in the market (for example, directing a generator on for system
security), dispatch is solved twice. The physical run (INTERVENTION = 1) sets what
plant actually does, and the pricing run (INTERVENTION = 0) sets the price as if
the intervention had not happened. So gold dispatch_price keeps the pricing run, and
gold dispatch_region keeps the physical run. No interventions occurred in the
window loaded on 5 October; the test stages one with an edited file.
5 · Current first, then the archive¶
sync(lake, fetcher, start, end) makes the lake hold every interval ending in
(start, end]:
- List the intervals the lake does not already have. A quarantined file counts as missing, so re-running the job is also how it recovers.
- Download them from
Current, every version AEMO issued, in sequence order. - For anything
Currentno longer holds, download the daily archive and load the wanted intervals from inside it. - Report what is still missing. It does not raise an error: a late file is something to alert on, not a reason to throw away the 287 intervals that did arrive.
The first real run, for the window from 3 October 00:00 to 5 October 13:45 AEST:
window (2026-10-03 00:00, 2026-10-05 13:45] AEST: 741 intervals expected
already in the lake: 0
loaded from Current: 576, from the daily archive: 165
downloads: 577 (18.0 MB) in 214.8 s
gold: dispatch_price 3,705 rows, dispatch_region 3,705 rows
Current held the newest 576 intervals (48 hours). The first 165 of the window had
already rolled off, so they came from the archive for 3 October: one 6 MB download
instead of 165 small ones. Gold has 741 × 5 regions = 3,705 rows, none missing. A
re-run six minutes later downloaded one new file and finished in 2.8 seconds.
Fetching goes through a three-line Fetcher protocol (get(url) -> bytes), so the tests
substitute a fake that serves canned bytes, removes files from Current, or fails on
demand. HttpFetcher adds timeouts, up to four attempts with exponential backoff and
jitter for 408/429/5xx and dropped connections, and no retry on 404.
Windows gotcha: antivirus and TLS
The first live run on Windows died at once with
OPENSSL_Uplink ... no OPENSSL_Applink. The cause was the antivirus product, which set
the environment variable SSLKEYLOGFILE to a named pipe; Python's
ssl.create_default_context() tries to open that file and OpenSSL aborts the
process. HttpFetcher builds its TLS context directly with
ssl.SSLContext(PROTOCOL_TLS_CLIENT) and load_default_certs(). That skips the key
log but keeps full certificate and host-name verification. The same antivirus also
intercepts HTTPS, which is why uv needed --native-tls (see the handover notes).
6 · Is the data fresh?¶
check_freshness(lake, now_utc) answers three questions in SQL, so the same queries can
feed a dashboard or an alert rule:
- Freshness. How far is the newest interval behind the clock? Because files arrive before their interval ends, a healthy lag is negative or a few minutes. The status is late after 10 minutes and stale after 30.
- Completeness. Which intervals in the last 24 hours are missing? Which have fewer than the five regions?
- Quarantine. Which broken files are still waiting for a clean copy?
The missing intervals come from comparing a generated time grid with what the lake holds:
WITH grid AS (SELECT unnest(generate_series(?, ?, INTERVAL 5 MINUTE)) AS t),
have AS (SELECT interval_end_aest AS t, count(DISTINCT region_id) AS n_regions
FROM silver.dispatch_price WHERE intervention = 0 GROUP BY 1)
SELECT grid.t, coalesce(n_regions, 0) FROM grid LEFT JOIN have USING (t)
WHERE coalesce(n_regions, 0) < 5
On the real lake, between the first sync and the incremental one:
OK: newest interval ends 2026-10-05 13:45 AEST, +6.8 min vs the clock (13:51 AEST)
window 2026-10-04 13:50 to 2026-10-05 13:50: 289 intervals, 1 missing, 0 with fewer than 5 regions
missing: 10-05 13:50
verdict: needs attention
The data was fresh enough by the lag threshold, but the interval ending 13:50 had been
published and not yet fetched. The incremental sync fetched it, and the verdict became
healthy. scripts/nem_sync.py exits 0 when healthy, 1 when the data needs attention,
and 2 when the sync itself failed. Those three exit codes are all a scheduler needs to
raise an alert.
7 · The one door models use¶
Models never open the lake's database or walk its folders. They ask for a dataset and a window:
from datetime import datetime
from energy_or.nem import load_dataset
rel = load_dataset(
"nem/dispatch_price", region="SA1", start=datetime(2026, 10, 3), end=datetime(2026, 10, 4)
) # one trading day
df = rel.df() # pandas; also rel.pl() for Polars, rel.fetchnumpy() for NumPy
load_dataset reads only the gold Parquet files, so a notebook can read while a
scheduled job is writing. Its window (start, end] is half-open on the interval end,
matching AEMO's labels, so midnight to midnight is exactly 288 intervals.
What 3–4 October 2026 looked like (576 intervals per region, source AEMO):
| Region | Mean $/MWh | Min | Max | Share of intervals below $0 |
|---|---|---|---|---|
| NSW1 | 85.48 | −5.59 | 130.71 | 0.3 % |
| QLD1 | 54.13 | −19.93 | 117.90 | 34.7 % |
| SA1 | 72.55 | −100.00 | 390.20 | 21.5 % |
| TAS1 | 46.72 | −46.74 | 149.91 | 12.3 % |
| VIC1 | 45.88 | −44.84 | 167.50 | 21.0 % |
And at 13:25–13:35 on 5 October, South Australia's total demand was negative:
−82, −99 and −137 MW. AEMO's TOTALDEMAND is the demand that the generators it
dispatches must meet; rooftop solar is not among them and shows up as lower demand. At
midday it more than covered the region's load, so the region needed less than nothing
from its dispatched generators. Chapters 2 and 9 priced a battery against synthetic
prices; these are the real conditions it would face.
Two details in the gold data deserve a model's attention:
- Price status. Six consecutive intervals on 4 October (09:30–09:55) were published
with
price_status = 'NOT FIRM', AEMO's flag for prices that may still be reviewed. Gold keeps the column, so a backtest can include or exclude those intervals on purpose. - Lineage. Every gold row carries its
source_file, so a surprising price can be traced to the exact AEMO file and its hash.
8 · What it costs¶
Lab 7 will budget the whole five-minute cycle; here are this lab's measured numbers:
| Step | Measured (5 Oct 2026, one laptop) |
|---|---|
| Download one five-minute file | about 0.3 s (sequential, one connection) |
| Parse + bronze + silver for one file | about 70 ms (the first file about 1.2 s: tables are created) |
| Incremental sync (one new file) + publish + health check | 2.8 s end to end |
| 2.6 days of history (741 intervals) | 215 s; 18 MB downloaded; raw 19 MB, DuckDB 38 MB, gold 0.2 MB |
The DuckDB file is dominated by bronze constraint rows (780,000 of 835,000 rows). That is the price of keeping every table. Exercise 4 asks when that stops being worth it.
Exercises¶
Guided
Load the three sample files in reverse order (notebook section 2). Is the lake identical? Which line of the upsert makes the order irrelevant?
Engineering
Add DISPATCH.INTERCONNECTORRES as a silver table (flow_mw, losses_mw,
marginal_value_per_mwh) keyed on interval, interconnector, run and intervention.
Write the test first: ingest twice, nothing changes.
Market
With a day of live data, find the intervals where SA1's price differs from VIC1's by more than $50/MWh. What does that tell you about the Heywood interconnector at those moments? Check your answer against the interconnector table in bronze.
Challenge
Bronze keeps about 300,000 constraint rows a day. Design a retention rule: which tables stay in bronze for how long, given that raw keeps every byte and bronze can be rebuilt from raw? What is the rebuild time, and who needs it?
Production
The daily archive keeps about a year; older days are only in the monthly MMSDM archive.
Extend sync to fall back to it. What should the manifest record about a file that
came from three levels of fallback, and how would a health check notice that a whole
month is missing?
Production perspective¶
- Scheduling.
scripts/nem_sync.py --hours 1every five minutes is enough: it is idempotent, it looks back over a whole hour (so a missed run heals itself), and it exits 1 or 2 for the scheduler to alert on. On Cloudflare, a Cron Trigger can call a container running it, with the lake on R2 instead of a local disk; only the paths change. - One writer, many readers. The DuckDB file has one writer, the sync job. Notebooks, backtests and the optimiser read gold Parquet only.
- What to alarm on. A non-zero exit; lag beyond 10 minutes; any missing interval in the last hour; any quarantined file; any schema-drift warning (AEMO changed a table and silver needs a new mapping).
- What it does not do yet. It loads prices and regional quantities only. NEMDE's constraint equations, bids and unit dispatch come in later labs. It also makes no claim to reproduce AEMO's dispatch: it stores what AEMO published.