Advanced Topics

Change Data Capture

Master Cassandra CDC to stream database changes in real-time for data pipelines, analytics, and event-driven architectures!

📡 What is Change Data Capture (CDC)?

The Data Sync Challenge 🔄

Imagine you're running an e-commerce platform with Cassandra as your main database:

  • 📊 Analytics team needs real-time data in their data warehouse
  • 🔍 Search team needs to update Elasticsearch whenever products change
  • 💾 Cache team needs to invalidate Redis when data updates
  • 📧 Email team needs to send confirmations when orders change status

How do you capture EVERY change happening in Cassandra and stream it to all these systems in real-time?

That's what CDC solves! It captures every INSERT, UPDATE, DELETE from Cassandra's commit log and streams them as events.

CDC in Simple Terms

Change Data Capture (CDC) = A mechanism to capture and stream database mutations (INSERT/UPDATE/DELETE) as they happen, enabling real-time data pipelines and event-driven architectures.

Think of CDC Like:

A transaction log reader that watches Cassandra's commit log (like a security camera watching every change), captures mutations, and streams them to downstream systems without impacting database performance.

CDC vs Traditional Data Sync

🐌

Traditional Batch Sync

  • Method: Poll database every N minutes
  • Latency: Minutes to hours
  • Load: Full table scans (heavy)
  • Missing: Deletes hard to track
  • Complexity: Track last_updated timestamps
⚡

CDC Streaming

  • Method: Stream from commit log
  • Latency: Sub-second real-time
  • Load: Minimal (log reading)
  • Captures: ALL changes including deletes
  • Simplicity: Event-driven, automatic

❓ Why Use CDC?

Common Use Cases

📊 Real-Time Analytics

Problem: Analytics warehouse needs fresh data

Solution: Stream changes to Kafka → Spark → Data Warehouse

  • Near real-time dashboards
  • Live business intelligence
  • Up-to-date reporting

🔍 Search Index Sync

Problem: Elasticsearch out of sync with Cassandra

Solution: CDC streams changes → Update ES indices

  • Auto-sync search indices
  • No stale search results
  • Real-time full-text search

💾 Cache Invalidation

Problem: Redis cache becomes stale

Solution: CDC triggers cache invalidation

  • Auto-invalidate on changes
  • No stale cached data
  • Event-driven updates

🔔 Event Notifications

Problem: Trigger actions on data changes

Solution: CDC events → Email/SMS/Webhooks

  • Order confirmation emails
  • Status change notifications
  • Webhook triggers

🔄 Data Replication

Problem: Replicate to other databases

Solution: CDC streams to PostgreSQL/MySQL

  • Multi-database architecture
  • Hybrid cloud setups
  • Database migration

📈 Audit Logging

Problem: Track all data changes for compliance

Solution: CDC → Audit log storage

  • Compliance requirements
  • Change history tracking
  • Security audits

⚙️ How CDC Works Internally

CDC Architecture

CDC works by intercepting writes at the CommitLog level - the lowest layer where ALL mutations are written before going to MemTable.

The Write Path with CDC

Normal Write Path: 1. Client sends write → Coordinator Node 2. Coordinator → CommitLog (durability) 3. Coordinator → MemTable (performance) 4. MemTable flushes → SSTable (disk) CDC-Enabled Write Path: 1. Client sends write → Coordinator Node 2. Coordinator → CDC CommitLog (CDC captures here!) 3. Coordinator → MemTable 4. MemTable flushes → SSTable 5. CDC Log Reader → Reads mutations → Streams to consumers Key Point: CDC operates at CommitLog level = Captures EVERY write before it's applied!

CDC Storage Structure

# Cassandra Data Directory /var/lib/cassandra/ ├── commitlog/ # Normal commit logs │ ├── CommitLog-7-*.log │ └── CommitLog-8-*.log │ ├── cdc_raw/ # CDC-specific commit logs │ ├── CommitLog-7-*.log # Contains CDC mutations │ └── CommitLog-8-*.log │ └── data/ └── keyspace/table/ # CDC logs are SEPARATE from normal commit logs # Prevents CDC from impacting normal operations

Important CDC Characteristics

  • 📝 CommitLog-based: Reads from commit log, not SSTables
  • ⚡ Low overhead: Minimal performance impact
  • 🔄 At-least-once: CDC guarantees at-least-once delivery
  • 📊 Table-level: Enable CDC per table, not globally
  • 💾 Disk space: CDC logs consume additional disk space
  • ⏰ Retention: CDC logs must be cleaned up by consumer

🔧 Setting Up CDC

Step 1: Enable CDC on Cassandra Cluster

# Edit cassandra.yaml cdc_enabled: true # Configure CDC directory cdc_raw_directory: /var/lib/cassandra/cdc_raw # Set CDC total space (important!) cdc_total_space_in_mb: 4096 # 4GB for CDC logs # Restart Cassandra for changes to take effect sudo systemctl restart cassandra

CDC Disk Space Protection

cdc_total_space_in_mb is a safety limit! If CDC logs exceed this limit, Cassandra will REJECT WRITES to CDC-enabled tables to prevent disk fill-up.

Solution: Your CDC consumer MUST read and delete CDC logs regularly!

Step 2: Enable CDC on Specific Tables

-- Enable CDC when creating table CREATE TABLE products ( product_id uuid PRIMARY KEY, name text, price decimal, stock int ) WITH cdc = true; -- Enable CDC on existing table ALTER TABLE products WITH cdc = true; -- Disable CDC on table ALTER TABLE products WITH cdc = false; -- Check if CDC is enabled DESC TABLE products; -- Look for: cdc = true

Step 3: Verify CDC is Working

# Check CDC directory ls -lh /var/lib/cassandra/cdc_raw/ # Output shows CDC commit logs: CommitLog-7-1234567890.log CommitLog-7-1234567891.log # Insert test data cqlsh> INSERT INTO products (product_id, name, price) VALUES (uuid(), 'Test Product', 99.99); # CDC log should grow in size ls -lh /var/lib/cassandra/cdc_raw/

📝 Understanding CommitLog CDC

What's in a CDC Log?

CDC CommitLog contains binary mutation records of every INSERT, UPDATE, DELETE on CDC-enabled tables.

CDC Mutation Record Structure

Each CDC record contains: { "keyspace": "ecommerce", "table": "products", "partition_key": "123e4567-e89b-12d3-a456-426614174000", "operation": "INSERT", // INSERT, UPDATE, DELETE "timestamp": 1640000000000000, // Microseconds "columns": { "product_id": "123e4567...", "name": "Laptop", "price": 999.99, "stock": 50 } } For DELETE: { "operation": "DELETE", "partition_key": "...", "clustering_keys": {...} // If composite key }

Reading CDC Logs

# CDC logs are in binary format (not human-readable) # You need a CDC consumer to parse them # Popular CDC consumers: 1. DataStax CDC for Apache Pulsar 2. Debezium for Cassandra (Kafka Connect) 3. Custom consumer using Cassandra CDC API # These tools read CDC logs and stream to: - Apache Kafka - Apache Pulsar - Amazon Kinesis - Custom applications

🌊 Building CDC Data Pipelines

Architecture: Cassandra → Kafka → Consumers

┌─────────────┐ │ Cassandra │ │ (CDC ON) │ └──────┬──────┘ │ CDC Logs ↓ ┌─────────────────┐ │ CDC Consumer │ ← Reads cdc_raw/ directory │ (Debezium/etc) │ ← Parses mutations └──────┬──────────┘ │ JSON events ↓ ┌─────────────┐ │ Kafka │ ← Streams to topics │ Topics │ ← ecommerce.products └──────┬──────┘ │ ├─────────────→ 📊 Analytics (Spark/Flink) │ ├─────────────→ 🔍 Search (Elasticsearch) │ ├─────────────→ 💾 Cache Invalidation (Redis) │ └─────────────→ 📧 Notifications (Email/SMS)

Example: Debezium CDC Connector

# Install Debezium Cassandra Connector for Kafka Connect # Configure connector (JSON) { "name": "cassandra-cdc-connector", "config": { "connector.class": "io.debezium.connector.cassandra.CassandraConnector", "cassandra.hosts": "localhost", "cassandra.port": "9042", "cassandra.cdc.dir": "/var/lib/cassandra/cdc_raw", "kafka.topic.prefix": "cdc", "cassandra.keyspaces": "ecommerce" } } # Debezium will stream changes to Kafka topics: # - cdc.ecommerce.products # - cdc.ecommerce.orders

Example: Consuming CDC Events from Kafka

# Python Kafka Consumer Example from kafka import KafkaConsumer import json consumer = KafkaConsumer( 'cdc.ecommerce.products', bootstrap_servers=['localhost:9092'], value_deserializer=lambda m: json.loads(m.decode('utf-8')) ) for message in consumer: event = message.value if event['operation'] == 'INSERT': # New product added - update Elasticsearch update_search_index(event['columns']) elif event['operation'] == 'UPDATE': # Product updated - invalidate cache invalidate_cache(event['partition_key']) elif event['operation'] == 'DELETE': # Product deleted - remove from search remove_from_search(event['partition_key'])

🎯 Complete Use Case Examples

Use Case 1: Real-Time Product Catalog Sync

Problem:

E-commerce site with Cassandra (write DB) and Elasticsearch (search). Product updates in Cassandra must appear in search instantly.

Solution with CDC:

1. Enable CDC on products table ALTER TABLE products WITH cdc = true; 2. CDC → Kafka → Python Consumer 3. Consumer updates Elasticsearch on every change # Python Consumer def sync_to_elasticsearch(event): product_id = event['partition_key'] if event['operation'] == 'DELETE': es.delete(index='products', id=product_id) else: es.index( index='products', id=product_id, body=event['columns'] )

✅ Result:

  • Search results always fresh (< 1 second lag)
  • No manual sync jobs needed
  • Automatic delete handling

Use Case 2: Cache Invalidation on Updates

# Problem: Redis cache becomes stale when Cassandra updates # CDC Solution: def handle_cdc_event(event): cache_key = f"product:{event['partition_key']}" if event['operation'] in ['UPDATE', 'DELETE']: # Invalidate cache on change redis.delete(cache_key) print(f"Invalidated cache for {cache_key}") elif event['operation'] == 'INSERT': # Pre-warm cache with new product redis.setex( cache_key, 3600, # 1 hour TTL json.dumps(event['columns']) )

Use Case 3: Data Warehouse Streaming

# Architecture: Cassandra CDC → Kafka → Spark Streaming → Snowflake/Redshift # Spark Streaming Job (Scala) val stream = spark .readStream .format("kafka") .option("kafka.bootstrap.servers", "localhost:9092") .option("subscribe", "cdc.ecommerce.*") .load() stream .selectExpr("CAST(value AS STRING) as json") .select(from_json($"json", schema).as("data")) .writeStream .format("snowflake") .option("dbtable", "products_changelog") .start() # Result: Near real-time analytics dashboard!

✅ CDC Best Practices

✅ DO These

  • Monitor CDC disk space usage
  • Set appropriate cdc_total_space_in_mb
  • Clean CDC logs after processing
  • Enable CDC only on needed tables
  • Use idempotent consumers
  • Handle at-least-once delivery
  • Monitor consumer lag
  • Test CDC pipeline thoroughly

❌ DON'T Do These

  • Enable CDC on ALL tables
  • Ignore CDC disk space limits
  • Leave CDC logs unprocessed
  • Assume exactly-once delivery
  • Process CDC logs manually
  • Skip consumer error handling
  • Ignore consumer monitoring
  • Mix CDC with batch sync

Critical CDC Warnings

  • ⚠️ Disk Space: If CDC logs fill disk, writes FAIL! Monitor closely!
  • ⚠️ Consumer Must Run: CDC logs accumulate if consumer stops
  • ⚠️ At-Least-Once: Duplicates possible - make consumers idempotent
  • ⚠️ Performance: Minimal overhead but monitor during peak load
  • ⚠️ Ordering: Events ordered per partition, not globally

Monitoring Checklist

  • ✅ CDC disk space used (< 80% of limit)
  • ✅ Consumer lag (< 5 seconds ideal)
  • ✅ CDC log file count (growing = consumer issue)
  • ✅ Event processing rate (events/sec)
  • ✅ Consumer errors and retries
  • ✅ Downstream system health (Kafka, ES, etc)

🎯 CDC Summary

You now understand Cassandra Change Data Capture!

📚 Key Takeaways:

  • 📡 CDC = Stream database changes in real-time
  • 📝 Works at CommitLog level (captures all mutations)
  • ⚡ Sub-second latency for downstream systems
  • 🔧 Enable per-table with cdc = true
  • 💾 Monitor CDC disk space to prevent write failures
  • 🌊 Popular consumers: Debezium, DataStax Pulsar CDC
  • 🎯 Perfect for: search sync, cache invalidation, analytics, audit logs

CDC enables event-driven architectures and real-time data pipelines! 📡🚀

Advertisement

Responsive Ad