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:

  1. ✅ Add spark-cassandra-connector dependency
  2. ✅ Configure connection in SparkSession
  3. ✅ Read data with .format("org.apache.spark.sql.cassandra")
  4. ✅ Transform using DataFrames or SQL
  5. ✅ Write results back to Cassandra
  6. ✅ Optimize with push-down and batching
  7. ✅ 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 📱