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

kitchin-services.cheme.cmu.edu, port 443, WebSocket path /mqtt

Topics

plant/tep/birth (one retained message), plant/tep/telemetry (samples)

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 XMEAS_* and 12 manipulated variables XMV_*

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

sourceTimestamp (per metric)

plant

event time, when the instrument measured it. Window on this.

timestamp (per message)

plant

when the batch was assembled

receivedAt (envelope)

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.

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

tag

string

metrics[].name

value

double

metrics[].value

status

string

metrics[].statusCode

event_time

timestamp, UTC

metrics[].sourceTimestamp

server_time

timestamp, UTC

the message timestamp

received_at

timestamp, UTC

the envelope receivedAt

seq

int

the message seq

is_historical

bool

the message isHistorical

  • 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 in results/explain.txt. Later stages may use pandas or Polars.

  • Getting at the nesting. Keep only plant/tep/telemetry lines. payload is a struct, so pl.col("payload").struct.field("seq") pulls out one field (.unnest("payload") pulls out all of them). metrics is a list of structs: .explode("metrics") gives one row per reading, and pl.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 with pd.json_normalize, if you want to see the shape of the result first.

  • Watch out: pl.scan_ndjson infers 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. Pass infer_schema_length=None or declare the schema.

3. Validate and quarantine#

  • Generate a pandera schema from the birth message’s tags list: per-tag dtype and plausible range, status in the set the stream uses, event_time inside 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 value against 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, receivedAt on 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_cases gives you the offending rows: its index column points back at them.

  • Write every Bad_* and Uncertain_* reading, and everything the schema rejects, to processed/quarantine.parquet with a reason column. Remove them from tidy.parquet and 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 on seq: 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_time sits behind the frontier, the newest event_time in 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.csv with 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 whose event_time is within the allowance of the newest event_time you received.

  • Compare clocks: compute one aggregate in 30-minute windows on event_time, and again on received_at converted to the plant clock (the formula is under Three clocks, above), and report where and by how much they disagree. Polars group_by_dynamic or pandas resample works 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_range or pd.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 _filled column per tag, or a long list of filled cells. Take the mask before you fill. Polars has forward_fill, backward_fill and interpolate; pandas has ffill, bfill and interpolate. 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 Bad null 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:

  1. The three clocks. Which one you windowed on and why. The largest sourceTimestamp to message timestamp gap you measured, and its cause.

  2. What you collected. Messages, records, duplicates removed, Uncertain and Bad counts, records quarantined and why, fraction late.

  3. 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.

  4. 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.

  5. Event time against processing time. The same aggregate both ways, the disagreement, where it was worst, and why.

  6. 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

raw/stream-*.ndjson, append-only

Pipeline code

src/…/*.py, or flat .py files

pandera schema

any .py (a notebook is allowed)

Tidy table

processed/tidy.parquet

Filled sample grid

processed/grid.parquet

Quarantined records

processed/quarantine.parquet

Watermark sweep

results/watermark.csv

Windowed aggregate

results/windows.csv

Report

REPORT.md

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 broken uv.lock in 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; receivedAt non-decreasing; raw file not rewritten

2

Decode and reshape

1.0

tidy.parquet with the eight columns and types; Polars LazyFrame with one .collect(); three clocks distinct and populated; more than one event_time per server_time; record count reconciles with the raw file

3

Validate

1.0

pandera schema generated from the birth tags; lazy=True; Bad/Uncertain in quarantine.parquet and not in tidy.parquet; report counts reconcile

4

Dedup and watermark

1.25

(tag, event_time) unique; duplicate count matches the script’s; watermark.csv with several allowances and completeness and still-waiting columns; clock comparison in the report

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 windows.csv on 30-minute event-time buckets that the script can recompute; incomplete last window handled; disturbance interval named with tags

6

REPORT.md

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 sourceTimestamp advances.

  • Replay the raw file one message at a time with a bounded buffer, and report what a watermark would have emitted and when.