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:
- ✅ Install kafka-python and cassandra-driver
- ✅ Create Kafka consumer with proper config
- ✅ Connect to Cassandra with connection pooling
- ✅ Use prepared statements and batching
- ✅ Implement error handling and retries
- ✅ Monitor consumer lag and throughput
- ✅ 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 📱