Real-time Streaming

Kafka + Cassandra

Build real-time streaming pipelines with Kafka and Cassandra!

🔄 Kafka + Cassandra Integration

Combine Kafka's real-time streaming with Cassandra's scalable storage for powerful event-driven architectures!

Why Kafka + Cassandra?

  • 🔄 Real-time Ingestion: Stream millions of events per second
  • 📊 Event Sourcing: Store every event durably in Cassandra
  • ⚡ Low Latency: Sub-millisecond reads and writes
  • 📈 Scalable: Both systems scale horizontally
  • 🔌 Decoupled: Kafka buffers, Cassandra persists
  • 💪 Resilient: Handle spikes and failures gracefully

Perfect Partnership!

Kafka: Stream processing & messaging

Cassandra: Long-term storage & queries

Together: Complete real-time data platform!

📦 Setup & Dependencies

1

Kafka Consumer Dependencies (Python)

# Install Kafka consumer $ pip install kafka-python # Install Cassandra driver $ pip install cassandra-driver # Or both at once $ pip install kafka-python cassandra-driver
2

Java/Scala Dependencies

<!-- Kafka Client --> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>3.6.1</version> </dependency> <!-- Cassandra Driver --> <dependency> <groupId>com.datastax.oss</groupId> <artifactId>java-driver-core</artifactId> <version>4.17.0</version> </dependency>

📥 Kafka → Cassandra (Consumer Pattern)

Python Consumer Writing to Cassandra

from kafka import KafkaConsumer from cassandra.cluster import Cluster import json # Connect to Cassandra cluster = Cluster(['127.0.0.1']) session = cluster.connect('myapp') # Prepare statement insert_stmt = session.prepare(""" INSERT INTO events (id, user_id, event_type, timestamp, data) VALUES (?, ?, ?, ?, ?) """) # Create Kafka consumer consumer = KafkaConsumer( 'user-events', bootstrap_servers=['localhost:9092'], auto_offset_reset='earliest', enable_auto_commit=True, group_id='cassandra-consumer', value_deserializer=lambda x: json.loads(x.decode('utf-8')) ) # Consume and write to Cassandra for message in consumer: event = message.value # Insert into Cassandra session.execute(insert_stmt, ( event['id'], event['user_id'], event['event_type'], event['timestamp'], json.dumps(event['data']) )) print(f"Stored event: {event['id']}")

Batch Processing for Performance

from kafka import KafkaConsumer from cassandra.cluster import Cluster from cassandra.query import BatchStatement import json cluster = Cluster(['127.0.0.1']) session = cluster.connect('myapp') insert_stmt = session.prepare(""" INSERT INTO events (id, user_id, event_type, timestamp) VALUES (?, ?, ?, ?) """) consumer = KafkaConsumer('user-events', ...) batch = BatchStatement() batch_count = 0 BATCH_SIZE = 100 for message in consumer: event = message.value # Add to batch batch.add(insert_stmt, ( event['id'], event['user_id'], event['event_type'], event['timestamp'] )) batch_count += 1 # Execute batch when full if batch_count >= BATCH_SIZE: session.execute(batch) batch = BatchStatement() batch_count = 0 print(f"Batch of {BATCH_SIZE} events written")

Java Consumer Example

import org.apache.kafka.clients.consumer.*; import com.datastax.oss.driver.api.core.CqlSession; public class KafkaToCassandra { public static void main(String[] args) { // Cassandra connection CqlSession session = CqlSession.builder().build(); // Kafka consumer config Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("group.id", "cassandra-consumer"); props.put("key.deserializer", StringDeserializer.class); props.put("value.deserializer", StringDeserializer.class); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); consumer.subscribe(Collections.singletonList("events")); // Consume and write while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); for (ConsumerRecord<String, String> record : records) { session.execute( "INSERT INTO events (id, data) VALUES (?, ?)", UUID.randomUUID(), record.value() ); } } } }

📤 Cassandra → Kafka (CDC Pattern)

Using Debezium CDC

# Debezium configuration for Cassandra CDC { "name": "cassandra-cdc-connector", "config": { "connector.class": "io.debezium.connector.cassandra.CassandraConnector", "cassandra.hosts": "localhost", "cassandra.port": "9042", "cassandra.keyspace": "myapp", "kafka.topic.prefix": "cdc", "snapshot.mode": "initial" } } # Result: Changes in Cassandra → Kafka topics # Table "users" → Topic "cdc.myapp.users"

Manual CDC with Polling

from kafka import KafkaProducer from cassandra.cluster import Cluster import json, time # Connections cluster = Cluster(['127.0.0.1']) session = cluster.connect('myapp') producer = KafkaProducer( bootstrap_servers=['localhost:9092'], value_serializer=lambda v: json.dumps(v).encode('utf-8') ) # Track last processed timestamp last_timestamp = None while True: # Query new records if last_timestamp: query = "SELECT * FROM users WHERE updated_at > ? ALLOW FILTERING" rows = session.execute(query, (last_timestamp,)) else: rows = session.execute("SELECT * FROM users") # Publish to Kafka for row in rows: event = { 'id': str(row.id), 'name': row.name, 'email': row.email, 'updated_at': str(row.updated_at) } producer.send('user-changes', event) last_timestamp = row.updated_at time.sleep(5) # Poll every 5 seconds

🔌 Kafka Connect

DataStax Kafka Connector

# Install connector $ confluent-hub install datastax/kafka-connect-cassandra-sink:2.0.0 # Connector configuration (JSON) { "name": "cassandra-sink", "config": { "connector.class": "com.datastax.oss.kafka.sink.CassandraSinkConnector", "tasks.max": "1", "topics": "user-events", // Cassandra connection "contactPoints": "localhost", "loadBalancing.localDc": "datacenter1", "port": 9042, // Mapping "topic.user-events.myapp.events.mapping": "id=value.id, user_id=value.user_id, event_type=value.type" } } # Deploy connector $ curl -X POST http://localhost:8083/connectors \ -H "Content-Type: application/json" \ -d @cassandra-sink.json

Connector Features

DataStax Connector Benefits

  • ✅ No code: Configuration-based integration
  • ✅ Schema mapping: Flexible field mappings
  • ✅ Error handling: Dead letter queues
  • ✅ Batching: Automatic batch optimization
  • ✅ Monitoring: JMX metrics and health checks

⚡ Stream Processing with Kafka Streams

Kafka Streams + Cassandra Lookups

import org.apache.kafka.streams.*; import com.datastax.oss.driver.api.core.CqlSession; public class EnrichmentStream { public static void main(String[] args) { CqlSession cassandra = CqlSession.builder().build(); StreamsBuilder builder = new StreamsBuilder(); // Read from Kafka KStream<String, String> events = builder.stream("raw-events"); // Enrich with Cassandra data KStream<String, String> enriched = events.mapValues(eventId -> { // Lookup in Cassandra Row row = cassandra.execute( "SELECT * FROM users WHERE id = ?", eventId ).one(); return enrichEvent(eventId, row); }); // Write back to Kafka enriched.to("enriched-events"); KafkaStreams streams = new KafkaStreams(builder.build(), props); streams.start(); } }

Python Stream Processing (Faust)

import faust from cassandra.cluster import Cluster # Cassandra connection cluster = Cluster(['127.0.0.1']) session = cluster.connect('myapp') # Faust app app = faust.App('myapp', broker='kafka://localhost:9092') # Define event model class Event(faust.Record): user_id: str event_type: str timestamp: int # Input topic events_topic = app.topic('events', value_type=Event) # Process stream @app.agent(events_topic) async def process_events(events): async for event in events: # Enrich from Cassandra user = session.execute( "SELECT * FROM users WHERE id = ?", (event.user_id,) ).one() # Write enriched data to Cassandra session.execute(""" INSERT INTO enriched_events (id, user_name, event_type) VALUES (?, ?, ?) """, (event.user_id, user.name, event.event_type))

🎯 Real-time Architecture Patterns

📥

Event Ingestion

High-volume write pattern

  • Kafka buffers spikes
  • Consumer writes to Cassandra
  • Backpressure handled
  • Perfect for IoT, logs
🔄

Event Sourcing

Store all events

  • Kafka: Event log
  • Cassandra: Event store
  • Replay & audit trail
  • Time travel queries
📊

CQRS Pattern

Separate read/write models

  • Write to Kafka
  • Multiple read models
  • Cassandra for queries
  • Eventual consistency
📡

CDC Streaming

Capture database changes

  • Cassandra changes → Kafka
  • Real-time replication
  • Microservice sync
  • Analytics pipelines

Complete Pipeline Example

# Architecture: IoT sensor data pipeline # # Sensors → Kafka → Stream Processing → Cassandra # ↓ # Real-time alerts # # 1. Producers send sensor data to Kafka from kafka import KafkaProducer producer = KafkaProducer(bootstrap_servers=['localhost:9092']) producer.send('sensor-data', { 'sensor_id': 'sensor-123', 'temperature': 25.3, 'timestamp': 1704067200 }) # 2. Stream processor checks thresholds for message in consumer: data = message.value if data['temperature'] > 30: # Send alert producer.send('alerts', data) # 3. Store in Cassandra session.execute(""" INSERT INTO sensor_readings (sensor_id, timestamp, temperature) VALUES (?, ?, ?) """, (data['sensor_id'], data['timestamp'], data['temperature']))

💡 Best Practices

Performance Optimization

  • ✅ Batch writes: Use Cassandra BatchStatement for multiple inserts
  • ✅ Consumer groups: Parallelize with multiple consumers
  • ✅ Prepared statements: Always prepare Cassandra queries
  • ✅ Connection pooling: Reuse Cassandra connections
  • ✅ Partition keys: Use Kafka key for Cassandra partition key

Common Pitfalls

  • ❌ No batching: Writing one record at a time = slow
  • ❌ Large batches: Same partition key required for batches
  • ❌ No error handling: Consumer crashes on write failures
  • ❌ Auto-commit too fast: May lose data on crash
  • ❌ No monitoring: Consumer lag can grow unnoticed

Production Configuration

# Production-ready consumer from kafka import KafkaConsumer from cassandra.cluster import Cluster from cassandra.policies import DCAwareRoundRobinPolicy # Cassandra with production settings cluster = Cluster( contact_points=['10.0.0.1', '10.0.0.2', '10.0.0.3'], load_balancing_policy=DCAwareRoundRobinPolicy(local_dc='DC1'), protocol_version=4 ) session = cluster.connect('myapp') # Kafka consumer with production settings consumer = KafkaConsumer( 'events', bootstrap_servers=['kafka1:9092', 'kafka2:9092', 'kafka3:9092'], group_id='cassandra-writer', auto_offset_reset='earliest', enable_auto_commit=False, # Manual commit for safety max_poll_records=500, # Batch size session_timeout_ms=30000, heartbeat_interval_ms=10000 ) # Process with error handling try: for message in consumer: try: # Write to Cassandra session.execute(...) # Commit offset after successful write consumer.commit() except Exception as e: print(f"Write failed: {e}") # Handle error - retry, DLQ, etc. finally: consumer.close() cluster.shutdown()

🎉 Build Real-time Pipelines!

You're now ready to build production streaming pipelines with Kafka and Cassandra!

🚀 Quick Start:

  1. ✅ Install kafka-python and cassandra-driver
  2. ✅ Create Kafka consumer with proper config
  3. ✅ Connect to Cassandra with connection pooling
  4. ✅ Use prepared statements and batching
  5. ✅ Implement error handling and retries
  6. ✅ Monitor consumer lag and throughput
  7. ✅ Deploy and scale! 🎯

💡 Key Benefits:

  • 🔄 Real-time: Sub-second latency pipelines
  • 📊 Scalable: Handle millions of events/second
  • ⚡ Resilient: Kafka buffers, Cassandra persists
  • 🔌 Decoupled: Independent scaling
  • 💪 Durable: No data loss
  • 📈 Flexible: Multiple consumption patterns

🎯 Use Cases:

  • 📊 Real-time analytics and dashboards
  • 🌐 IoT sensor data ingestion
  • 📱 User activity tracking
  • 🔔 Event-driven microservices
  • 📡 Change data capture (CDC)
  • 🔄 Event sourcing architectures

🔄 Stream all the things with Kafka + Cassandra! 🚀

Advertisement

📱 Responsive Ad 📱