Internet of Things
IoT Sensor Data
Handle millions of sensor readings per second with Cassandra!
📡 IoT Sensor Data with Cassandra
Cassandra is perfect for IoT sensor data - high write throughput, time-series storage, and horizontal scaling!
Why Cassandra for IoT?
- ⚡ High Write Throughput: Millions of writes per second
- 📈 Time-series Data: Built-in support for temporal data
- 🌍 Distributed: Deploy sensors globally, store locally
- 💾 Compression: Efficient storage for time-series
- 📊 Fast Queries: Sub-millisecond reads by sensor/time
- 🔄 No Downtime: Always-on for continuous ingestion
Real-world Scale
Production deployments handle:
• Apple: 160,000+ nodes for IoT/wearables
• Uber: 100M+ reads/writes per second
• Netflix: 1 trillion+ requests per day
🏗️ Time-series Data Modeling
Typical IoT Queries
# What we need to query:
Q1: Get latest reading from sensor
Q2: Get all readings from sensor in last 24 hours
Q3: Get readings from sensor between timestamps
Q4: Get average temperature per hour
Q5: Get all sensors reporting anomalies
Partitioning Strategy
BAD: One Partition
PRIMARY KEY (sensor_id, timestamp)
-- All data in one partition!
-- Gets huge over time ❌
GOOD: Time Buckets
PRIMARY KEY ((sensor_id, bucket), timestamp)
-- Bucket = day/hour
-- Limited partition size ✅
Partition Size Matters!
Rule of Thumb: Keep partitions under 100MB
- High-frequency sensors (1/sec) → Bucket by hour
- Medium-frequency (1/min) → Bucket by day
- Low-frequency (1/hour) → Bucket by month
📋 Schema Design
1
Basic Sensor Readings Table
CREATE TABLE sensor_readings (
sensor_id text,
bucket text, -- 'YYYY-MM-DD' or 'YYYY-MM-DD-HH'
timestamp timestamp,
temperature decimal,
humidity decimal,
pressure decimal,
battery_level int,
PRIMARY KEY ((sensor_id, bucket), timestamp)
) WITH CLUSTERING ORDER BY (timestamp DESC)
AND compaction = {
'class': 'TimeWindowCompactionStrategy',
'compaction_window_size': '24',
'compaction_window_unit': 'HOURS'
};
-- Why this works:
-- ✅ Partition key: (sensor_id, bucket) - distributes data
-- ✅ Clustering: timestamp DESC - newest first
-- ✅ Compaction: TWCS - perfect for time-series
2
Sensor Metadata Table
CREATE TABLE sensors (
sensor_id text PRIMARY KEY,
name text,
location text,
type text, -- temperature, humidity, pressure
install_date timestamp,
last_reading timestamp,
status text, -- active, inactive, error
metadata map<text, text>
);
3
Hourly Aggregates Table
CREATE TABLE hourly_aggregates (
sensor_id text,
hour timestamp, -- truncated to hour
avg_temperature decimal,
min_temperature decimal,
max_temperature decimal,
reading_count counter,
PRIMARY KEY (sensor_id, hour)
) WITH CLUSTERING ORDER BY (hour DESC);
📥 Data Ingestion Patterns
Python: MQTT → Cassandra
import paho.mqtt.client as mqtt
from cassandra.cluster import Cluster
from datetime import datetime
import json
# Connect to Cassandra
cluster = Cluster(['127.0.0.1'])
session = cluster.connect('iot')
# Prepared statement
insert_stmt = session.prepare("""
INSERT INTO sensor_readings
(sensor_id, bucket, timestamp, temperature, humidity, pressure)
VALUES (?, ?, ?, ?, ?, ?)
""")
def on_message(client, userdata, message):
data = json.loads(message.payload)
# Calculate bucket (by day)
now = datetime.now()
bucket = now.strftime('%Y-%m-%d')
# Insert to Cassandra
session.execute(insert_stmt, (
data['sensor_id'],
bucket,
now,
data['temperature'],
data['humidity'],
data['pressure']
))
# MQTT client
client = mqtt.Client()
client.on_message = on_message
client.connect("mqtt.broker.com", 1883)
client.subscribe("sensors/#")
client.loop_forever()
Batch Ingestion for Performance
from cassandra.query import BatchStatement, SimpleStatement
batch = BatchStatement()
insert_query = "INSERT INTO sensor_readings (...) VALUES (?, ?, ?, ?, ?, ?)"
readings_buffer = []
def on_message(client, userdata, message):
data = json.loads(message.payload)
readings_buffer.append(data)
# Flush every 100 readings
if len(readings_buffer) >= 100:
batch = BatchStatement()
for reading in readings_buffer:
now = datetime.now()
bucket = now.strftime('%Y-%m-%d')
batch.add(insert_stmt, (
reading['sensor_id'], bucket, now,
reading['temperature'], reading['humidity'], reading['pressure']
))
session.execute(batch)
readings_buffer.clear()
print(f"Inserted batch of 100 readings")
Node.js: HTTP API → Cassandra
const express = require('express');
const cassandra = require('cassandra-driver');
const client = new cassandra.Client({
contactPoints: ['127.0.0.1'],
localDataCenter: 'datacenter1',
keyspace: 'iot'
});
const app = express();
app.use(express.json());
// Endpoint for sensor data
app.post('/sensor/reading', async (req, res) => {
const { sensor_id, temperature, humidity, pressure } = req.body;
const now = new Date();
const bucket = now.toISOString().split('T')[0];
const query = `
INSERT INTO sensor_readings
(sensor_id, bucket, timestamp, temperature, humidity, pressure)
VALUES (?, ?, ?, ?, ?, ?)
`;
await client.execute(query,
[sensor_id, bucket, now, temperature, humidity, pressure],
{ prepare: true }
);
res.json({ status: 'success' });
});
app.listen(3000);
🔍 Query Patterns
Get Latest Reading
# Python
from datetime import datetime
sensor_id = 'sensor-123'
bucket = datetime.now().strftime('%Y-%m-%d')
query = """
SELECT * FROM sensor_readings
WHERE sensor_id = ? AND bucket = ?
LIMIT 1
"""
result = session.execute(query, (sensor_id, bucket))
latest = result.one()
print(f"Latest temp: {latest.temperature}°C")
Get Last 24 Hours
from datetime import datetime, timedelta
sensor_id = 'sensor-123'
now = datetime.now()
# Query today and yesterday buckets
buckets = [
now.strftime('%Y-%m-%d'),
(now - timedelta(days=1)).strftime('%Y-%m-%d')
]
readings = []
for bucket in buckets:
query = """
SELECT * FROM sensor_readings
WHERE sensor_id = ? AND bucket = ?
AND timestamp >= ?
"""
start_time = now - timedelta(hours=24)
result = session.execute(query, (sensor_id, bucket, start_time))
readings.extend(result)
print(f"Found {len(readings)} readings in last 24 hours")
Time Range Query
// Node.js - Query specific time range
async function getReadingsInRange(sensorId, startTime, endTime) {
const query = `
SELECT * FROM sensor_readings
WHERE sensor_id = ?
AND bucket = ?
AND timestamp >= ?
AND timestamp <= ?
`;
// Calculate bucket
const bucket = startTime.toISOString().split('T')[0];
const result = await client.execute(query,
[sensorId, bucket, startTime, endTime],
{ prepare: true }
);
return result.rows;
}
📊 Real-time Aggregation
Calculate Hourly Averages
import schedule
from datetime import datetime, timedelta
def calculate_hourly_aggregates():
# Get last hour's data
now = datetime.now()
hour_start = now.replace(minute=0, second=0, microsecond=0) - timedelta(hours=1)
hour_end = hour_start + timedelta(hours=1)
bucket = hour_start.strftime('%Y-%m-%d')
# Get all sensors
sensors = session.execute("SELECT sensor_id FROM sensors")
for sensor in sensors:
# Query readings for this hour
query = """
SELECT temperature FROM sensor_readings
WHERE sensor_id = ? AND bucket = ?
AND timestamp >= ? AND timestamp < ?
"""
readings = session.execute(query,
(sensor.sensor_id, bucket, hour_start, hour_end)
)
temps = [r.temperature for r in readings]
if temps:
# Calculate aggregates
avg_temp = sum(temps) / len(temps)
min_temp = min(temps)
max_temp = max(temps)
# Store aggregates
insert_agg = """
INSERT INTO hourly_aggregates
(sensor_id, hour, avg_temperature, min_temperature, max_temperature)
VALUES (?, ?, ?, ?, ?)
"""
session.execute(insert_agg, (
sensor.sensor_id, hour_start,
avg_temp, min_temp, max_temp
))
print(f"Aggregated {len(temps)} readings for {sensor.sensor_id}")
# Run every hour
schedule.every().hour.at(":05").do(calculate_hourly_aggregates)
Stream Processing with Spark
from pyspark.sql import SparkSession
from pyspark.sql.functions import avg, min, max, window
spark = SparkSession.builder \
.appName("IoT Aggregation") \
.config("spark.cassandra.connection.host", "127.0.0.1") \
.getOrCreate()
# Read from Cassandra
df = spark.read \
.format("org.apache.spark.sql.cassandra") \
.options(table="sensor_readings", keyspace="iot") \
.load()
# Aggregate by hour
hourly = df.groupBy(
"sensor_id",
window("timestamp", "1 hour")
).agg(
avg("temperature").alias("avg_temp"),
min("temperature").alias("min_temp"),
max("temperature").alias("max_temp")
)
# Write back to Cassandra
hourly.write \
.format("org.apache.spark.sql.cassandra") \
.options(table="hourly_aggregates", keyspace="iot") \
.mode("append") \
.save()
🗄️ Data Retention & TTL
Automatic Data Expiration
-- Set TTL on insert (7 days)
INSERT INTO sensor_readings (...) VALUES (...)
USING TTL 604800; -- 7 days in seconds
-- Set default TTL on table
ALTER TABLE sensor_readings
WITH default_time_to_live = 2592000; -- 30 days
Time-Window Compaction Strategy
-- Perfect for time-series data!
ALTER TABLE sensor_readings
WITH compaction = {
'class': 'TimeWindowCompactionStrategy',
'compaction_window_size': '24',
'compaction_window_unit': 'HOURS'
};
-- Benefits:
-- ✅ Better compression for similar timestamps
-- ✅ Efficient TTL expiration
-- ✅ Faster queries on recent data
Retention Strategy
- Raw data: 7-30 days with TTL
- Hourly aggregates: 1 year
- Daily aggregates: 5 years
- Monthly aggregates: Forever
💡 Best Practices
Performance Optimization
- ✅ Time buckets: Partition by sensor + time bucket
- ✅ TWCS compaction: Use TimeWindowCompactionStrategy
- ✅ Batch writes: Group 50-100 inserts together
- ✅ Prepared statements: Always prepare your queries
- ✅ TTL: Auto-expire old data to save storage
- ✅ Compression: LZ4 or Snappy for time-series
Common Mistakes
- ❌ No time buckets: Partitions grow forever
- ❌ Wrong compaction: Using SizeTiered for time-series
- ❌ No TTL: Storage costs explode
- ❌ Single writes: Batch for 10x better performance
- ❌ Querying old data: Use aggregates instead
Production Checklist
Before Going Live
- ☑️ Time bucketing implemented (hour/day)
- ☑️ TWCS compaction configured
- ☑️ TTL set on raw data table
- ☑️ Prepared statements for all queries
- ☑️ Batch writes for ingestion
- ☑️ Monitoring on write latency
- ☑️ Aggregation pipeline (Spark/scheduled)
- ☑️ Backup strategy in place
- ☑️ Tested at expected write volume
- ☑️ Capacity planning done
🎉 Build Your IoT Platform!
You're now ready to handle millions of sensor readings with Cassandra!
🚀 Implementation Steps:
- ✅ Design schema with time buckets
- ✅ Set up TWCS compaction
- ✅ Configure TTL for data retention
- ✅ Implement batch ingestion
- ✅ Build aggregation pipeline
- ✅ Add monitoring and alerts
- ✅ Scale as needed! 📈
💡 Key Takeaways:
- 📊 Time buckets: Essential for bounded partitions
- ⚡ TWCS: Perfect compaction for time-series
- ⏱️ TTL: Auto-expire old data automatically
- 📦 Batch writes: 10x better throughput
- 📈 Aggregate: Pre-calculate metrics for fast queries
- 🔧 Monitor: Watch write latency and partition sizes
📡 Real-world Examples:
- 🌡️ Temperature monitoring systems
- 🏭 Industrial equipment sensors
- 🚗 Connected vehicle telemetry
- 🏠 Smart home devices
- ⚡ Energy grid monitoring
- 🌾 Agricultural sensors
📡 Power your IoT platform with Cassandra! 🚀
Advertisement
📱 Responsive Ad 📱