Distributed Logging
Logging System
Build a scalable logging platform handling billions of log entries per day!
📜 Distributed Logging with Cassandra
Cassandra is ideal for centralized logging - extreme write throughput, time-series storage, and query flexibility!
Why Cassandra for Logging?
- ⚡ Write-Optimized: Handle millions of logs per second
- 📈 Time-Series Native: Perfect for log timestamp ordering
- 🌍 Distributed: Collect logs from global services
- 💾 Cost-Effective: Efficient storage with compression & TTL
- 📊 Fast Queries: Find logs by service, level, or time range
- 🔄 High Availability: Never lose logs during outages
Production Examples
Companies using Cassandra for logging:
• Instagram: Billions of events per day
• Spotify: User activity & playback logs
• eBay: 250+ TB of operational logs
• Discord: Chat messages & user events
🏗️ Log Data Modeling
Typical Logging Queries
# What we need to query:
Q1: Get latest ERROR logs from service X
Q2: Get all logs from service X in last hour
Q3: Get logs by trace_id (distributed tracing)
Q4: Count ERROR/WARNING logs per hour
Q5: Search logs by message content
Q6: Get logs by user_id or request_id
Partitioning Strategy
BAD: Service Only
PRIMARY KEY (service_name, timestamp)
-- Partition grows forever!
-- Hot partitions ❌
GOOD: Service + Time Bucket
PRIMARY KEY ((service, bucket), timestamp)
-- Bounded partitions
-- Even distribution ✅
Partition Size Planning!
Log Volume Determines Bucket Size:
- High-traffic services (1000+ logs/sec) → Bucket by hour
- Medium traffic (100-1000/sec) → Bucket by day
- Low traffic (<100/sec) → Bucket by day or week
- Keep partitions under 100MB for optimal performance
📋 Schema Design
1
Main Logs Table
CREATE TABLE logs (
service_name text,
bucket text, -- 'YYYY-MM-DD' or 'YYYY-MM-DD-HH'
timestamp timestamp,
log_level text, -- ERROR, WARN, INFO, DEBUG
message text,
trace_id text,
request_id text,
user_id text,
host text,
environment text, -- prod, staging, dev
metadata map<text, text>,
PRIMARY KEY ((service_name, bucket), timestamp, log_level)
) WITH CLUSTERING ORDER BY (timestamp DESC, log_level ASC)
AND compaction = {
'class': 'TimeWindowCompactionStrategy',
'compaction_window_size': '24',
'compaction_window_unit': 'HOURS'
}
AND default_time_to_live = 2592000; -- 30 days
-- Why this works:
-- ✅ Partition: (service, bucket) - bounded size
-- ✅ Clustering: timestamp DESC - newest first
-- ✅ TWCS compaction - optimal for time-series
-- ✅ Auto-expire with TTL - saves storage
2
Logs by Trace ID (Distributed Tracing)
CREATE TABLE logs_by_trace (
trace_id text,
timestamp timestamp,
service_name text,
log_level text,
message text,
span_id text,
PRIMARY KEY (trace_id, timestamp)
) WITH CLUSTERING ORDER BY (timestamp ASC)
AND default_time_to_live = 604800; -- 7 days
-- Use case: Track a request across microservices
3
Error Summary Table
CREATE TABLE error_summary (
service_name text,
hour timestamp, -- truncated to hour
error_count counter,
warning_count counter,
PRIMARY KEY (service_name, hour)
) WITH CLUSTERING ORDER BY (hour DESC);
4
Logs by User (Optional)
CREATE TABLE logs_by_user (
user_id text,
bucket text,
timestamp timestamp,
service_name text,
log_level text,
message text,
PRIMARY KEY ((user_id, bucket), timestamp)
) WITH CLUSTERING ORDER BY (timestamp DESC)
AND default_time_to_live = 604800; -- 7 days
-- Use case: Debugging user-specific issues
📥 Log Ingestion Patterns
Python: Direct Logging to Cassandra
import logging
from cassandra.cluster import Cluster
from datetime import datetime
import json
class CassandraLogHandler(logging.Handler):
def __init__(self, service_name):
super().__init__()
self.service_name = service_name
# Connect to Cassandra
cluster = Cluster(['127.0.0.1'])
self.session = cluster.connect('logging')
# Prepared statement
self.insert_stmt = self.session.prepare("""
INSERT INTO logs
(service_name, bucket, timestamp, log_level, message, host)
VALUES (?, ?, ?, ?, ?, ?)
""")
def emit(self, record):
now = datetime.now()
bucket = now.strftime('%Y-%m-%d')
self.session.execute(self.insert_stmt, (
self.service_name,
bucket,
now,
record.levelname,
record.getMessage(),
record.hostname if hasattr(record, 'hostname') else 'unknown'
))
# Usage
logger = logging.getLogger('my-service')
logger.addHandler(CassandraLogHandler('my-service'))
logger.info("User logged in successfully")
logger.error("Database connection failed")
Kafka → Cassandra Pipeline
from kafka import KafkaConsumer
from cassandra.cluster import Cluster
from cassandra.query import BatchStatement
import json
from datetime import datetime
# Cassandra setup
cluster = Cluster(['127.0.0.1'])
session = cluster.connect('logging')
insert_stmt = session.prepare("""
INSERT INTO logs (service_name, bucket, timestamp, log_level, message,
trace_id, request_id, host, metadata)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
""")
# Kafka consumer
consumer = KafkaConsumer(
'application-logs',
bootstrap_servers=['localhost:9092'],
value_deserializer=lambda m: json.loads(m.decode('utf-8'))
)
batch = BatchStatement()
batch_size = 0
for message in consumer:
log = message.value
timestamp = datetime.fromisoformat(log['timestamp'])
bucket = timestamp.strftime('%Y-%m-%d')
batch.add(insert_stmt, (
log['service'],
bucket,
timestamp,
log['level'],
log['message'],
log.get('trace_id'),
log.get('request_id'),
log['host'],
log.get('metadata', {})
))
batch_size += 1
# Execute batch every 100 logs
if batch_size >= 100:
session.execute(batch)
batch = BatchStatement()
batch_size = 0
print(f"Inserted batch of 100 logs")
Node.js: Winston → Cassandra
const winston = require('winston');
const cassandra = require('cassandra-driver');
const client = new cassandra.Client({
contactPoints: ['127.0.0.1'],
localDataCenter: 'datacenter1',
keyspace: 'logging'
});
class CassandraTransport extends winston.Transport {
constructor(options) {
super(options);
this.serviceName = options.serviceName;
this.query = `
INSERT INTO logs
(service_name, bucket, timestamp, log_level, message, host)
VALUES (?, ?, ?, ?, ?, ?)
`;
}
log(info, callback) {
const now = new Date();
const bucket = now.toISOString().split('T')[0];
client.execute(this.query,
[this.serviceName, bucket, now, info.level, info.message, require('os').hostname()],
{ prepare: true }
)
.then(() => callback())
.catch((err) => callback(err));
}
}
// Create logger
const logger = winston.createLogger({
transports: [
new CassandraTransport({ serviceName: 'api-service' })
]
});
logger.info('Server started on port 3000');
logger.error('Failed to connect to Redis');
🔍 Query Patterns
Get Latest Logs from Service
# Python - Get last 100 logs from api-service
from datetime import datetime
service = 'api-service'
bucket = datetime.now().strftime('%Y-%m-%d')
query = """
SELECT * FROM logs
WHERE service_name = ? AND bucket = ?
LIMIT 100
"""
result = session.execute(query, (service, bucket))
for log in result:
print(f"[{log.timestamp}] {log.log_level}: {log.message}")
Get ERROR Logs Only
# Get only ERROR logs from last hour
from datetime import datetime, timedelta
service = 'payment-service'
now = datetime.now()
bucket = now.strftime('%Y-%m-%d')
one_hour_ago = now - timedelta(hours=1)
query = """
SELECT * FROM logs
WHERE service_name = ?
AND bucket = ?
AND timestamp >= ?
AND log_level = 'ERROR'
"""
errors = session.execute(query, (service, bucket, one_hour_ago))
print(f"Found {len(list(errors))} errors in last hour")
Query by Trace ID (Distributed Tracing)
// Node.js - Get all logs for a trace
async function getTraceLogs(traceId) {
const query = `
SELECT * FROM logs_by_trace
WHERE trace_id = ?
ORDER BY timestamp ASC
`;
const result = await client.execute(query, [traceId], { prepare: true });
// Shows complete request flow across services
return result.rows;
}
const trace = await getTraceLogs('trace-12345');
trace.forEach(log => {
console.log(`[${log.service_name}] ${log.message}`);
});
Time Range Query
# Get logs between specific timestamps
query = """
SELECT * FROM logs
WHERE service_name = ?
AND bucket = ?
AND timestamp >= ?
AND timestamp <= ?
"""
start_time = datetime(2025, 1, 10, 14, 0)
end_time = datetime(2025, 1, 10, 15, 0)
bucket = '2025-01-10'
logs = session.execute(query, ('api-service', bucket, start_time, end_time))
📊 Log Analysis & Aggregation
Count Errors Per Hour
import schedule
from datetime import datetime, timedelta
def aggregate_errors():
# Get last hour
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 services
services = session.execute("SELECT DISTINCT service_name FROM logs")
for service in services:
# Count errors and warnings
error_query = """
SELECT COUNT(*) as count FROM logs
WHERE service_name = ? AND bucket = ?
AND timestamp >= ? AND timestamp < ?
AND log_level = ?
"""
errors = session.execute(error_query,
(service.service_name, bucket, hour_start, hour_end, 'ERROR')
).one().count
warnings = session.execute(error_query,
(service.service_name, bucket, hour_start, hour_end, 'WARN')
).one().count
# Update counter table
update = """
UPDATE error_summary
SET error_count = error_count + ?,
warning_count = warning_count + ?
WHERE service_name = ? AND hour = ?
"""
session.execute(update, (errors, warnings, service.service_name, hour_start))
print(f"{service.service_name}: {errors} errors, {warnings} warnings")
# Run every hour
schedule.every().hour.at(":05").do(aggregate_errors)
Real-time Error Detection
from collections import defaultdict
import time
# Track error rate in memory
error_counts = defaultdict(lambda: {'count': 0, 'timestamp': time.time()})
def check_error_rate(service, log_level):
if log_level == 'ERROR':
error_counts[service]['count'] += 1
# Check if more than 100 errors in last minute
now = time.time()
if now - error_counts[service]['timestamp'] < 60:
if error_counts[service]['count'] > 100:
# ALERT! High error rate
send_alert(f"{service} has {error_counts[service]['count']} errors/min!")
else:
# Reset counter
error_counts[service] = {'count': 1, 'timestamp': now}
Spark Analytics on Logs
from pyspark.sql import SparkSession
from pyspark.sql.functions import count, window, col
spark = SparkSession.builder \
.appName("Log Analytics") \
.config("spark.cassandra.connection.host", "127.0.0.1") \
.getOrCreate()
# Read logs from Cassandra
df = spark.read \
.format("org.apache.spark.sql.cassandra") \
.options(table="logs", keyspace="logging") \
.load()
# Find most common errors
top_errors = df \
.filter(col("log_level") == "ERROR") \
.groupBy("message") \
.agg(count("*").alias("count")) \
.orderBy(col("count").desc()) \
.limit(10)
top_errors.show(truncate=False)
🗄️ Data Retention & TTL
Automatic Expiration with TTL
-- Different retention for different log levels
INSERT INTO logs (...) VALUES (...)
USING TTL 2592000; -- 30 days for INFO/DEBUG
INSERT INTO logs (...) VALUES (...)
USING TTL 7776000; -- 90 days for ERROR/WARN
-- Set default TTL on table
ALTER TABLE logs
WITH default_time_to_live = 2592000; -- 30 days default
Retention in Code
# Python - Set TTL based on log level
def get_ttl(log_level):
ttl_map = {
'ERROR': 7776000, # 90 days
'WARN': 5184000, # 60 days
'INFO': 2592000, # 30 days
'DEBUG': 604800, # 7 days
}
return ttl_map.get(log_level, 2592000)
# Insert with dynamic TTL
insert_query = """
INSERT INTO logs (service_name, bucket, timestamp, log_level, message)
VALUES (?, ?, ?, ?, ?)
USING TTL ?
"""
ttl = get_ttl(log_level)
session.execute(insert_query, (service, bucket, now, log_level, message, ttl))
Recommended Retention
- DEBUG logs: 7 days (high volume, low value)
- INFO logs: 30 days (normal operations)
- WARN logs: 60 days (potential issues)
- ERROR logs: 90 days (critical debugging)
- Aggregated metrics: 1-2 years
💡 Best Practices
Performance Optimization
- ✅ Time buckets: Partition by service + date/hour
- ✅ TWCS compaction: Essential for time-series logs
- ✅ Batch writes: Use Kafka/message queue for buffering
- ✅ Prepared statements: Cache query compilation
- ✅ TTL: Auto-expire based on log level importance
- ✅ Multiple tables: Different queries = different tables
Common Mistakes
- ❌ No time bucketing: Unbounded partition growth
- ❌ Full table scans: Always include partition key
- ❌ No TTL: Log storage costs spiral out of control
- ❌ Synchronous logging: Slows down application
- ❌ Text search: Use Elasticsearch for full-text search
- ❌ Single table: Create query-specific tables
Production Checklist
Before Going Live
- ☑️ Time bucketing by service + date/hour
- ☑️ TWCS compaction configured
- ☑️ TTL set based on log level
- ☑️ Async logging (Kafka, RabbitMQ, etc.)
- ☑️ Batch writes implemented
- ☑️ Multiple tables for different queries
- ☑️ Error alerting pipeline (>X errors/min)
- ☑️ Log aggregation scheduled jobs
- ☑️ Distributed tracing integration
- ☑️ Storage capacity monitoring
Architecture Patterns
High-Volume Setup
For 1M+ logs/second:
- Kafka for buffering
- Consumer groups for parallel writes
- Hourly time buckets
- 7-day TTL for most logs
- Separate ERROR table
Hybrid Approach
Best of both worlds:
- Cassandra for structured logs
- Elasticsearch for text search
- S3 for long-term archive
- Grafana for visualization
- Alerts via Prometheus
Cost-Optimized
Reduce storage costs:
- Aggressive TTL (7 days)
- LZ4 compression enabled
- Sample DEBUG logs (10%)
- Archive to object storage
- Smaller cluster size
🎉 Build Your Logging Platform!
You're now ready to build a production-grade distributed logging system!
🚀 Implementation Roadmap:
- ✅ Design schema with time bucketing strategy
- ✅ Set up Kafka for async log ingestion
- ✅ Configure TWCS compaction & TTL
- ✅ Build consumer workers for batching
- ✅ Create query-specific tables (by trace, by user)
- ✅ Add real-time error alerting
- ✅ Set up Spark for analytics
- ✅ Integrate with monitoring (Grafana)
- ✅ Scale horizontally as needed! 📈
💡 Key Takeaways:
- 📊 Time bucketing: Essential for bounded partitions
- ⚡ Async ingestion: Never block your application
- ⏱️ TTL strategy: Different retention per log level
- 📦 Batch writes: Kafka + batch consumers
- 🔍 Query tables: Create tables per query pattern
- 🚨 Real-time alerts: Monitor error rates continuously
📜 Production Use Cases:
- 🌐 Microservices centralized logging
- 🔍 Distributed request tracing
- 📊 Application performance monitoring (APM)
- 🚨 Real-time error detection & alerting
- 📈 Log analytics & trending
- 🔐 Security audit logs & compliance
📜 Never lose a log entry with Cassandra! 🚀
Advertisement
📱 Responsive Ad 📱