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:

  1. ✅ Design schema with time bucketing strategy
  2. ✅ Set up Kafka for async log ingestion
  3. ✅ Configure TWCS compaction & TTL
  4. ✅ Build consumer workers for batching
  5. ✅ Create query-specific tables (by trace, by user)
  6. ✅ Add real-time error alerting
  7. ✅ Set up Spark for analytics
  8. ✅ Integrate with monitoring (Grafana)
  9. ✅ 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 📱