Big Data Analytics
Spark + Cassandra
Powerful analytics on Cassandra data with Apache Spark!
⚡ Spark Cassandra Connector
Integrate Apache Spark with Cassandra for distributed analytics, ETL, and big data processing!
Why Spark + Cassandra?
- ⚡ Fast Analytics: Process billions of rows in seconds
- 📊 SQL Queries: Use Spark SQL on Cassandra tables
- 🔄 ETL Workflows: Extract, transform, load at scale
- 📈 Aggregations: Complex analytics and machine learning
- 💻 Multiple Languages: Scala, Python, Java, R support
- 🚀 Distributed: Parallel processing across cluster
Connector Version: 3.x
Compatible with Spark 3.x and Cassandra 3.x/4.x
Open source project maintained by DataStax community.
📦 Setup & Configuration
1
Add Dependency
# Maven (pom.xml)
<dependency>
<groupId>com.datastax.spark</groupId>
<artifactId>spark-cassandra-connector_2.12</artifactId>
<version>3.4.1</version>
</dependency>
# SBT (build.sbt)
libraryDependencies += "com.datastax.spark" %% "spark-cassandra-connector" % "3.4.1"
# spark-submit with package
$ spark-submit --packages com.datastax.spark:spark-cassandra-connector_2.12:3.4.1 app.py
2
Configure Connection (Python/PySpark)
from pyspark.sql import SparkSession
# Create Spark session
spark = SparkSession.builder \
.appName("CassandraApp") \
.config("spark.cassandra.connection.host", "127.0.0.1") \
.config("spark.cassandra.connection.port", "9042") \
.config("spark.cassandra.auth.username", "cassandra") \
.config("spark.cassandra.auth.password", "cassandra") \
.getOrCreate()
print("Connected to Cassandra!")
3
Configure Connection (Scala)
import org.apache.spark.sql.SparkSession
val spark = SparkSession.builder()
.appName("CassandraApp")
.config("spark.cassandra.connection.host", "127.0.0.1")
.config("spark.cassandra.connection.port", "9042")
.getOrCreate()
📖 Reading Data from Cassandra
Python/PySpark
# Read entire table
df = spark.read \
.format("org.apache.spark.sql.cassandra") \
.options(table="users", keyspace="myapp") \
.load()
# Show data
df.show()
# Count rows
print(f"Total users: {df.count()}")
# Select specific columns
df.select("id", "name", "email").show()
# Filter data
active_users = df.filter(df.status == "active")
active_users.show()
With WHERE Clause (Push-down)
# Push filter to Cassandra (efficient!)
df = spark.read \
.format("org.apache.spark.sql.cassandra") \
.options(
table="posts",
keyspace="myapp",
pushdown="true"
) \
.load() \
.filter("user_id = 'uuid-123'")
# Cassandra executes WHERE clause before sending data!
# Much faster than filtering in Spark
Scala Examples
import org.apache.spark.sql._
// Read table
val df = spark.read
.format("org.apache.spark.sql.cassandra")
.options(Map("table" -> "users", "keyspace" -> "myapp"))
.load()
// Show schema
df.printSchema()
// Filter and select
df.filter("age > 18").select("name", "email").show()
💾 Writing Data to Cassandra
Insert/Update Data
# Create DataFrame
data = [
("uuid-1", "Alice", "alice@example.com"),
("uuid-2", "Bob", "bob@example.com"),
("uuid-3", "Charlie", "charlie@example.com")
]
columns = ["id", "name", "email"]
df = spark.createDataFrame(data, columns)
# Write to Cassandra
df.write \
.format("org.apache.spark.sql.cassandra") \
.options(table="users", keyspace="myapp") \
.mode("append") \
.save()
print("Data written!")
Write Modes
# Append (default) - INSERT
df.write.format("org.apache.spark.sql.cassandra") \
.mode("append") \
.options(table="users", keyspace="myapp") \
.save()
# Overwrite - TRUNCATE then INSERT
df.write.format("org.apache.spark.sql.cassandra") \
.mode("overwrite") \
.options(table="users", keyspace="myapp") \
.save()
# Ignore - Skip if exists
df.write.format("org.apache.spark.sql.cassandra") \
.mode("ignore") \
.options(table="users", keyspace="myapp") \
.save()
Batch Writing with Options
# Write with custom options
df.write \
.format("org.apache.spark.sql.cassandra") \
.options(
table="users",
keyspace="myapp",
confirm_truncate="true",
ttl="86400", # 1 day TTL
writeTime="current_timestamp"
) \
.mode("append") \
.save()
📊 DataFrames & Spark SQL
Register as Temp Table
# Read from Cassandra
df = spark.read \
.format("org.apache.spark.sql.cassandra") \
.options(table="users", keyspace="myapp") \
.load()
# Register as temporary view
df.createOrReplaceTempView("users")
# Run SQL queries!
result = spark.sql("""
SELECT age, COUNT(*) as count
FROM users
WHERE age > 18
GROUP BY age
ORDER BY age
""")
result.show()
Complex SQL Queries
# Join multiple tables
users_df = spark.read.format("org.apache.spark.sql.cassandra") \
.options(table="users", keyspace="myapp").load()
posts_df = spark.read.format("org.apache.spark.sql.cassandra") \
.options(table="posts", keyspace="myapp").load()
# Register tables
users_df.createOrReplaceTempView("users")
posts_df.createOrReplaceTempView("posts")
# Join query
result = spark.sql("""
SELECT u.name, COUNT(p.id) as post_count
FROM users u
LEFT JOIN posts p ON u.id = p.user_id
GROUP BY u.name
ORDER BY post_count DESC
LIMIT 10
""")
result.show()
DataFrame Operations
from pyspark.sql.functions import col, count, avg, max, min
# Aggregations
df.groupBy("country") \
.agg(
count("*").alias("total_users"),
avg("age").alias("avg_age"),
max("age").alias("max_age")
) \
.show()
# Window functions
from pyspark.sql.window import Window
window = Window.partitionBy("country").orderBy(col("age").desc())
df.withColumn("rank", rank().over(window)).show()
📈 Analytics Operations
Aggregations
Group by and aggregate
# User activity by country
df.groupBy("country") \
.count() \
.orderBy("count", ascending=False) \
.show(10)
# Revenue by product
df.groupBy("product") \
.agg(sum("amount")) \
.show()
Filtering & Joins
Complex queries
# Filter and join
active = users.filter("status = 'active'")
result = active.join(
orders,
active.id == orders.user_id
)
result.show()
Time Series
Temporal analytics
# Daily metrics
df.withColumn("date", to_date("timestamp")) \
.groupBy("date") \
.agg(count("*")) \
.orderBy("date") \
.show()
ETL Pipelines
Extract-Transform-Load
# Read → Transform → Write
spark.read.cassandraFormat(...) \
.load() \
.filter(...) \
.groupBy(...) \
.agg(...) \
.write.cassandraFormat(...) \
.save()
Complete Analytics Example
from pyspark.sql.functions import *
# Read user events
events = spark.read \
.format("org.apache.spark.sql.cassandra") \
.options(table="user_events", keyspace="analytics") \
.load()
# Calculate daily active users
dau = events \
.withColumn("date", to_date("timestamp")) \
.groupBy("date") \
.agg(countDistinct("user_id").alias("active_users")) \
.orderBy("date")
# Show results
dau.show(30)
# Save back to Cassandra
dau.write \
.format("org.apache.spark.sql.cassandra") \
.options(table="daily_metrics", keyspace="analytics") \
.mode("append") \
.save()
🚀 Performance Optimization
Push-Down Optimizations
Connector automatically pushes operations to Cassandra:
- ✅ WHERE clauses: Filters executed in Cassandra
- ✅ Column selection: Only requested columns fetched
- ✅ Aggregations: COUNT pushed to Cassandra when possible
- ⚡ Result: Much less data transferred!
Partition Awareness
# Connector is partition-aware!
# Spark executors read from local Cassandra nodes
# Configure partition size
spark = SparkSession.builder \
.config("spark.cassandra.input.split.sizeInMB", "64") \
.getOrCreate()
# Result: Minimal network traffic, maximum locality
Batch Size Tuning
# Configure write batch size
df.write \
.format("org.apache.spark.sql.cassandra") \
.options(
table="users",
keyspace="myapp",
spark_cassandra_output_batch_size_rows="auto",
spark_cassandra_output_batch_size_bytes="1024",
spark_cassandra_output_concurrent_writes="5"
) \
.save()
Connection Pooling
# Configure connection pools
spark = SparkSession.builder \
.config("spark.cassandra.connection.connections_per_executor_max", "10") \
.config("spark.cassandra.connection.keep_alive_ms", "60000") \
.getOrCreate()
💡 Best Practices
Query Optimization
- ✅ Filter early: Use WHERE on partition keys
- ✅ Select specific columns: Don't SELECT *
- ✅ Partition awareness: Co-locate Spark and Cassandra
- ✅ Batch writes: Tune batch sizes for throughput
- ✅ Use DataFrames: Better optimization than RDDs
Common Pitfalls
- ❌ Full table scans: Always filter when possible
- ❌ Small batches: Too many small writes = slow
- ❌ No push-down: Forgetting to enable push-down
- ❌ Network separation: Spark & Cassandra on different networks
- ❌ No caching: Caching frequently used DataFrames helps
Production Configuration
# Production-ready Spark session
spark = SparkSession.builder \
.appName("CassandraETL") \
# Connection
.config("spark.cassandra.connection.host", "10.0.0.1,10.0.0.2,10.0.0.3") \
.config("spark.cassandra.connection.port", "9042") \
.config("spark.cassandra.auth.username", "user") \
.config("spark.cassandra.auth.password", "pass") \
# Performance
.config("spark.cassandra.input.split.sizeInMB", "64") \
.config("spark.cassandra.output.batch.size.rows", "auto") \
.config("spark.cassandra.output.concurrent.writes", "5") \
# Connection pooling
.config("spark.cassandra.connection.connections_per_executor_max", "10") \
.config("spark.cassandra.connection.keep_alive_ms", "60000") \
.getOrCreate()
🎉 Master Spark + Cassandra!
You're now ready to perform powerful analytics on Cassandra data!
🚀 Quick Start:
- ✅ Add spark-cassandra-connector dependency
- ✅ Configure connection in SparkSession
- ✅ Read data with
.format("org.apache.spark.sql.cassandra") - ✅ Transform using DataFrames or SQL
- ✅ Write results back to Cassandra
- ✅ Optimize with push-down and batching
- ✅ Deploy and scale! 🎯
💡 Key Benefits:
- ⚡ Fast Analytics: Distributed processing at scale
- 📊 SQL Queries: Familiar SQL on NoSQL data
- 🔄 ETL Made Easy: Read, transform, write
- 📈 Rich APIs: DataFrames, SQL, RDDs, ML
- 🚀 Optimized: Push-down, partition awareness
- 💻 Multi-Language: Python, Scala, Java, R
📚 Use Cases:
- 📊 Real-time analytics dashboards
- 🔄 ETL pipelines for data warehouses
- 📈 Machine learning feature engineering
- 📉 Time series analysis and forecasting
- 🔍 Log analysis and aggregation
- 💼 Business intelligence reporting
⚡ Unlock big data analytics with Spark! 🚀
Advertisement
📱 Responsive Ad 📱