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
- Know your queries first! Model data around access patterns
- Denormalize: Duplicate data for query efficiency
- Partition key: Distributes data across nodes
- Clustering key: Sorts data within partition
- 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:
- Data model first: Know your queries before schema
- Time-series design: Partition by sensor, cluster by time DESC
- TTL is essential: Auto-expire old data (90 days)
- TWCS for time-series: 10x better than STCS
- Pre-aggregate: Hourly stats for fast analytics
- Batch writes: 50-100 inserts per batch
- 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