Section 11: Projects

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

πŸ’‘ What Makes This Project Unique?

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.

6
Collections
30+
Features
80+
Queries
Time-Series
Optimized

πŸ“– 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}
🎯 Key Design Decisions:
  • 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

Initialize SmartFactory Database
use smartFactory

// Verify creation
db.getName()
// Output: "smartFactory"

2Create Time-Series Collection

Critical: Time-series collection for sensor data
// 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 })
βœ… Setup Complete!

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

Add IoT Devices to Registry
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

High-Frequency Data Ingestion
// 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

Aggregation on Time-Series Data
// 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)

Calculate 5-Reading Moving Average
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

Define Threshold-Based 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

Create Alert When Rule Violated
// 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

Alert Summary and Trends
// 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

Downsample Data to Hourly Averages
// 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

Pre-compute for Dashboard Performance
// 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

Calculate Rate of Change
// 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

Update Device Settings Remotely
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

Fleet Health Summary
// 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

Calculate Power Usage from Voltage/Current
// 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

High Throughput Data Ingestion
// 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()
⚠️ Performance Best Practices:
  • 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

Multi-Tier Data Storage
// 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

πŸ’‘ Challenge Hints:
  • 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