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:

  1. ✅ Design schema with time buckets
  2. ✅ Set up TWCS compaction
  3. ✅ Configure TTL for data retention
  4. ✅ Implement batch ingestion
  5. ✅ Build aggregation pipeline
  6. ✅ Add monitoring and alerts
  7. ✅ 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 📱