Assignment 3: A plant stream, collected and made trustworthy#
Module: Week 03 · Released: L6 · Due: ~1 week later · 6 points on Canvas.
Goal#
Collect ten minutes of a live chemical plant stream, then build a batch pipeline that turns it into a table you would compute on. The stream arrives late, out of order, duplicated, with gaps and suspect readings. You will land it, validate it, deduplicate it, fill a regular grid, and find the disturbance in it.
You need nothing beyond your project and uv. There is no account, and the connection runs over port 443, so it works from any network.
You should be able to:
land an unbounded stream durably before processing it;
window over event time, and choose a watermark and state its cost;
make a pipeline idempotent under at-least-once delivery;
validate with pandera and quarantine bad records instead of deleting them;
reshape long to wide, fill a ragged grid, and record which cells you invented;
write a stage against a Polars lazy query plan and read
.explain().
The stream#
A Tennessee Eastman plant simulation, the same plant L5 used as files, published live over MQTT.
Host |
|
Topics |
|
Speed |
plant runs 144× wall clock; one sample per 180 plant seconds, one message per 1.25 s |
Ten minutes |
one plant day: 480 messages, ~25,000 readings, ~3 MB |
Tags |
53: 41 measurements |
Disturbances |
one episode every 12 plant hours, lasting 1.5 to 4 h, then clears; your ten minutes holds two |
L5’s files had 52 columns because they drop XMV_12, the agitator speed, which sits constant. Expect it here.
The birth message is your data contract#
It lists every tag with its description and unit, plus the clock conversion. Build your schema from it rather than typing the tag list.
{"type": "birth", "plant": "tep", "accelerationFactor": 144,
"plantEpoch": "2026-09-09T14:55:46.480Z", "sampleIntervalSeconds": 180,
"tags": [{"name": "XMEAS_1_A_feed", "description": "A Feed (stream 1)", "unit": "kscmh"},
{"name": "XMEAS_2_D_feed", "description": "D Feed (stream 2)", "unit": "kg/hr"}]}
Two of the 53 tags are shown; the real message carries all of them.
Use the descriptions in your report: “the separator pressure moved” is a finding, “column 13 moved” is not.
A telemetry message, as the collector saves it#
{"receivedAt": "2026-09-09T14:56:29.712Z", "topic": "plant/tep/telemetry",
"payload": {"seq": 197, "timestamp": "2026-09-09T16:37:45.479Z", "isHistorical": false,
"metrics": [
{"name": "XMEAS_1_A_feed", "value": 0.253256, "statusCode": "Good",
"sourceTimestamp": "2026-09-09T16:37:45.479Z"},
{"name": "XMEAS_23_feed_A", "value": 31.551008, "statusCode": "Good",
"sourceTimestamp": "2026-09-09T16:34:45.479Z"}]}}
Three clocks#
Field |
Clock |
What it is |
|---|---|---|
|
plant |
event time, when the instrument measured it. Window on this. |
|
plant |
when the batch was assembled |
|
your laptop |
processing time, when it landed |
The first two run on the plant clock, 144× fast from plantEpoch. Subtracting receivedAt from either gives a meaningless number until you put receivedAt on the plant clock too. Both clocks agree at plantEpoch, and every wall second after it is 144 plant seconds:
receivedAt on the plant clock = plantEpoch + (receivedAt − plantEpoch) × accelerationFactor
For the message below, receivedAt is 43.232 wall seconds after plantEpoch, which is 6,225.4 plant seconds, so the message landed at 16:39:31.9 on the plant clock, 106 plant seconds after its timestamp. In Polars a datetime minus a datetime is a duration, and a duration times an integer is a duration, so this is one expression.
Nineteen tags report the past. XMEAS_23 to XMEAS_41 are composition analysers. Each reports its previous sample, so its sourceTimestamp lags the message by 0 to 720 s and repeats across several messages. Assume one timestamp per message and those nineteen channels are wrong by up to twelve minutes.
What goes wrong on the link#
The link drops out, the publisher buffers, and the backlog drains when it reconnects, so late data comes in bursts. Over 6,000 messages at these settings:
records arriving after a newer one |
13.6% |
lateness, plant minutes |
median 51, p95 117, max 198 |
samples never delivered |
20 in 5,961 |
duplicate deliveries |
59 |
|
648 / 146 in 318,000 |
Your numbers will differ. Every late record carries isHistorical: true.
Tasks#
1. Collect#
Run the collector for at least ten minutes:
uv run --no-project https://kitchingroup.cheme.cmu.edu/f26-06763/a03-collect.py --minutes 10
It appends to raw/stream-<date>.ndjson and never truncates, so rerunning adds data. Treat that file as read-only. Every later stage replays from it, and it is the one thing you cannot regenerate.
2. Decode and reshape (Polars, lazy)#
Explode metrics[] into one row per reading and write processed/tidy.parquet with exactly these columns:
Column |
Type |
From |
|---|---|---|
|
string |
|
|
double |
|
|
string |
|
|
timestamp, UTC |
|
|
timestamp, UTC |
the message |
|
timestamp, UTC |
the envelope |
|
int |
the message |
|
bool |
the message |
Write this stage in Polars on a
LazyFrame: scan, express the explode and casts as expressions, and.collect()once. Print.explain()and keep the output, for example inresults/explain.txt. Later stages may use pandas or Polars.Getting at the nesting. Keep only
plant/tep/telemetrylines.payloadis a struct, sopl.col("payload").struct.field("seq")pulls out one field (.unnest("payload")pulls out all of them).metricsis a list of structs:.explode("metrics")gives one row per reading, andpl.col("metrics").struct.field("name")then reads each field. Parse the three timestamps with.str.to_datetime(...)to UTC. The L6 notebook does the same reshape in pandas withpd.json_normalize, if you want to see the shape of the result first.Watch out:
pl.scan_ndjsoninfers the schema from the first 100 lines. Your file mixes birth and telemetry messages, and a birth message past line 100 loses its fields silently. Passinfer_schema_length=Noneor declare the schema.
3. Validate and quarantine#
Generate a pandera schema from the birth message’s
tagslist: per-tag dtype and plausible range,statusin the set the stream uses,event_timeinside your collection window.A reading’s range depends on its unit, and the birth message gives each tag’s unit, so build a lookup from tag to bounds and check
valueagainst it with a check on the whole frame (pa.DataFrameSchema(columns, checks=pa.Check(fn))), as the L6 notebook does with its unit ranges.Take the collection window from when you collected,
receivedAton the plant clock, with some slack for the replayed backlog. A window computed from the event times it checks can never fail.Validate with
lazy=True, so one run reports every failure.SchemaErrors.failure_casesgives you the offending rows: itsindexcolumn points back at them.Write every
Bad_*andUncertain_*reading, and everything the schema rejects, toprocessed/quarantine.parquetwith areasoncolumn. Remove them fromtidy.parquetand report the counts.Show the gate works: copy a few hundred clean rows, break a handful on purpose (a tag the birth message does not declare, a value outside its unit’s range, a missing value, an unknown status, a timestamp from another year), validate the copy with
lazy=True, and show which checks caught which rows.
Install pandera[pandas] and import pandera.pandas as pa. The old import pandera as pa still works but warns that it will stop.
import pandera.pandas as pa
schema = pa.DataFrameSchema({
"tag": pa.Column(str, pa.Check.isin(tag_names)),
"value": pa.Column(float, nullable=False),
"status": pa.Column(str),
})
try:
schema.validate(df, lazy=True)
except pa.errors.SchemaErrors as e:
bad = e.failure_cases # column, check, failure_case, index
4. Deduplicate and choose a watermark#
Deduplicate on
(tag, event_time)and report how many rows that removed. Do not key onseq: it is one byte and wraps at 256 within your collection. Expect fewer than 53 rows per sample afterwards, because a latched analyser’s repeats collapse to one measurement.Show idempotence: run the whole pipeline twice over the same file and show nothing changes.
Measure lateness. A reading’s lateness is how far its
event_timesits behind the frontier, the newestevent_timein any message that arrived before its own, or zero if it is not behind. The L6 notebook computes exactly this in pandas.Sweep a watermark over allowances of 0, 15, 30, 60 and 120 plant minutes (more if you like). Write
results/watermark.csvwith one row per allowance and at least these columns:allowance_plant_min, the allowance;completeness_pct, the percentage of readings whose lateness is within the allowance, so they arrive in time to be counted;still_waiting, the readings not yet final when you stopped collecting: those whoseevent_timeis within the allowance of the newestevent_timeyou received.
Compare clocks: compute one aggregate in 30-minute windows on
event_time, and again onreceived_atconverted to the plant clock (the formula is under Three clocks, above), and report where and by how much they disagree. Polarsgroup_by_dynamicor pandasresampleworks on either column.
5. Fill the grid, window, find the disturbance#
Grid. Pivot the deduplicated, quarantined table onto a regular grid, one row per 180-plant-second instant from first to last sample and one column per tag. Save it as processed/grid.parquet.
A pivot only has rows for instants that delivered something. Build the spine of instants first (
pl.datetime_rangeorpd.date_range) and join the pivot onto it, so a lost sample is an empty row rather than a missing one, and add a column for any tag that never arrived.Before filling, count empty cells separately for the 34 continuous channels and the 19 analysers. The birth message marks analysers with
analyserIntervalHours. Expect roughly 1% and 58% empty. Most continuous holes sit in the grid’s first few rows, which exist only because lagged analyser samples carry earlier event times than the first continuous sample; lost messages and quarantined readings make up the rest. Say where yours came from.Fill every cell, and record which cells you filled: a boolean mask, a
_filledcolumn per tag, or a long list of filled cells. Take the mask before you fill. Polars hasforward_fill,backward_fillandinterpolate; pandas hasffill,bfillandinterpolate. The first rows have nothing before them, so decide what they get.The two kinds of hole differ. An empty analyser cell means no new analysis yet, so carrying forward repeats what the instrument says. An empty continuous cell means a lost message, so carrying forward invents a value. Justify your fill for each.
Quarantine before you fill, or a
Badnull becomes an invented number.
Windows. Aggregate the filled grid into tumbling 30-plant-minute windows over event time: pandas pd.Grouper(key="event_time", freq="30min") or .resample(), Polars group_by_dynamic. Write results/windows.csv in wide form, one row per window and one column per tag. Your last window is incomplete; count the samples in each window and say how you handled it.
Disturbance. Find when the plant was disturbed, when it recovered, and which tags moved. A per-window score against a quiet stretch is enough: z-score each continuous tag’s window mean against the mean and standard deviation of windows you believe are quiet, average the absolute z-scores across tags, then look at the windows that stand out and the tags with the largest scores in them. A baseline that includes the disturbance inflates its own standard deviation and scores the disturbance as nearly normal. You are graded on method and honesty, not on naming the fault. Some disturbances are slow and genuinely hard to separate from noise, and saying so is a good answer.
6. Write REPORT.md#
Two pages maximum, six sections in this order:
The three clocks. Which one you windowed on and why. The largest
sourceTimestampto messagetimestampgap you measured, and its cause.What you collected. Messages, records, duplicates removed,
UncertainandBadcounts, records quarantined and why, fraction late.What you filled. Empty cells before filling, continuous against analyser, and the cause of each. Your fill for each, with the argument. How a reader tells measured from invented.
The watermark. Your sweep table, the allowance you would pick for a real plant, and why. What being wrong costs in each direction. A choice with no reason gets about half.
Event time against processing time. The same aggregate both ways, the disagreement, where it was worst, and why.
The disturbance. What, when, which tags, how confident, and what you could not tell from noise.
Also include two or three sentences on one thing .explain() showed the planner doing that your code did not ask for, and one line disclosing generative-AI use.
Names to use#
Thing |
Name |
|---|---|
Raw capture |
|
Pipeline code |
|
pandera schema |
any |
Tidy table |
|
Filled sample grid |
|
Quarantined records |
|
Watermark sweep |
|
Windowed aggregate |
|
Report |
|
If you used other names, the script searches for your files and prints what it found on page one. Correct it with a flag if needed:
uv run --no-project a03-evidence.py --andrew-id yourid --tidy out/readings.parquet
Submit#
From your project root, after running the pipeline:
uv run --no-project https://kitchingroup.cheme.cmu.edu/f26-06763/a03-evidence.py \
--andrew-id yourid --name "Your Name"
Upload the resulting evidence.pdf to Canvas. There is nothing to push.
Keep
--no-project, so a brokenuv.lockin your project cannot stop the script.The script opens your files read-only. It re-parses your raw file independently, checks your tables against it, and embeds your code, schema and
REPORT.md.Read the PDF before uploading. It prints your score by group. If a check fails, fix it and rerun, and submit anyway if you run out of time.
The PDF prints the script’s sha256, which matches https://kitchingroup.cheme.cmu.edu/f26-06763/a03-evidence.py.sha256 when run from the URL.
Grading#
Five groups are scored by the script and the sixth by your TA. Within a group, checks are equally weighted.
# |
Group |
Pts |
Checks |
|---|---|---|---|
1 |
Collection |
0.75 |
raw NDJSON with the birth message; at least 20 plant hours of event time; |
2 |
Decode and reshape |
1.0 |
|
3 |
Validate |
1.0 |
pandera schema generated from the birth tags; |
4 |
Dedup and watermark |
1.25 |
|
5 |
Grid and windows |
1.0 |
regular 180 s grid, one column per tag, no empty cells, fill record; pre-fill hole counts reconcile; wide |
6 |
1.0 |
the six sections and the AI-use line, read by your TA |
A check the script cannot decide is held for your TA rather than lost. Your TA can adjust in either direction when the work is clearly better or worse than the checks show. Report numbers that do not reconcile with your raw file will be asked about.
AI use#
Generative AI is allowed with disclosure in REPORT.md. Editing the generated PDF by hand is falsifying a submission. You must be able to explain your dedup key, your watermark choice, and why windowing on received_at gives a different answer.
Stretch (not graded)#
Collect for an hour and check whether your watermark still holds over six plant days.
Collect at the same time as a classmate and compare files: same plant, different losses.
Reconstruct each analyser’s sampling schedule from when its
sourceTimestampadvances.Replay the raw file one message at a time with a bounded buffer, and report what a watermark would have emitted and when.