Building a Real-Time Monitoring and Alerting Pipeline for 10,000 Sensors with DolphinDB
Picture a production line with 10,000 sensors, each reporting every 30 seconds. That's over 300 messages arriving per second, around the clock — and buried in that stream are the moments that actually matter: a temperature spike, a sudden current surge, a device silently going offline. Catch one of those the moment it happens, and you prevent downtime. Catch it in tomorrow's report, and the damage is already done.
That's the core challenge real-time stream processing is built to solve. Unlike a scheduled batch job that looks back at a finished dataset, a streaming pipeline never stops — data arrives, and cleaning, aggregation, evaluation, and alerting all have to happen immediately, in one continuous, low-latency chain.
This article walks through building a complete real-time monitoring system with DolphinDB's streaming engines — covering data quality handling, real-time alerting, and in-stream inference — and ties it together with an end-to-end example.
Real-Time Data Cleaning: Making Incoming Data Usable
Raw sensor data is rarely clean. Network jitter produces nulls, sensor drift produces outliers, and device restarts produce spurious jumps. If this "dirty" data flows straight into your alerting logic, you get false alarms at best and mis-triggered actions at worst. That's why the first checkpoint in any streaming pipeline is data cleaning.
Take voltage and current readings as an example: a production line reports voltage and current once per second. Under normal conditions, voltage should stay above 122V and current should never be null. Occasionally, though, sensor sampling glitches or electromagnetic interference cause voltage to report as 0 or current to go missing. That data needs to be filtered out before it reaches downstream computation.
Implementation
In DolphinDB, this is done by subscribing to a stream table and attaching a user-defined cleaning function — filtering happens inline as data arrives, and only clean records get persisted:
// user-defined handler: filters out records where voltage<=122 or current is NULL.
def append_after_filtering(inputTable, msg){
t = select * from msg where voltage>122, isValid(current)
if(size(t)>0){
insert into inputTable values(
t.timev,
t.voltage,
t.current)
}
}
electricityAggregator = createTimeSeriesEngine(
name = "electricityAggregator",
windowSize = 6,
step = 3,
metrics = < [avg(voltage), avg(current)] >,
dummyTable = electricity,
outputTable = outputTable,
timeColumn = `timev,
garbageSize = 2000
)
subscribeTable(
tableName = "electricity",
actionName = "avgElectricity",
offset = 0,
handler = append_after_filtering{electricityAggregator},
msgAsTable = true
)select * from electricity shows the raw incoming stream; select * from outputTable shows the aggregated results computed from the filtered data.
Raw incoming data
Data after real-time cleaning
The handler parameter is where the cleaning function plugs in — every record is checked against the filter conditions before it's allowed into the downstream aggregation pipeline.
This design decouples cleaning from business logic. The filtering rules can be modified and debugged independently without touching the downstream aggregation engine, and it's straightforward to add more rules later — rate-of-change checks, range validation, and so on.
Real-Time Temperature Alerting
Once data is clean, it flows into the real-time computation pipeline — and the core requirement here is alerting. Threshold breaches are the most common alert pattern on the factory floor: temperature above 80°C, pressure below a set point, current jumping more than 100% — all of these need to be evaluated the moment new data lands.
Take temperature monitoring as an example: if a device has crossed 40°C more than twice, and crossed 30°C more than three times, within a rolling 3-minute window, that triggers a temperature anomaly alert. The window slides every 30 seconds for continuous monitoring.
Implementation
DolphinDB's anomaly detection engine (createAnomalyDetectionEngine) handles this natively — grouping by device and evaluating complex conditions over sliding time windows:
// Anomaly detection engine for real-time temperature alerting
engine = createAnomalyDetectionEngine(
name = "engine1",
metrics = <[
sum(temperature > 40) > 2
&&
sum(temperature > 30) > 3
]>,
dummyTable = sensor,
outputTable = warningTable,
timeColumn = `ts,
keyColumn = `device_id,
windowSize = 180,
step = 30
)
subscribeTable(
tableName = "sensor",
actionName = "sensorAnomalyDetection",
offset = 0,
handler = append!{engine},
msgAsTable = true
)As data streams in, devices in an anomalous state show up in warningTable in real time.
Devices flagged with temperature anomalies
The key parameters here are windowSize=180 and step=30: the engine maintains a 180-second sliding window per device_id, re-evaluating the condition every 30 seconds. The metrics expression defines the alert logic — an alert only fires when both temperature conditions are satisfied at once.
This solves something batch processing structurally can't. If you only run a nightly script asking "did any device overheat yesterday?", the machine may have already tripped a breaker hours earlier. A streaming pipeline checks every 30 seconds instead, and the moment conditions are met, an alert fires — putting notification latency in the range of minutes, not a full day.
Machine Learning: Predicting Ahead of the Threshold
Threshold-based alerting only catches problems that have already crossed the line. Industrial settings often need something more forward-looking: predictive maintenance — flagging equipment that's likely to fail before it actually breaches a threshold. In DolphinDB, this is done by running machine learning models directly inside the streaming pipeline: data flows in continuously, the model keeps updating, and predictions come out in real time.
Take a wind turbine fleet as an example. Given current wind speed, humidity, air pressure, temperature, and equipment age, the system predicts whether current power output looks normal. If the predicted value deviates significantly from the actual reading, that's a sign the turbine may be in a degraded ("sub-healthy") state and worth inspecting early.
Approach
DolphinDB's streaming framework supports in-stream inference end to end:
- Ingestion — high-frequency sensor data (wind speed, humidity, pressure, temperature) flows into a stream table.
- Time-series aggregation — features are averaged per device every 10 seconds to reduce noise and data rate.
- Online training — each time a new aggregated data point arrives, a KNN regression model is retrained on the last 60 seconds of aggregated history (roughly 600 records).
- Real-time prediction — the current feature set is fed into the model to predict expected power output.
- Deviation detection — predicted vs. actual output is compared; a deviation greater than 0.1% writes a record to the alert table.
At a fleet of 100 turbines, this approach can sustain up to 100,000 raw records per second, with the model continuously retraining in-stream to keep pace with real-world conditions.
Compared with a traditional scheduled batch pipeline, the advantage of in-stream inference is that the model keeps adapting to the latest data — capturing equipment drift as it happens, rather than judging current conditions against a model trained on data from a week ago.
(The full implementation — KNN training, stream table subscriptions, time-series engine configuration, and more — is fairly extensive; Email us at info@dolphindb.com for the complete walkthrough.)
Closing Thoughts
The real shift here isn't any single engine or API — it's a change in what "knowing" means. Batch tells you what happened. Stream tells you what's happening. On a factory floor, that difference is the gap between a report you read tomorrow and a breaker that never trips tonight.
Want to see this pipeline running end to end — cleaning, alerting, and prediction wired together on live data? Try DolphinDB: https://dolphindb.com/.