Real-World Project

IoT Sensor Data Platform

Build a production-ready time-series data system with Cassandra!

📖 Project Overview

Build a real-world IoT platform!

The Challenge: Smart Factory IoT

SmartFactory Inc. operates 500 factories worldwide. Each factory has 1,000 IoT sensors monitoring:

  • 🌡️ Temperature sensors: Every 10 seconds
  • ⚡ Power consumption: Every 5 seconds
  • 🏭 Production metrics: Every 30 seconds
  • 🔔 Alert sensors: Real-time events

Scale Requirements

Total Sensors: 500,000 (500 factories × 1,000 sensors)
Data Points/Second: ~100,000 writes/second
Data Points/Day: 8.64 billion data points
Retention: 90 days (detailed), 2 years (aggregated)
Query Pattern: Recent data (last 24 hours most common)

Why Cassandra?

  • ✅ High write throughput: 100K+ writes/second
  • ✅ Time-series optimized: Perfect for sensor data
  • ✅ Linear scalability: Add nodes as needed
  • ✅ No single point of failure: Always available
  • ✅ TTL support: Automatic data expiration

🎯 Requirements & Use Cases

What we need to build!

📝

Write Operations

  • Insert sensor reading (100K/sec)
  • Batch insert multiple readings
  • Insert with TTL (90 days)
  • Handle late-arriving data
  • Write alerts/anomalies
🔍

Read Operations

  • Get latest reading for sensor
  • Get readings for time range
  • Get all sensors in factory
  • Get hourly/daily aggregates
  • Real-time dashboard queries
📊

Analytics Queries

  • Average temperature per hour
  • Peak power consumption
  • Anomaly detection
  • Trend analysis
  • Comparative reports
🚨

Operational

  • Auto-expire old data (TTL)
  • Handle node failures
  • Scale horizontally
  • Monitor performance
  • Backup critical data

🗺️ Data Model Design

Cassandra data modeling!

Cassandra Data Modeling Rules

  1. Know your queries first! Model data around access patterns
  2. Denormalize: Duplicate data for query efficiency
  3. Partition key: Distributes data across nodes
  4. Clustering key: Sorts data within partition
  5. One partition per query: Avoid multi-partition reads

Query-Driven Design

Query Pattern Access Pattern Table Needed
Q1: Latest reading for sensor By sensor_id, recent first sensor_data_by_sensor
Q2: Readings by time range By sensor_id + time sensor_data_by_sensor
Q3: All sensors in factory By factory_id sensors_by_factory
Q4: Hourly aggregates By sensor_id + hour sensor_hourly_stats
Q5: Recent alerts By factory_id + time alerts_by_factory

📋 Schema Design

Complete table schemas!

1

Keyspace Creation

-- Create keyspace CREATE KEYSPACE iot_sensors WITH replication = { 'class': 'NetworkTopologyStrategy', 'datacenter1': 3, -- 3 replicas 'datacenter2': 3 -- Multi-DC for HA }; USE iot_sensors;
2

Main Sensor Data Table

-- Raw sensor readings (high write throughput) CREATE TABLE sensor_data_by_sensor ( sensor_id UUID, -- Partition key reading_time TIMESTAMP, -- Clustering key (DESC) value DOUBLE, unit TEXT, quality TEXT, -- 'good', 'bad', 'uncertain' PRIMARY KEY (sensor_id, reading_time) ) WITH CLUSTERING ORDER BY ) WITH CLUSTERING ORDER BY (reading_time DESC) AND default_time_to_live = 7776000 -- 90 days in seconds AND compaction = { 'class': 'TimeWindowCompactionStrategy', 'compaction_window_unit': 'DAYS', 'compaction_window_size': 1 }; -- Why this design? -- ✅ Partition key: sensor_id (distributes data evenly) -- ✅ Clustering: reading_time DESC (recent data first) -- ✅ TTL: Auto-delete after 90 days -- ✅ TWCS: Optimized for time-series data
3

Sensors Registry Table

-- Sensor metadata and lookup CREATE TABLE sensors_by_factory ( factory_id UUID, -- Partition key sensor_id UUID, -- Clustering key sensor_type TEXT, -- 'temperature', 'power', etc location TEXT, installation_date TIMESTAMP, status TEXT, -- 'active', 'maintenance', 'offline' PRIMARY KEY (factory_id, sensor_id) ); -- Query: Get all sensors in factory SELECT * FROM sensors_by_factory WHERE factory_id = ?;
4

Hourly Aggregates Table

-- Pre-aggregated hourly statistics CREATE TABLE sensor_hourly_stats ( sensor_id UUID, hour_bucket TIMESTAMP, -- Rounded to hour avg_value DOUBLE, min_value DOUBLE, max_value DOUBLE, count COUNTER, PRIMARY KEY (sensor_id, hour_bucket) ) WITH CLUSTERING ORDER BY (hour_bucket DESC) AND default_time_to_live = 63072000; -- 2 years -- Query: Get hourly stats for date range SELECT * FROM sensor_hourly_stats WHERE sensor_id = ? AND hour_bucket >= ? AND hour_bucket <= ?;
5

Alerts Table

-- Critical alerts and anomalies CREATE TABLE alerts_by_factory ( factory_id UUID, alert_time TIMESTAMP, alert_id UUID, sensor_id UUID, alert_type TEXT, -- 'temperature_high', 'power_spike' severity TEXT, -- 'low', 'medium', 'high', 'critical' value DOUBLE, threshold DOUBLE, message TEXT, acknowledged BOOLEAN, PRIMARY KEY (factory_id, alert_time, alert_id) ) WITH CLUSTERING ORDER BY (alert_time DESC, alert_id ASC) AND default_time_to_live = 7776000; -- 90 days -- Query: Get recent alerts for factory SELECT * FROM alerts_by_factory WHERE factory_id = ? AND alert_time >= ? LIMIT 100;

💻 Implementation

Python implementation!

Data Ingestion Service

# sensor_ingestion.py from cassandra.cluster import Cluster from cassandra.query import BatchStatement from datetime import datetime import uuid class SensorDataIngestion: def __init__(self, contact_points): cluster = Cluster(contact_points) self.session = cluster.connect('iot_sensors') # Prepare statements for performance self.insert_reading = self.session.prepare(""" INSERT INTO sensor_data_by_sensor (sensor_id, reading_time, value, unit, quality) VALUES (?, ?, ?, ?, ?) """) self.insert_alert = self.session.prepare(""" INSERT INTO alerts_by_factory (factory_id, alert_time, alert_id, sensor_id, alert_type, severity, value, threshold, message, acknowledged) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """) def insert_sensor_reading(self, sensor_id, value, unit, quality='good'): """Insert single sensor reading""" reading_time = datetime.now() self.session.execute( self.insert_reading, (sensor_id, reading_time, value, unit, quality) ) return reading_time def insert_batch_readings(self, readings): """Batch insert for high throughput""" batch = BatchStatement() for reading in readings: batch.add( self.insert_reading, ( reading['sensor_id'], reading.get('reading_time', datetime.now()), reading['value'], reading['unit'], reading.get('quality', 'good') ) ) self.session.execute(batch) def create_alert(self, factory_id, sensor_id, alert_type, value, threshold, severity='medium'): """Create alert for anomaly""" alert_id = uuid.uuid4() alert_time = datetime.now() message = f"{alert_type}: {value} exceeds {threshold}" self.session.execute( self.insert_alert, (factory_id, alert_time, alert_id, sensor_id, alert_type, severity, value, threshold, message, False) ) return alert_id # Usage example if __name__ == '__main__': ingestion = SensorDataIngestion(['localhost']) # Insert single reading sensor_id = uuid.UUID('123e4567-e89b-12d3-a456-426614174000') ingestion.insert_sensor_reading( sensor_id=sensor_id, value=23.5, unit='celsius' ) # Batch insert (efficient!) readings = [ {'sensor_id': sensor_id, 'value': 23.6, 'unit': 'celsius'}, {'sensor_id': sensor_id, 'value': 23.7, 'unit': 'celsius'}, {'sensor_id': sensor_id, 'value': 23.8, 'unit': 'celsius'}, ] ingestion.insert_batch_readings(readings)

Query Service

# sensor_queries.py from cassandra.cluster import Cluster from datetime import datetime, timedelta class SensorQueries: def __init__(self, contact_points): cluster = Cluster(contact_points) self.session = cluster.connect('iot_sensors') def get_latest_reading(self, sensor_id): """Get most recent reading for sensor""" query = """ SELECT * FROM sensor_data_by_sensor WHERE sensor_id = ? LIMIT 1 """ row = self.session.execute(query, (sensor_id,)).one() if row: return { 'sensor_id': row.sensor_id, 'reading_time': row.reading_time, 'value': row.value, 'unit': row.unit, 'quality': row.quality } return None def get_readings_range(self, sensor_id, start_time, end_time): """Get readings for time range""" query = """ SELECT * FROM sensor_data_by_sensor WHERE sensor_id = ? AND reading_time >= ? AND reading_time <= ? """ rows = self.session.execute(query, (sensor_id, start_time, end_time)) return [ { 'reading_time': row.reading_time, 'value': row.value, 'unit': row.unit, 'quality': row.quality } for row in rows ] def get_last_24_hours(self, sensor_id): """Get readings from last 24 hours""" end_time = datetime.now() start_time = end_time - timedelta(hours=24) return self.get_readings_range(sensor_id, start_time, end_time) def get_factory_sensors(self, factory_id): """Get all sensors in factory""" query = """ SELECT * FROM sensors_by_factory WHERE factory_id = ? """ rows = self.session.execute(query, (factory_id,)) return [ { 'sensor_id': row.sensor_id, 'sensor_type': row.sensor_type, 'location': row.location, 'status': row.status } for row in rows ] def get_recent_alerts(self, factory_id, hours=24): """Get recent alerts for factory""" start_time = datetime.now() - timedelta(hours=hours) query = """ SELECT * FROM alerts_by_factory WHERE factory_id = ? AND alert_time >= ? """ rows = self.session.execute(query, (factory_id, start_time)) return [ { 'alert_id': row.alert_id, 'alert_time': row.alert_time, 'sensor_id': row.sensor_id, 'alert_type': row.alert_type, 'severity': row.severity, 'value': row.value, 'message': row.message, 'acknowledged': row.acknowledged } for row in rows ] def get_hourly_stats(self, sensor_id, start_hour, end_hour): """Get pre-aggregated hourly statistics""" query = """ SELECT * FROM sensor_hourly_stats WHERE sensor_id = ? AND hour_bucket >= ? AND hour_bucket <= ? """ rows = self.session.execute(query, (sensor_id, start_hour, end_hour)) return [ { 'hour': row.hour_bucket, 'avg': row.avg_value, 'min': row.min_value, 'max': row.max_value, 'count': row.count } for row in rows ] # Usage queries = SensorQueries(['localhost']) # Get latest reading latest = queries.get_latest_reading(sensor_id) print(f"Latest: {latest['value']} {latest['unit']}") # Get last 24 hours readings = queries.get_last_24_hours(sensor_id) print(f"Readings in last 24h: {len(readings)}")

🔍 Query Patterns

Common query examples!

Real-Time Dashboard

-- Get current status for all sensors in factory -- Step 1: Get all sensors SELECT sensor_id, sensor_type, location FROM sensors_by_factory WHERE factory_id = 'factory-uuid'; -- Step 2: Get latest reading for each sensor (in parallel) SELECT reading_time, value, unit, quality FROM sensor_data_by_sensor WHERE sensor_id = ? LIMIT 1; -- Fast: Each query hits ONE partition!

Historical Trend Analysis

-- Get hourly averages for last week SELECT hour_bucket, avg_value, min_value, max_value FROM sensor_hourly_stats WHERE sensor_id = ? AND hour_bucket >= '2024-01-01 00:00:00' AND hour_bucket <= '2024-01-07 23:00:00'; -- Returns: 168 rows (7 days × 24 hours) -- Much faster than aggregating raw data!

Anomaly Detection

-- Get recent critical alerts SELECT alert_time, sensor_id, alert_type, value, message FROM alerts_by_factory WHERE factory_id = ? AND alert_time >= ? AND severity = 'critical' ALLOW FILTERING; -- Note: ALLOW FILTERING on small result sets OK -- Better: Create alerts_by_severity table for large scale

Time Range Queries

-- Get all readings for specific time window SELECT reading_time, value FROM sensor_data_by_sensor WHERE sensor_id = ? AND reading_time >= '2024-01-09 08:00:00' AND reading_time <= '2024-01-09 17:00:00'; -- Efficient: Scans ONE partition with time range -- Clustering by time DESC makes recent queries fast!

⚡ Performance Optimization

Production optimization techniques!

📊

Write Optimization

  • Batch writes: 50-100 per batch
  • Async inserts: Non-blocking writes
  • Prepared statements: 10x faster
  • Write consistency: ONE for speed
  • No read-before-write: Direct inserts
  • TTL: Auto-cleanup old data
🔍

Read Optimization

  • Single partition: One query = one partition
  • LIMIT clause: Fetch only needed rows
  • Pre-aggregation: Hourly stats table
  • Read consistency: ONE for speed
  • Paging: Avoid large result sets
  • Caching: Recent data in Redis
🗜️

Storage Optimization

  • TWCS: Time-window compaction
  • Compression: LZ4 compression
  • TTL: 90 days auto-expiration
  • Tombstone GC: Clean old deletes
  • Bloom filters: Reduce disk reads
  • Row cache: Hot data in memory
⚙️

Compaction Strategy

  • Strategy: TimeWindowCompactionStrategy
  • Window: 1 day buckets
  • Why TWCS? Time-series optimized
  • Benefit: Drop old windows entirely
  • TTL sync: Align with TTL
  • Performance: 10x faster than STCS

Aggregation Pipeline

# aggregation_job.py - Hourly aggregation from datetime import datetime, timedelta import schedule class HourlyAggregation: def __init__(self, session): self.session = session def aggregate_hour(self, sensor_id, hour_start): """Aggregate one hour of data""" hour_end = hour_start + timedelta(hours=1) # Get raw readings for hour query = """ SELECT value FROM sensor_data_by_sensor WHERE sensor_id = ? AND reading_time >= ? AND reading_time < ? """ rows = self.session.execute(query, (sensor_id, hour_start, hour_end)) values = [row.value for row in rows] if not values: return # Calculate statistics avg_value = sum(values) / len(values) min_value = min(values) max_value = max(values) # Insert aggregated stats insert = """ INSERT INTO sensor_hourly_stats (sensor_id, hour_bucket, avg_value, min_value, max_value, count) VALUES (?, ?, ?, ?, ?, ?) """ self.session.execute( insert, (sensor_id, hour_start, avg_value, min_value, max_value, len(values)) ) print(f"Aggregated {len(values)} readings for {hour_start}") def run_hourly_job(self): """Run every hour to aggregate previous hour""" # Aggregate last complete hour now = datetime.now() last_hour = now.replace(minute=0, second=0, microsecond=0) - timedelta(hours=1) # Get all active sensors sensors = self.get_all_sensors() for sensor_id in sensors: self.aggregate_hour(sensor_id, last_hour) # Schedule to run every hour schedule.every().hour.at(":05").do(aggregator.run_hourly_job)

🚀 Deployment Architecture

Production deployment!

Multi-Datacenter Architecture

                        ┌─────────────────────────────────────┐
                        │      Load Balancer (Global)         │
                        └──────────────┬──────────────────────┘
                                       │
                ┌──────────────────────┼──────────────────────┐
                │                      │                      │
                ▼                      ▼                      ▼
        ┌───────────────┐      ┌───────────────┐      ┌───────────────┐
        │ Datacenter 1  │      │ Datacenter 2  │      │ Datacenter 3  │
        │   (US-East)   │      │   (EU-West)   │      │   (AP-South)  │
        └───────────────┘      └───────────────┘      └───────────────┘
                │                      │                      │
        ┌───────┴─────────┐    ┌───────┴─────────┐    ┌───────┴─────────┐
        │                 │    │                 │    │                 │
    ┌───▼───┐ ┌───▼───┐  │┌───▼───┐ ┌───▼───┐  │┌───▼───┐ ┌───▼───┐  │
    │ C* N1 │ │ C* N2 │  ││ C* N1 │ │ C* N2 │  ││ C* N1 │ │ C* N2 │  │
    └───────┘ └───────┘  │└───────┘ └───────┘  │└───────┘ └───────┘  │
    ┌───▼───┐            │┌───▼───┐            │┌───▼───┐            │
    │ C* N3 │            ││ C* N3 │            ││ C* N3 │            │
    └───────┘            │└───────┘            │└───────┘            │
                         │                     │                     │
    ┌─────────────────┐  │┌─────────────────┐  │┌─────────────────┐  │
    │ Ingestion API   │  ││ Ingestion API   │  ││ Ingestion API   │  │
    │ Query API       │  ││ Query API       │  ││ Query API       │  │
    │ Aggregation Job │  ││ Aggregation Job │  ││ Aggregation Job │  │
    └─────────────────┘  │└─────────────────┘  │└─────────────────┘  │
                         │                     │                     │
    ┌─────────────────┐  │┌─────────────────┐  │┌─────────────────┐  │
    │   Monitoring    │  ││   Monitoring    │  ││   Monitoring    │  │
    │  (Prometheus)   │  ││  (Prometheus)   │  ││  (Prometheus)   │  │
    └─────────────────┘  │└─────────────────┘  │└─────────────────┘  │
                         └                     └                     └

    Replication: NetworkTopologyStrategy
    - DC1: 3 replicas
    - DC2: 3 replicas
    - DC3: 3 replicas
    
    Consistency: 
    - Writes: LOCAL_QUORUM (fast, local)
    - Reads: LOCAL_ONE (fast, eventual)
        

Deployment Configuration

# cassandra.yaml - Production settings # Cluster name cluster_name: 'IoT Sensors Production' # Listen addresses listen_address: 10.0.1.10 broadcast_address: 203.0.113.10 # Seed nodes (2-3 per DC) seed_provider: - class_name: org.apache.cassandra.locator.SimpleSeedProvider parameters: - seeds: "10.0.1.10,10.0.1.11,10.0.2.10,10.0.2.11" # Compaction concurrent_compactors: 4 # Memory memtable_heap_space_in_mb: 2048 memtable_offheap_space_in_mb: 2048 # Commit log commitlog_sync: periodic commitlog_sync_period_in_ms: 10000 # Monitoring jmx_port: 7199

Production Checklist

Component Configuration Status
Cluster Size 9 nodes (3 per DC) ✅
Replication RF=3 per DC ✅
Write Consistency LOCAL_QUORUM ✅
Read Consistency LOCAL_ONE ✅
TTL 90 days (raw), 2 years (aggregates) ✅
Compaction TimeWindowCompactionStrategy ✅
Monitoring Prometheus + Grafana ✅
Backup Daily snapshots to S3 ✅

🎉 Project Complete!

You built a production-ready IoT sensor data platform!

🎓 What You Built:

  • 📊 Scale: 100K writes/second, 8.64B data points/day
  • 🗺️ Data model: Query-driven design with 5 optimized tables
  • 📋 Schema: Time-series optimized with TWCS + TTL
  • 💻 Implementation: Python ingestion + query services
  • 🔍 Queries: Real-time dashboard, trends, alerts, aggregates
  • ⚡ Optimization: Batch writes, pre-aggregation, caching
  • 🚀 Deployment: Multi-DC with 9 nodes (3 per DC)

💡 Key Learnings:

  1. Data model first: Know your queries before schema
  2. Time-series design: Partition by sensor, cluster by time DESC
  3. TTL is essential: Auto-expire old data (90 days)
  4. TWCS for time-series: 10x better than STCS
  5. Pre-aggregate: Hourly stats for fast analytics
  6. Batch writes: 50-100 inserts per batch
  7. Multi-DC: Global deployment with LOCAL consistency

📊 Architecture Highlights:

  • ✅ 5 tables: sensor_data, sensors_by_factory, hourly_stats, alerts
  • ✅ Write path: 100K/sec with batch inserts
  • ✅ Read path: Single-partition queries (fast!)
  • ✅ Storage: TWCS + LZ4 + 90-day TTL
  • ✅ Aggregation: Hourly job for analytics
  • ✅ HA: RF=3 per DC, LOCAL_QUORUM writes

🚀 Next Steps:

  • 🔧 Add real-time anomaly detection with Spark
  • 📈 Build Grafana dashboards for visualization
  • 🔔 Implement alert notification system
  • 💾 Set up automated backups to S3
  • 📱 Create mobile app for factory managers
  • 🤖 Add machine learning for predictive maintenance

🎯 You're ready to build production IoT systems! 📊
Apply this pattern to ANY time-series use case!

Advertisement

Responsive Ad