← Personal projects Data engineering Graduate work · 2026

A streaming sensor pipeline, from MQTT to dashboard

End-to-end device telemetry: MQTT ingestion, schema validation at the edge, MongoDB as the raw store, and an aggregation layer that keeps queries fast as the collection grows. Built around the failure cases first, because those are the only parts that matter at 3 a.m.

  • Python
  • MQTT
  • MongoDB
  • Docker
MQTT → Mongo
Ingestion path
Idempotent
Write semantics
Raw + rolled-up collections
Store

Architecture

  1. 01Devicespublish to topics
  2. 02BrokerMQTT, QoS 1
  3. 03Consumervalidate + tag
  4. 04MongoDBraw, indexed
  5. 05Aggregationrolled-up reads

The problem

Sensor data is deceptively easy to demo and genuinely hard to run. A publisher, a broker, and a subscriber that prints readings is twenty lines of Python. What makes it a pipeline rather than a demo is everything that happens when the happy path stops holding.

So the design question I actually cared about was: what does this do when the broker drops, when a device’s clock is wrong, when the same reading arrives twice, and when the collection is large enough that a naive query stops returning in reasonable time?

Shape of it

Devices publish to topics on an MQTT broker. A consumer subscribes, validates each message against an expected shape before anything is written, tags it with ingestion metadata, and writes it into MongoDB as an immutable raw record. A separate aggregation layer produces the rolled-up documents that reads actually hit.

The split between raw and aggregated is deliberate. Raw records are never mutated, which means a bug in aggregation logic is recoverable — you fix the logic and recompute — rather than permanent.

The failure cases, and what each one forced

Duplicate delivery. MQTT at QoS 1 guarantees at least once, which means the pipeline has to be built for more than once. Writes are keyed on a deterministic identity — device, metric, and reading timestamp — so a redelivered message overwrites itself instead of double-counting. This is the single design decision that makes replay safe, and it has to be made at the start; retrofitting idempotency onto an append-only store is a migration, not a patch.

Late arrival. A device that loses connectivity and reconnects flushes its queue, and those readings arrive with timestamps well behind the current window. Records carry both an event timestamp and an ingestion timestamp, so a late reading lands in the right bucket rather than the current one, and it’s still possible to ask “what did we know at the time” separately from “what was actually true.”

Malformed payloads. Validation happens before the write, not after. Anything that fails lands in a dead-letter path with the raw bytes intact — because a malformed message is a signal about a device, and discarding it discards the signal.

Query cost. Time-series reads are almost always “recent window, one device” or “aggregate over range,” and both want compound indexing on device plus event time. The aggregation layer exists so that dashboard reads don’t scan raw at all.

Running it

Packaged with Docker Compose — broker, consumer, and database come up together, which makes the whole thing reproducible on any machine and made testing failure modes tractable: killing the broker container is a one-line way to check that the consumer reconnects and doesn’t lose or duplicate anything on the way back.

What I’d do differently

MongoDB was the right call for the assignment and the wrong call for the shape of the data. This is time-series, and a purpose-built time-series store — Timescale, or Mongo’s own time-series collections — gives you the partitioning, retention, and downsampling for free instead of hand-rolling the aggregation layer. If I rebuilt this, that’s the first change, and it would delete more code than it added.

Contact

Want the longer version?

Happy to walk through any of this in detail — the parts that broke are usually the interesting bit.