π IoT Platform Project
Build a Scalable MongoDB-Powered IoT Data Management System
π― Project Overview
Welcome to the IoT Platform project! In this comprehensive hands-on tutorial, you'll build a production-ready Internet of Things data management system similar to AWS IoT, Azure IoT Hub, or Google Cloud IoT. You'll implement device registration, real-time sensor data collection, time-series analysis, alert systems, remote device control, and comprehensive monitoring dashboards.
This project introduces MongoDB's time-series collections for efficient IoT data storage, window functions for advanced analytics, TTL indexes for automatic data expiration, and aggregation pipelines for real-time monitoring. You'll learn patterns used by industrial IoT platforms managing millions of devices.
π Meet SmartFactory: Your IoT Management Platform
The Challenge: Sarah is deploying SmartFactory, an IoT platform for industrial monitoring. The system needs to manage thousands of sensors across multiple facilities, collect temperature, humidity, pressure, and vibration data in real-time, detect anomalies and trigger automated alerts, store historical data for trend analysis, enable remote device configuration and control, and generate insights from billions of sensor readings.
Your Mission: As the database architect, you'll design a MongoDB backend that supports:
- β Device registration and lifecycle management
- β High-frequency sensor data ingestion (1000s per second)
- β Real-time alerting based on threshold rules
- β Time-series data analysis and aggregation
- β Remote device commands and firmware updates
- β Historical data retention and archival
- β Dashboard analytics and reporting
π Learning Objectives
By completing this project, you will master:
π Time-Series Collections
Use MongoDB's specialized time-series collections for optimized IoT data storage
πͺ Window Functions
Calculate moving averages, running totals, and trend analysis
β±οΈ TTL Indexes
Automatically expire old data to manage storage costs
π Alert Systems
Build rule-based alerting for threshold violations
π Data Aggregation
Create complex pipelines for downsampling and analytics
β‘ High Throughput
Optimize for millions of writes per day
π Database Schema Design
Our IoT platform consists of 6 collections optimized for device and sensor data:
π± devices - Device Registry
Store device metadata, configuration, and status
- _id ObjectId
- deviceId String (unique, indexed)
- name String
- type String (sensor/actuator/gateway)
- manufacturer String
- model String
- firmwareVersion String
- location Object {facility, zone, coordinates}
- status String (online/offline/maintenance) (indexed)
- lastHeartbeat Date (indexed)
- configuration Object {samplingRate, thresholds}
- metadata Object (custom fields)
- registeredAt Date
- isActive Boolean
π‘ sensorReadings - Time-Series Data
High-frequency sensor measurements (Time-Series Collection)
- timestamp Date (timeField)
- deviceId String (metaField)
- temperature Number
- humidity Number
- pressure Number
- vibration Number
- voltage Number
- current Number
- metadata Object {facility, zone, sensorType}
π alerts - Alert Definitions & History
Alert rules and triggered alerts
- _id ObjectId
- alertId String (unique)
- deviceId String (indexed)
- ruleId ObjectId (reference)
- severity String (critical/warning/info)
- metric String (temperature/humidity/etc)
- value Number
- threshold Number
- message String
- status String (active/acknowledged/resolved)
- triggeredAt Date (indexed)
- acknowledgedAt Date
- resolvedAt Date
βοΈ alertRules - Alert Configuration
Define threshold-based alert rules
- _id ObjectId
- name String
- description String
- deviceFilter Object {type, facility, zone}
- metric String
- condition String (gt/lt/eq/gte/lte)
- threshold Number
- duration Number (seconds)
- severity String
- isEnabled Boolean (indexed)
- createdAt Date
π€ deviceCommands - Remote Control
Commands sent to devices for remote control
- _id ObjectId
- commandId String (unique)
- deviceId String (indexed)
- commandType String (config/reboot/update)
- payload Object
- status String (pending/sent/ack/failed)
- createdAt Date
- sentAt Date
- acknowledgedAt Date
- response Object
π deviceMetrics - Aggregated Statistics
Pre-computed hourly/daily aggregates for performance
- _id ObjectId
- deviceId String (indexed)
- date Date (indexed)
- interval String (hourly/daily)
- metrics Object {temperature, humidity, etc}
- statistics Object {min, max, avg, count}
- Time-Series Collection: sensorReadings uses time-series for optimal storage and query performance
- Separate Metrics: Pre-aggregated data for fast dashboard queries
- TTL on Readings: Auto-expire raw data after retention period (e.g., 90 days)
- Device Heartbeat: Track online/offline status with lastHeartbeat
- Command Queue: Async command delivery with status tracking
π§ Initial Setup
1Create the Database
use smartFactory // Verify creation db.getName() // Output: "smartFactory"
2Create Time-Series Collection
// Create time-series collection for sensor readings
db.createCollection("sensorReadings", {
timeseries: {
timeField: "timestamp",
metaField: "deviceId",
granularity: "seconds" // seconds, minutes, or hours
},
expireAfterSeconds: 7776000 // 90 days (90 * 24 * 60 * 60)
})
3Create Regular Collections
db.createCollection("devices")
db.createCollection("alerts")
db.createCollection("alertRules")
db.createCollection("deviceCommands")
db.createCollection("deviceMetrics")
4Create Indexes
// Devices
db.devices.createIndex({ deviceId: 1 }, { unique: true })
db.devices.createIndex({ status: 1 })
db.devices.createIndex({ lastHeartbeat: -1 })
db.devices.createIndex({ "location.facility": 1, "location.zone": 1 })
db.devices.createIndex({ type: 1 })
// Alerts
db.alerts.createIndex({ deviceId: 1, triggeredAt: -1 })
db.alerts.createIndex({ status: 1, severity: 1 })
db.alerts.createIndex({ triggeredAt: -1 })
db.alerts.createIndex({ alertId: 1 }, { unique: true })
// Alert Rules
db.alertRules.createIndex({ isEnabled: 1 })
db.alertRules.createIndex({ "deviceFilter.type": 1 })
// Device Commands
db.deviceCommands.createIndex({ commandId: 1 }, { unique: true })
db.deviceCommands.createIndex({ deviceId: 1, createdAt: -1 })
db.deviceCommands.createIndex({ status: 1 })
// Device Metrics
db.deviceMetrics.createIndex({ deviceId: 1, date: -1 })
db.deviceMetrics.createIndex({ date: -1, interval: 1 })
Your SmartFactory database is ready with time-series collections and all necessary indexes for high-performance IoT data management!
π± Phase 1: Device Management
Register and manage IoT devices across facilities.
Register Devices
db.devices.insertMany([
{
deviceId: "TEMP-SENSOR-001",
name: "Assembly Line Temperature Sensor",
type: "sensor",
manufacturer: "IndustrialSensors Inc",
model: "TS-4000",
firmwareVersion: "2.1.3",
location: {
facility: "Factory-A",
zone: "Assembly-Line-1",
coordinates: { lat: 37.7749, lng: -122.4194 }
},
status: "online",
lastHeartbeat: new Date(),
configuration: {
samplingRate: 5, // seconds
thresholds: {
temperature: { min: 15, max: 30 },
humidity: { min: 30, max: 70 }
}
},
metadata: {
installDate: new Date("2024-01-15"),
warrantyExpires: new Date("2026-01-15")
},
registeredAt: new Date(),
isActive: true
},
{
deviceId: "HUM-SENSOR-002",
name: "Warehouse Humidity Monitor",
type: "sensor",
manufacturer: "IndustrialSensors Inc",
model: "HM-3500",
firmwareVersion: "1.8.2",
location: {
facility: "Warehouse-B",
zone: "Storage-A",
coordinates: { lat: 37.7750, lng: -122.4195 }
},
status: "online",
lastHeartbeat: new Date(),
configuration: {
samplingRate: 10,
thresholds: {
humidity: { min: 40, max: 60 },
temperature: { min: 18, max: 25 }
}
},
metadata: {
installDate: new Date("2024-02-01"),
warrantyExpires: new Date("2026-02-01")
},
registeredAt: new Date(),
isActive: true
},
{
deviceId: "PRESS-SENSOR-003",
name: "Hydraulic Pressure Monitor",
type: "sensor",
manufacturer: "PressureTech",
model: "PT-9000",
firmwareVersion: "3.2.1",
location: {
facility: "Factory-A",
zone: "Hydraulics-Bay",
coordinates: { lat: 37.7751, lng: -122.4196 }
},
status: "online",
lastHeartbeat: new Date(),
configuration: {
samplingRate: 2,
thresholds: {
pressure: { min: 2000, max: 5000 } // PSI
}
},
metadata: {
installDate: new Date("2023-11-10"),
warrantyExpires: new Date("2025-11-10"),
calibrationDate: new Date("2024-11-01")
},
registeredAt: new Date(),
isActive: true
},
{
deviceId: "VIB-SENSOR-004",
name: "Motor Vibration Sensor",
type: "sensor",
manufacturer: "VibrationSystems",
model: "VS-2200",
firmwareVersion: "2.0.5",
location: {
facility: "Factory-A",
zone: "Motor-Room-1",
coordinates: { lat: 37.7749, lng: -122.4193 }
},
status: "online",
lastHeartbeat: new Date(),
configuration: {
samplingRate: 1, // High frequency for vibration
thresholds: {
vibration: { min: 0, max: 15 } // mm/s
}
},
metadata: {
installDate: new Date("2024-03-20"),
warrantyExpires: new Date("2026-03-20")
},
registeredAt: new Date(),
isActive: true
},
{
deviceId: "GATEWAY-001",
name: "Main Gateway - Factory A",
type: "gateway",
manufacturer: "IoTGateways Corp",
model: "GW-5000",
firmwareVersion: "4.1.0",
location: {
facility: "Factory-A",
zone: "Server-Room",
coordinates: { lat: 37.7750, lng: -122.4194 }
},
status: "online",
lastHeartbeat: new Date(),
configuration: {
connectedDevices: 25,
protocol: "MQTT"
},
metadata: {
installDate: new Date("2023-10-01"),
uptime: "99.98%"
},
registeredAt: new Date(),
isActive: true
}
])
Query Devices
// Get all online sensors
db.devices.find({
type: "sensor",
status: "online"
})
// Find devices by facility and zone
db.devices.find({
"location.facility": "Factory-A",
"location.zone": "Assembly-Line-1"
})
// Get devices that haven't sent heartbeat in 5 minutes
var fiveMinutesAgo = new Date(Date.now() - 5 * 60000)
db.devices.find({
lastHeartbeat: { $lt: fiveMinutesAgo },
isActive: true
})
Update Device Status
// Update heartbeat (device check-in)
db.devices.updateOne(
{ deviceId: "TEMP-SENSOR-001" },
{
$set: {
lastHeartbeat: new Date(),
status: "online"
}
}
)
// Mark device offline
db.devices.updateOne(
{ deviceId: "HUM-SENSOR-002" },
{
$set: {
status: "offline"
}
}
)
// Update firmware version
db.devices.updateOne(
{ deviceId: "TEMP-SENSOR-001" },
{
$set: {
firmwareVersion: "2.2.0",
"metadata.lastFirmwareUpdate": new Date()
}
}
)
// Update configuration
db.devices.updateOne(
{ deviceId: "TEMP-SENSOR-001" },
{
$set: {
"configuration.samplingRate": 3,
"configuration.thresholds.temperature.max": 35
}
}
)
Decommission Device
// Soft delete device
db.devices.updateOne(
{ deviceId: "TEMP-SENSOR-001" },
{
$set: {
isActive: false,
status: "decommissioned",
"metadata.decommissionedAt": new Date()
}
}
)
π‘ Phase 2: Sensor Data Collection
Ingest and query time-series sensor data at scale.
Insert Sensor Readings
// Single reading
db.sensorReadings.insertOne({
timestamp: new Date(),
deviceId: "TEMP-SENSOR-001",
temperature: 23.5,
humidity: 45.2,
metadata: {
facility: "Factory-A",
zone: "Assembly-Line-1",
sensorType: "temperature"
}
})
// Batch insert (more efficient for high throughput)
db.sensorReadings.insertMany([
{
timestamp: new Date(),
deviceId: "TEMP-SENSOR-001",
temperature: 24.1,
humidity: 46.3,
metadata: {
facility: "Factory-A",
zone: "Assembly-Line-1",
sensorType: "temperature"
}
},
{
timestamp: new Date(),
deviceId: "HUM-SENSOR-002",
temperature: 22.8,
humidity: 52.7,
metadata: {
facility: "Warehouse-B",
zone: "Storage-A",
sensorType: "humidity"
}
},
{
timestamp: new Date(),
deviceId: "PRESS-SENSOR-003",
pressure: 3245.5,
temperature: 28.3,
metadata: {
facility: "Factory-A",
zone: "Hydraulics-Bay",
sensorType: "pressure"
}
},
{
timestamp: new Date(),
deviceId: "VIB-SENSOR-004",
vibration: 8.2,
temperature: 35.6,
metadata: {
facility: "Factory-A",
zone: "Motor-Room-1",
sensorType: "vibration"
}
}
])
// Simulate continuous readings over time
for (let i = 0; i < 100; i++) {
let timestamp = new Date(Date.now() - (100 - i) * 5000) // 5 second intervals
db.sensorReadings.insertOne({
timestamp: timestamp,
deviceId: "TEMP-SENSOR-001",
temperature: 22 + Math.random() * 8, // Random between 22-30
humidity: 40 + Math.random() * 20, // Random between 40-60
metadata: {
facility: "Factory-A",
zone: "Assembly-Line-1",
sensorType: "temperature"
}
})
}
Query Recent Readings
// Get last 10 readings for a device
db.sensorReadings.find({
deviceId: "TEMP-SENSOR-001"
})
.sort({ timestamp: -1 })
.limit(10)
// Get readings from last hour
var oneHourAgo = new Date(Date.now() - 60 * 60000)
db.sensorReadings.find({
deviceId: "TEMP-SENSOR-001",
timestamp: { $gte: oneHourAgo }
})
.sort({ timestamp: -1 })
// Get readings in time range
db.sensorReadings.find({
deviceId: "TEMP-SENSOR-001",
timestamp: {
$gte: new Date("2024-12-22T10:00:00Z"),
$lte: new Date("2024-12-22T11:00:00Z")
}
})
Calculate Statistics
// Average temperature over last hour
var oneHourAgo = new Date(Date.now() - 60 * 60000)
db.sensorReadings.aggregate([
{
$match: {
deviceId: "TEMP-SENSOR-001",
timestamp: { $gte: oneHourAgo }
}
},
{
$group: {
_id: null,
avgTemp: { $avg: "$temperature" },
minTemp: { $min: "$temperature" },
maxTemp: { $max: "$temperature" },
avgHumidity: { $avg: "$humidity" },
count: { $sum: 1 }
}
}
])
// Hourly averages for the day
var startOfDay = new Date()
startOfDay.setHours(0, 0, 0, 0)
db.sensorReadings.aggregate([
{
$match: {
deviceId: "TEMP-SENSOR-001",
timestamp: { $gte: startOfDay }
}
},
{
$group: {
_id: {
year: { $year: "$timestamp" },
month: { $month: "$timestamp" },
day: { $dayOfMonth: "$timestamp" },
hour: { $hour: "$timestamp" }
},
avgTemp: { $avg: "$temperature" },
avgHumidity: { $avg: "$humidity" },
count: { $sum: 1 }
}
},
{
$project: {
_id: 0,
hour: "$_id.hour",
avgTemp: { $round: ["$avgTemp", 2] },
avgHumidity: { $round: ["$avgHumidity", 2] },
readingCount: "$count"
}
},
{ $sort: { hour: 1 } }
])
Detect Anomalies
// Find readings outside normal range
db.sensorReadings.find({
deviceId: "TEMP-SENSOR-001",
$or: [
{ temperature: { $gt: 30 } }, // Too hot
{ temperature: { $lt: 15 } }, // Too cold
{ humidity: { $gt: 70 } }, // Too humid
{ humidity: { $lt: 30 } } // Too dry
]
})
.sort({ timestamp: -1 })
.limit(20)
Moving Average (Window Function)
db.sensorReadings.aggregate([
{
$match: {
deviceId: "TEMP-SENSOR-001"
}
},
{ $sort: { timestamp: 1 } },
{
$setWindowFields: {
partitionBy: "$deviceId",
sortBy: { timestamp: 1 },
output: {
movingAvgTemp: {
$avg: "$temperature",
window: {
documents: [-2, 2] // 2 before, current, 2 after = 5 readings
}
},
movingAvgHumidity: {
$avg: "$humidity",
window: {
documents: [-2, 2]
}
}
}
}
},
{
$project: {
timestamp: 1,
temperature: { $round: ["$temperature", 2] },
humidity: { $round: ["$humidity", 2] },
movingAvgTemp: { $round: ["$movingAvgTemp", 2] },
movingAvgHumidity: { $round: ["$movingAvgHumidity", 2] }
}
},
{ $limit: 50 }
])
π Phase 3: Alerts & Monitoring Rules
Create automated alert rules and track triggered alerts.
Create Alert Rules
db.alertRules.insertMany([
{
name: "High Temperature Alert",
description: "Trigger when temperature exceeds 30Β°C",
deviceFilter: {
type: "sensor",
facility: "Factory-A"
},
metric: "temperature",
condition: "gt", // greater than
threshold: 30,
duration: 60, // Must persist for 60 seconds
severity: "warning",
isEnabled: true,
createdAt: new Date()
},
{
name: "Critical Temperature Alert",
description: "Critical alert when temperature exceeds 35Β°C",
deviceFilter: {
type: "sensor",
facility: "Factory-A"
},
metric: "temperature",
condition: "gt",
threshold: 35,
duration: 30,
severity: "critical",
isEnabled: true,
createdAt: new Date()
},
{
name: "Low Humidity Alert",
description: "Alert when humidity drops below 30%",
deviceFilter: {
type: "sensor"
},
metric: "humidity",
condition: "lt", // less than
threshold: 30,
duration: 120,
severity: "info",
isEnabled: true,
createdAt: new Date()
},
{
name: "High Pressure Alert",
description: "Pressure exceeds safe operating threshold",
deviceFilter: {
type: "sensor",
zone: "Hydraulics-Bay"
},
metric: "pressure",
condition: "gt",
threshold: 5000,
duration: 10,
severity: "critical",
isEnabled: true,
createdAt: new Date()
},
{
name: "Excessive Vibration Alert",
description: "Motor vibration indicates potential failure",
deviceFilter: {
type: "sensor",
zone: "Motor-Room-1"
},
metric: "vibration",
condition: "gt",
threshold: 15,
duration: 5,
severity: "critical",
isEnabled: true,
createdAt: new Date()
}
])
Trigger Alerts
// Simulate alert trigger
var rule = db.alertRules.findOne({ name: "High Temperature Alert" })
db.alerts.insertOne({
alertId: "ALERT-" + new Date().getTime(),
deviceId: "TEMP-SENSOR-001",
ruleId: rule._id,
severity: rule.severity,
metric: rule.metric,
value: 32.5,
threshold: rule.threshold,
message: "Temperature 32.5Β°C exceeds threshold of 30Β°C",
status: "active",
triggeredAt: new Date(),
acknowledgedAt: null,
resolvedAt: null
})
// Critical pressure alert
var pressureRule = db.alertRules.findOne({ name: "High Pressure Alert" })
db.alerts.insertOne({
alertId: "ALERT-" + new Date().getTime(),
deviceId: "PRESS-SENSOR-003",
ruleId: pressureRule._id,
severity: "critical",
metric: "pressure",
value: 5250,
threshold: 5000,
message: "CRITICAL: Pressure 5250 PSI exceeds safe limit of 5000 PSI",
status: "active",
triggeredAt: new Date(),
acknowledgedAt: null,
resolvedAt: null
})
Query Active Alerts
// Get all active alerts
db.alerts.find({
status: "active"
})
.sort({ triggeredAt: -1 })
// Get critical alerts only
db.alerts.find({
status: "active",
severity: "critical"
})
.sort({ triggeredAt: -1 })
// Get alerts for specific device
db.alerts.find({
deviceId: "TEMP-SENSOR-001"
})
.sort({ triggeredAt: -1 })
// Get unacknowledged alerts
db.alerts.find({
status: "active",
acknowledgedAt: null
})
.sort({ severity: 1, triggeredAt: -1 })
Acknowledge and Resolve Alerts
// Acknowledge alert
db.alerts.updateOne(
{ alertId: "ALERT-1734876543210" },
{
$set: {
status: "acknowledged",
acknowledgedAt: new Date(),
acknowledgedBy: "operator@smartfactory.com"
}
}
)
// Resolve alert
db.alerts.updateOne(
{ alertId: "ALERT-1734876543210" },
{
$set: {
status: "resolved",
resolvedAt: new Date(),
resolvedBy: "operator@smartfactory.com",
resolution: "Temperature normalized after HVAC adjustment"
}
}
)
Alert Analytics
// Alerts by severity
db.alerts.aggregate([
{
$group: {
_id: "$severity",
count: { $sum: 1 },
activeCount: {
$sum: { $cond: [{ $eq: ["$status", "active"] }, 1, 0] }
}
}
},
{ $sort: { count: -1 } }
])
// Most problematic devices
db.alerts.aggregate([
{
$group: {
_id: "$deviceId",
alertCount: { $sum: 1 },
criticalCount: {
$sum: { $cond: [{ $eq: ["$severity", "critical"] }, 1, 0] }
}
}
},
{
$lookup: {
from: "devices",
localField: "_id",
foreignField: "deviceId",
as: "device"
}
},
{ $unwind: "$device" },
{
$project: {
_id: 0,
deviceId: "$_id",
deviceName: "$device.name",
alertCount: 1,
criticalCount: 1
}
},
{ $sort: { alertCount: -1 } },
{ $limit: 10 }
])
π Phase 4: Time-Series Analytics
Advanced analytics on sensor data using aggregation pipelines.
Hourly Aggregations
// Calculate hourly averages for last 24 hours
var oneDayAgo = new Date(Date.now() - 24 * 60 * 60000)
db.sensorReadings.aggregate([
{
$match: {
deviceId: "TEMP-SENSOR-001",
timestamp: { $gte: oneDayAgo }
}
},
{
$group: {
_id: {
year: { $year: "$timestamp" },
month: { $month: "$timestamp" },
day: { $dayOfMonth: "$timestamp" },
hour: { $hour: "$timestamp" }
},
avgTemperature: { $avg: "$temperature" },
minTemperature: { $min: "$temperature" },
maxTemperature: { $max: "$temperature" },
avgHumidity: { $avg: "$humidity" },
readingCount: { $sum: 1 }
}
},
{
$project: {
_id: 0,
timestamp: {
$dateFromParts: {
year: "$_id.year",
month: "$_id.month",
day: "$_id.day",
hour: "$_id.hour"
}
},
avgTemperature: { $round: ["$avgTemperature", 2] },
minTemperature: { $round: ["$minTemperature", 2] },
maxTemperature: { $round: ["$maxTemperature", 2] },
avgHumidity: { $round: ["$avgHumidity", 2] },
readingCount: 1
}
},
{ $sort: { timestamp: 1 } }
])
Store Aggregated Metrics
// Save hourly metrics
var startOfHour = new Date()
startOfHour.setMinutes(0, 0, 0)
var endOfHour = new Date(startOfHour.getTime() + 60 * 60000)
var hourlyStats = db.sensorReadings.aggregate([
{
$match: {
deviceId: "TEMP-SENSOR-001",
timestamp: {
$gte: startOfHour,
$lt: endOfHour
}
}
},
{
$group: {
_id: null,
avgTemp: { $avg: "$temperature" },
minTemp: { $min: "$temperature" },
maxTemp: { $max: "$temperature" },
avgHumidity: { $avg: "$humidity" },
count: { $sum: 1 }
}
}
]).toArray()[0]
if (hourlyStats) {
db.deviceMetrics.insertOne({
deviceId: "TEMP-SENSOR-001",
date: startOfHour,
interval: "hourly",
metrics: {
temperature: hourlyStats.avgTemp,
humidity: hourlyStats.avgHumidity
},
statistics: {
min: hourlyStats.minTemp,
max: hourlyStats.maxTemp,
avg: hourlyStats.avgTemp,
count: hourlyStats.count
}
})
}
Multi-Device Comparison
// Compare all sensors in a facility
var oneHourAgo = new Date(Date.now() - 60 * 60000)
db.sensorReadings.aggregate([
{
$match: {
"metadata.facility": "Factory-A",
timestamp: { $gte: oneHourAgo }
}
},
{
$group: {
_id: "$deviceId",
avgTemp: { $avg: "$temperature" },
avgHumidity: { $avg: "$humidity" },
avgPressure: { $avg: "$pressure" },
avgVibration: { $avg: "$vibration" },
readingCount: { $sum: 1 }
}
},
{
$lookup: {
from: "devices",
localField: "_id",
foreignField: "deviceId",
as: "device"
}
},
{ $unwind: "$device" },
{
$project: {
_id: 0,
deviceId: "$_id",
deviceName: "$device.name",
zone: "$device.location.zone",
avgTemp: { $round: ["$avgTemp", 2] },
avgHumidity: { $round: ["$avgHumidity", 2] },
avgPressure: { $round: ["$avgPressure", 2] },
avgVibration: { $round: ["$avgVibration", 2] },
readingCount: 1
}
},
{ $sort: { avgTemp: -1 } }
])
Trend Analysis
// Temperature trend using window functions
db.sensorReadings.aggregate([
{
$match: {
deviceId: "TEMP-SENSOR-001"
}
},
{ $sort: { timestamp: 1 } },
{ $limit: 100 },
{
$setWindowFields: {
partitionBy: "$deviceId",
sortBy: { timestamp: 1 },
output: {
previousTemp: {
$shift: {
output: "$temperature",
by: -1
}
},
previousTime: {
$shift: {
output: "$timestamp",
by: -1
}
}
}
}
},
{
$project: {
timestamp: 1,
temperature: 1,
tempChange: {
$cond: [
{ $ne: ["$previousTemp", null] },
{ $subtract: ["$temperature", "$previousTemp"] },
0
]
},
timeElapsed: {
$cond: [
{ $ne: ["$previousTime", null] },
{
$divide: [
{ $subtract: ["$timestamp", "$previousTime"] },
1000 // Convert to seconds
]
},
0
]
}
}
},
{
$project: {
timestamp: 1,
temperature: { $round: ["$temperature", 2] },
tempChange: { $round: ["$tempChange", 2] },
rateOfChange: {
$cond: [
{ $gt: ["$timeElapsed", 0] },
{ $round: [{ $divide: ["$tempChange", "$timeElapsed"] }, 4] },
0
]
}
}
}
])
Percentile Analysis
// Calculate temperature percentiles
db.sensorReadings.aggregate([
{
$match: {
deviceId: "TEMP-SENSOR-001",
timestamp: { $gte: new Date(Date.now() - 24 * 60 * 60000) }
}
},
{
$group: {
_id: "$deviceId",
temps: { $push: "$temperature" }
}
},
{
$project: {
_id: 0,
deviceId: "$_id",
p50: { $percentile: { input: "$temps", p: [0.5], method: "approximate" } },
p75: { $percentile: { input: "$temps", p: [0.75], method: "approximate" } },
p90: { $percentile: { input: "$temps", p: [0.90], method: "approximate" } },
p95: { $percentile: { input: "$temps", p: [0.95], method: "approximate" } },
p99: { $percentile: { input: "$temps", p: [0.99], method: "approximate" } }
}
}
])
π€ Phase 5: Device Commands & Control
Send remote commands to devices and track execution.
Send Configuration Command
db.deviceCommands.insertOne({
commandId: "CMD-" + new Date().getTime(),
deviceId: "TEMP-SENSOR-001",
commandType: "config",
payload: {
samplingRate: 3, // Change from 5 to 3 seconds
thresholds: {
temperature: { min: 15, max: 35 }
}
},
status: "pending",
createdAt: new Date(),
sentAt: null,
acknowledgedAt: null,
response: null
})
Send Reboot Command
db.deviceCommands.insertOne({
commandId: "CMD-" + new Date().getTime(),
deviceId: "GATEWAY-001",
commandType: "reboot",
payload: {
reason: "Scheduled maintenance",
delay: 60 // Wait 60 seconds before reboot
},
status: "pending",
createdAt: new Date(),
sentAt: null,
acknowledgedAt: null,
response: null
})
Send Firmware Update
db.deviceCommands.insertOne({
commandId: "CMD-" + new Date().getTime(),
deviceId: "TEMP-SENSOR-001",
commandType: "update",
payload: {
firmwareVersion: "2.3.0",
downloadUrl: "https://updates.smartfactory.com/firmware/ts-4000-v2.3.0.bin",
checksum: "a3b5c7d9e1f2..."
},
status: "pending",
createdAt: new Date(),
sentAt: null,
acknowledgedAt: null,
response: null
})
Update Command Status
var cmdId = "CMD-1734876543210"
// Mark as sent
db.deviceCommands.updateOne(
{ commandId: cmdId },
{
$set: {
status: "sent",
sentAt: new Date()
}
}
)
// Device acknowledges command
db.deviceCommands.updateOne(
{ commandId: cmdId },
{
$set: {
status: "acknowledged",
acknowledgedAt: new Date(),
response: {
message: "Configuration updated successfully",
appliedAt: new Date()
}
}
}
)
// Command failed
db.deviceCommands.updateOne(
{ commandId: cmdId },
{
$set: {
status: "failed",
response: {
error: "Invalid parameter: samplingRate must be >= 1",
failedAt: new Date()
}
}
}
)
Query Commands
// Get pending commands for device
db.deviceCommands.find({
deviceId: "TEMP-SENSOR-001",
status: "pending"
})
.sort({ createdAt: 1 })
// Get command history
db.deviceCommands.find({
deviceId: "TEMP-SENSOR-001"
})
.sort({ createdAt: -1 })
.limit(20)
// Failed commands
db.deviceCommands.find({
status: "failed"
})
.sort({ createdAt: -1 })
π Phase 6: Dashboard Analytics
Generate comprehensive insights for monitoring dashboards.
System Overview
// Device status summary
db.devices.aggregate([
{
$group: {
_id: "$status",
count: { $sum: 1 }
}
},
{
$group: {
_id: null,
statusCounts: {
$push: { status: "$_id", count: "$count" }
},
totalDevices: { $sum: "$count" }
}
}
])
// Devices by type and status
db.devices.aggregate([
{
$group: {
_id: { type: "$type", status: "$status" },
count: { $sum: 1 }
}
},
{
$project: {
_id: 0,
type: "$_id.type",
status: "$_id.status",
count: 1
}
},
{ $sort: { type: 1, status: 1 } }
])
Facility Performance
// Average conditions by facility
var oneHourAgo = new Date(Date.now() - 60 * 60000)
db.sensorReadings.aggregate([
{
$match: {
timestamp: { $gte: oneHourAgo }
}
},
{
$group: {
_id: "$metadata.facility",
avgTemp: { $avg: "$temperature" },
avgHumidity: { $avg: "$humidity" },
avgPressure: { $avg: "$pressure" },
readingCount: { $sum: 1 },
deviceCount: { $addToSet: "$deviceId" }
}
},
{
$project: {
_id: 0,
facility: "$_id",
avgTemp: { $round: ["$avgTemp", 2] },
avgHumidity: { $round: ["$avgHumidity", 2] },
avgPressure: { $round: ["$avgPressure", 2] },
readingCount: 1,
activeDevices: { $size: "$deviceCount" }
}
},
{ $sort: { facility: 1 } }
])
Data Quality Metrics
// Calculate expected vs actual readings
db.devices.aggregate([
{
$match: {
type: "sensor",
isActive: true
}
},
{
$lookup: {
from: "sensorReadings",
let: { devId: "$deviceId", rate: "$configuration.samplingRate" },
pipeline: [
{
$match: {
$expr: {
$and: [
{ $eq: ["$deviceId", "$$devId"] },
{ $gte: ["$timestamp", new Date(Date.now() - 60 * 60000)] }
]
}
}
},
{
$group: {
_id: null,
actualReadings: { $sum: 1 }
}
}
],
as: "readings"
}
},
{
$project: {
deviceId: 1,
name: 1,
samplingRate: "$configuration.samplingRate",
expectedReadings: {
$divide: [3600, "$configuration.samplingRate"] // 1 hour / rate
},
actualReadings: {
$ifNull: [{ $arrayElemAt: ["$readings.actualReadings", 0] }, 0]
}
}
},
{
$project: {
deviceId: 1,
name: 1,
expectedReadings: { $round: ["$expectedReadings", 0] },
actualReadings: 1,
dataQuality: {
$round: [
{
$multiply: [
{ $divide: ["$actualReadings", "$expectedReadings"] },
100
]
},
2
]
}
}
},
{ $sort: { dataQuality: 1 } }
])
Alert Response Times
// Average time to acknowledge and resolve alerts
db.alerts.aggregate([
{
$match: {
acknowledgedAt: { $ne: null }
}
},
{
$project: {
severity: 1,
timeToAcknowledge: {
$divide: [
{ $subtract: ["$acknowledgedAt", "$triggeredAt"] },
60000 // Convert to minutes
]
},
timeToResolve: {
$cond: [
{ $ne: ["$resolvedAt", null] },
{
$divide: [
{ $subtract: ["$resolvedAt", "$triggeredAt"] },
60000
]
},
null
]
}
}
},
{
$group: {
_id: "$severity",
avgAckTime: { $avg: "$timeToAcknowledge" },
avgResolveTime: { $avg: "$timeToResolve" },
count: { $sum: 1 }
}
},
{
$project: {
_id: 0,
severity: "$_id",
avgAckTimeMinutes: { $round: ["$avgAckTime", 2] },
avgResolveTimeMinutes: { $round: ["$avgResolveTime", 2] },
alertCount: "$count"
}
},
{ $sort: { severity: 1 } }
])
Device Uptime Analysis
// Calculate uptime percentage based on heartbeat
var last7Days = new Date(Date.now() - 7 * 24 * 60 * 60000)
db.devices.aggregate([
{
$match: {
type: "sensor",
registeredAt: { $lte: last7Days }
}
},
{
$project: {
deviceId: 1,
name: 1,
daysSinceRegistration: {
$divide: [
{ $subtract: [new Date(), "$registeredAt"] },
86400000 // milliseconds in a day
]
},
hoursSinceHeartbeat: {
$divide: [
{ $subtract: [new Date(), "$lastHeartbeat"] },
3600000 // milliseconds in an hour
]
},
isCurrentlyOnline: { $eq: ["$status", "online"] }
}
},
{
$project: {
deviceId: 1,
name: 1,
isCurrentlyOnline: 1,
estimatedUptime: {
$cond: [
{ $lt: ["$hoursSinceHeartbeat", 1] },
99.9, // Recently checked in
{
$multiply: [
{
$divide: [
{
$subtract: [
"$daysSinceRegistration",
{ $divide: ["$hoursSinceHeartbeat", 24] }
]
},
"$daysSinceRegistration"
]
},
100
]
}
]
}
}
},
{
$project: {
deviceId: 1,
name: 1,
isCurrentlyOnline: 1,
uptimePercentage: { $round: ["$estimatedUptime", 2] }
}
},
{ $sort: { uptimePercentage: 1 } }
])
Energy Consumption Estimate
// Devices with voltage and current readings
var oneDayAgo = new Date(Date.now() - 24 * 60 * 60000)
db.sensorReadings.aggregate([
{
$match: {
voltage: { $exists: true },
current: { $exists: true },
timestamp: { $gte: oneDayAgo }
}
},
{
$project: {
deviceId: 1,
timestamp: 1,
powerWatts: { $multiply: ["$voltage", "$current"] }
}
},
{
$group: {
_id: "$deviceId",
avgPowerWatts: { $avg: "$powerWatts" },
maxPowerWatts: { $max: "$powerWatts" }
}
},
{
$project: {
_id: 0,
deviceId: "$_id",
avgPowerWatts: { $round: ["$avgPowerWatts", 2] },
maxPowerWatts: { $round: ["$maxPowerWatts", 2] },
dailyEnergyKWh: {
$round: [
{ $multiply: ["$avgPowerWatts", 24, 0.001] }, // 24 hours, convert to kWh
2
]
}
}
},
{ $sort: { dailyEnergyKWh: -1 } }
])
β‘ Performance Optimization
π Time-Series Collections
Optimized storage and compression for high-frequency sensor data
β±οΈ TTL Indexes
Automatic expiration of old data prevents unbounded growth
π Compound Indexes
{deviceId: 1, timestamp: -1} for efficient time-range queries
π¦ Pre-Aggregation
Store hourly/daily metrics for fast dashboard loading
Batch Insert Optimization
// Use ordered: false for parallel inserts
var readings = []
for (let i = 0; i < 1000; i++) {
readings.push({
timestamp: new Date(Date.now() - i * 5000),
deviceId: "TEMP-SENSOR-001",
temperature: 22 + Math.random() * 8,
humidity: 40 + Math.random() * 20,
metadata: {
facility: "Factory-A",
zone: "Assembly-Line-1",
sensorType: "temperature"
}
})
}
db.sensorReadings.insertMany(readings, { ordered: false })
Query Performance Analysis
// Check query execution plan
db.sensorReadings.find({
deviceId: "TEMP-SENSOR-001",
timestamp: { $gte: new Date(Date.now() - 60 * 60000) }
})
.sort({ timestamp: -1 })
.explain("executionStats")
// Verify time-series collection benefits
db.sensorReadings.stats()
- Use time-series collections for all sensor data - significant storage savings
- Set appropriate TTL for data retention (90 days for raw data)
- Pre-aggregate data into hourly/daily summaries for dashboards
- Use bulk inserts (insertMany) instead of individual inserts
- Archive historical data to cold storage after 6-12 months
- Monitor collection size and adjust granularity if needed
- Use projection to limit fields returned in queries
Data Retention Strategy
// Raw data: 90 days (auto-expires via TTL)
// Already set: expireAfterSeconds: 7776000
// Hourly aggregates: 1 year
db.deviceMetrics.createIndex(
{ date: 1 },
{ expireAfterSeconds: 31536000 } // 365 days
)
// Daily aggregates: Keep forever (small size)
// No TTL - permanent retention
π Advanced Challenges
Challenge 1: Predictive Maintenance
Use vibration trends to predict equipment failure 24 hours in advance
Challenge 2: Anomaly Detection
Implement statistical outlier detection using standard deviation
Challenge 3: Energy Optimization
Identify devices consuming more power than expected based on usage patterns
Challenge 4: Geofencing
Alert when mobile sensors leave designated facility boundaries
Challenge 5: Data Streaming
Use Change Streams to trigger real-time alerts on data ingestion
Challenge 6: Custom Metrics
Build a flexible metrics engine supporting user-defined formulas
- Predictive Maintenance: Use $setWindowFields to calculate moving average and standard deviation of vibration, flag when current reading exceeds avg + 2*stddev
- Anomaly Detection: Aggregate last 1000 readings, calculate mean/stddev, identify readings outside 3 standard deviations
- Energy Optimization: Compare actual power consumption vs baseline, use $lookup to join device metadata with consumption data
- Geofencing: Add 2dsphere index on device location, use $geoWithin to check if coordinates are inside facility polygon
- Data Streaming: Create change stream on sensorReadings, evaluate alert rules in real-time using watch()
- Custom Metrics: Store formula as string in configuration, use $function aggregation operator to evaluate JavaScript expressions