Skip to main content

Ten Thousand Sensors, One Question: Can I Trust the Alert?

00:08:34:40

I want you take a momment and imagine you're in a room where the temperature sensor is reporting 25.5 degrees. After an hour you check, and it's still 25.5. And every hour you check it's still the same.

Is the room stable, or has the sensor stopped working?

This is an interesting question because the number itself doesn't give you any straight-forward information. You will need to know what came before it, whether the readings are actualy arriving on time, and how the sensor actually works.

That is the part of anomaly detection I find interesting: how much has to happen before a system can reasonably say, you should go check this.

My IoT anomaly detection project brings together Kafka, Spark Structured Streaming, an LSTM autoencoder, and a few supporting services. The workload target is 10,000 simulated sensors, each sending one reading per second. It gives me a concrete setting for exploring streaming systems and machine learning together.

But the architecture becomes much more interesting when I follow one reading through it. Every stage makes a decision about what to keep, what to forget, and how long to wait. Those decisions eventually become the alert.

Before the model sees anything

The producer simulates temperature readings and injects variations intended to represent spikes, drift, stuck values, and noise. Each message carries a sensor ID, a timestamp, a value, and some metadata.

Synthetic data is useful here because I control the inputs. I can introduce something unusual and follow its path. It also gives the experiment a boundary: detecting patterns I generated myself says little about how the system would handle an unfamiliar piece of equipment.

The producer uses asynchronous sends, gzip compression, and a short batching delay. That delay is a small but revealing tradeoff. Waiting briefly lets the producer group messages more efficiently; sending immediately would reduce that waiting time but create more overhead.

At the target rate, the arithmetic is 10,000 events per second, or 864 million per day. That is a workload calculation, not a measured throughput result. It is enough to make retention and replay practical concerns, though. How much history should I keep? How quickly could I catch up after an outage?

Kafka gives the readings somewhere to wait while processing catches up. Messages are keyed by sensor ID, which keeps a sensor's records together under a stable partition layout. Ordering within a partition is useful, although it does not tell me whether the device's clock was correct or whether a reading arrived late.

Here is the intended route through the system:

text
Simulated sensors
        |
        v
Kafka: raw-readings
        |
        v
Spark: windowed features
        |
        v
TensorFlow Serving: inference
        |
        +--> Kafka: anomalies
        |
        +--> MinIO: features and scores

I find the arrows more useful than the boxes. An arrow hides serialization, waiting, retries, and assumptions about the shape of the data. It is also where two individually reasonable components can disagree.

A minute of context changes the question

One reading is a measurement. A series of readings can describe behavior.

The Spark job groups readings by sensor into 60-second windows, sliding forward every 10 seconds. It calculates six features: mean, standard deviation, minimum, maximum, count, and range.

Those features give me several ways to look at the same minute. The mean describes its general level. The maximum can preserve evidence of a brief spike that the mean dilutes. Standard deviation and range describe variation. Count tells me how much evidence went into the summary.

That last one is easy to overlook. A quiet sensor and a sensor with missing readings can both produce a reassuring average. I would want to know whether that average came from sixty observations or three.

The windows overlap. A reading can contribute to six windows, so the system revisits much of the same history as time advances. This makes the summaries more continuous, but it also creates state to maintain and correlated outputs. Six neighboring alerts may describe one incident.

There is another detail in the code that deserves more attention than its single line suggests:

python
.withWatermark("timestamp", "30 seconds")

The watermark lets Spark manage old window state as event time advances, with an allowance for late data. It is not a promise that every reading receives an answer within thirty seconds. For finalized window results, the system must wait for the watermark to pass the window's end; if event time stops advancing, that wait can stretch too.

This is why I would be careful with a phrase like “one-second detection.” One second spent computing is different from one second between a physical event and an actionable alert. The window, late-data allowance, processing backlog, and inference call all contribute to the latter.

For a sudden dangerous temperature, I would consider a direct threshold check alongside the windowed detector. For a subtle change in behavior, waiting for context may be worthwhile. The useful deadline depends on what someone needs to do with the answer.

What the autoencoder actually knows

The model is an LSTM autoencoder. It takes a sequence of feature vectors, compresses it, and tries to reconstruct it. The reconstruction error measures how far its output is from the input.

The idea is appealing: learn representative behavior, then investigate sequences the model struggles to reproduce.

But a large reconstruction error only means the input was difficult for this model to reconstruct. It does not tell me a machine is broken. A new operating regime, a different sensor baseline, or a preprocessing mistake could also produce that result.

The training implementation makes this distinction especially important. It standardizes six features and builds sequences of fifty rows. However, those rows come from shuffled synthetic feature summaries. They are not fifty chronological windows from one sensor, and the training split includes generated anomalies.

So I cannot treat the presence of LSTM layers as evidence that the model has learned how a real sensor evolves over time. The sequence needs to carry meaningful history first.

For the next version, I would build training examples from ordered readings passed through the same windowing logic as the live stream. I would keep each sensor's history together and evaluate on later time periods. That would make the experiment answer a more useful question: can the model recognize unfamiliar behavior after learning from the past?

I would also compare it with simple rules for high values, low variation, and changes from a rolling baseline. A neural network should earn its complexity through a useful improvement. I want that comparison in the project.

The handoff I need to get right

Looking at training and serving side by side exposes the most important unfinished part of this system.

Training expects a scaled sequence with shape (batch, 50, 6) and produces a reconstruction. The current streaming path sends a single unscaled six-feature vector and treats the first returned value as an anomaly score. Those are different interfaces. The streaming path also uses a configured threshold instead of loading the threshold selected during evaluation.

This is a concrete gap between the intended architecture and the implementation. Getting the services running would not resolve it.

The serving boundary needs to own a complete, explicit transformation:

text
Ordered windows for one sensor
  -> apply the training scaler
  -> assemble a 50-window sequence
  -> reconstruct it
  -> calculate reconstruction error
  -> compare with the matching threshold

Feature order, scaling, sequence length, model weights, and threshold belong together as one versioned unit. Changing any of them changes the meaning of the result.

Even a threshold of 0.95 needs an explanation. It is not automatically 95% confidence or 95% accuracy. Here, the intended score is reconstruction error, whose scale depends on preprocessing and the model. The threshold needs evidence from validation data and a decision about the cost of missed incidents versus unnecessary alerts.

There is a second behavior I would change before relying on the output: the streaming code assigns a score of zero when inference fails. With a positive threshold, that becomes a non-anomaly.

I want an explicit unscored state. A failed request tells me something about the detector's availability; it tells me nothing about the sensor's health. Quiet dashboards should never make that distinction disappear.

Where the system can fall behind

Spark distributes the feature aggregation, but the current scoring callback calls collect(), brings the batch onto the driver, and sends an HTTP inference request for each row in sequence.

That is an easy path to read and a clear place to improve. It concentrates both memory use and network waiting in one process. Adding Spark workers would not remove that bottleneck.

Suppose a batch contained 10,000 feature rows and each request took 5 milliseconds. Serial requests alone would take about 50 seconds. This is an illustrative calculation, not a benchmark of the project, but it shows why I would examine the serving boundary before adding infrastructure.

I would start with bounded batches of inference requests, then measure batch duration, memory, and Kafka lag. Distribution only helps when the work that takes time can actually run in parallel.

Recovery needs the same care. The callback writes anomalies to Kafka and appends records to MinIO. If one write succeeds and the other fails, retrying the batch can repeat a side effect. A checkpoint helps recover streaming progress; these writes still need a strategy for retries and duplicates.

One option I would explore is a stable result identity built from sensor ID, window boundaries, and model version, paired with storage or consumers that deduplicate on that identity. I want replay to be something I can use deliberately, including when comparing a new model against an older one.

The alert is where I would judge the project

The repository includes deployment and monitoring configuration, but I would judge readiness by more than healthy containers. I would want to see whether the stream is falling behind, how old the latest scored window is, how many windows remain unscored, and whether repeated flags are becoming repeated notifications.

Model evaluation also needs to describe the experience of using the detector. Precision asks how many flags were useful. Recall asks how many known incidents it caught. Detection delay asks whether it caught them in time. Alerts per sensor per day makes the interruption cost visible.

An overall accuracy percentage cannot answer all of that. Neither can a sensor count in a README.

For now, I think of this project as a place to work through those questions with real components and synthetic inputs. It has a concrete pipeline design, and it has gaps I can name: the training and serving contract, chronological evaluation, inference batching, and recovery behavior. Naming them gives me a much better next step than another round of hyperparameter tuning.

The thing I keep coming back to is that sensor reporting 25.5 degrees.

Before I ask someone to investigate it, I want to be able to explain which readings informed the decision, which model evaluated them, and whether the detector was working at the time. That is the kind of system I want to build: one whose answer I can follow all the way back to the observation.

You can explore the source on GitHub, or visit the project walkthrough for the architecture and setup details.

100kgs of Love,

Yassine X