Advanced Topics

Analytics Workloads

Run complex analytics on Cassandra data with Spark, Presto, aggregations, and real-time OLAP!

📊 What are Analytics Workloads?

The Business Intelligence Problem 💼

You're running an e-commerce platform storing billions of orders in Cassandra. Your CEO asks:

  • 📈 "What's our revenue by product category this quarter?"
  • 🌍 "Which regions have highest customer lifetime value?"
  • 👥 "Show me cohort analysis of user signups by month"
  • 🔍 "Find patterns in customer purchase behavior"
  • 💡 "Identify our top 100 most valuable customers"

Problem: Cassandra is optimized for fast point reads, not aggregations across billions of rows!

Analytics Workloads Explained

Analytics = Complex queries that aggregate, filter, group, and analyze large datasets to extract insights - think SQL GROUP BY, SUM, AVG across millions/billions of rows.

Typical Analytics Queries:
  • 📊 Aggregations: SUM, AVG, COUNT, MIN, MAX
  • 🔢 Group By: Revenue by category, users by region
  • 📅 Time Series: Daily/weekly/monthly trends
  • 🔍 Filtering: Complex WHERE clauses
  • 🔗 Joins: Combine data from multiple tables
  • 📈 Window Functions: Running totals, rankings

Analytics Use Cases

📊

Business Intelligence

  • Revenue reports
  • Sales dashboards
  • KPI tracking
  • Executive summaries
  • Trend analysis
🔬

Data Science

  • ML training data
  • Feature engineering
  • Statistical analysis
  • Cohort analysis
  • A/B testing
💡

Customer Insights

  • Behavior patterns
  • Segmentation
  • Churn prediction
  • Lifetime value
  • Recommendation tuning
⚙️

Operations

  • Log analysis
  • Performance metrics
  • Error tracking
  • Capacity planning
  • Cost optimization

⚡ OLTP vs OLAP: Understanding the Difference

Critical: Cassandra is OLTP, Not OLAP!

Cassandra excels at OLTP (Online Transaction Processing) - fast reads/writes by key. It's NOT designed for OLAP (Online Analytical Processing) - complex aggregations. For analytics, use complementary tools!

OLTP vs OLAP Comparison

Aspect OLTP (Cassandra) OLAP (Analytics)
Purpose Transaction processing Analysis & reporting
Query Type Simple (by key) Complex (aggregations)
Data Volume Single/few rows Millions/billions of rows
Response Time < 10ms Seconds to minutes
Users Thousands (app users) Tens (analysts)
Example Get user profile Monthly revenue report

The Analytics Stack

Typical Architecture:

Application → Cassandra (OLTP: fast reads/writes)
          ↓
   Apache Spark / Presto (OLAP: analytics)
          ↓
   Data Warehouse (Redshift, Snowflake)
          ↓
   BI Tools (Tableau, Looker)

Key Insight: Cassandra stores operational data, analytics tools process it!

🔥 Apache Spark + Cassandra

What is Spark?

Apache Spark = Distributed computing engine for large-scale data processing. The Spark Cassandra Connector enables reading Cassandra data directly into Spark for analytics.

Why Spark + Cassandra?

✅ Perfect Fit

  • Distributed processing matches distributed storage
  • Data locality (process where data lives)
  • Push-down predicates (filter in Cassandra)
  • Parallel reads across all nodes
  • Handles billions of rows
  • SQL, Python, Scala support

Common Use Cases

  • ETL pipelines (transform data)
  • Aggregation reports
  • ML model training
  • Data migration
  • Batch processing
  • Real-time streaming (Spark Streaming)

Spark Setup & Configuration

# 1. Add Spark Cassandra Connector dependency spark-shell --packages com.datastax.spark:spark-cassandra-connector_2.12:3.4.0 # 2. Configure Spark session from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("CassandraAnalytics") \ .config("spark.cassandra.connection.host", "localhost") \ .config("spark.cassandra.connection.port", "9042") \ .config("spark.cassandra.auth.username", "cassandra") \ .config("spark.cassandra.auth.password", "cassandra") \ .getOrCreate()

Reading Data from Cassandra

# Read entire table into Spark DataFrame df = spark.read \ .format("org.apache.spark.sql.cassandra") \ .options(table="orders", keyspace="ecommerce") \ .load() df.show(10) # Read with filter (pushed down to Cassandra!) df_filtered = spark.read \ .format("org.apache.spark.sql.cassandra") \ .options(table="orders", keyspace="ecommerce") \ .load() \ .filter("order_date >= '2024-01-01'") # Select specific columns (reduces data transfer) df_cols = spark.read \ .format("org.apache.spark.sql.cassandra") \ .options(table="orders", keyspace="ecommerce") \ .load() \ .select("order_id", "user_id", "total_amount")

Analytics Examples

# Example 1: Revenue by product category from pyspark.sql.functions import sum, avg, count revenue_by_category = df.groupBy("category") \ .agg( sum("total_amount").alias("total_revenue"), count("*").alias("order_count"), avg("total_amount").alias("avg_order_value") ) \ .orderBy("total_revenue", ascending=False) revenue_by_category.show() # Example 2: Top 100 customers by lifetime value top_customers = df.groupBy("user_id") \ .agg( sum("total_amount").alias("lifetime_value"), count("*").alias("order_count") ) \ .orderBy("lifetime_value", ascending=False) \ .limit(100) top_customers.show() # Example 3: Monthly cohort analysis from pyspark.sql.functions import year, month monthly_cohorts = df.groupBy( year("created_at").alias("year"), month("created_at").alias("month") ).agg( count("user_id").alias("new_users"), sum("total_amount").alias("revenue") ).orderBy("year", "month") monthly_cohorts.show()

Writing Back to Cassandra

# Write aggregated results back to Cassandra revenue_by_category.write \ .format("org.apache.spark.sql.cassandra") \ .options(table="revenue_reports", keyspace="analytics") \ .mode("append") \ .save() # ✅ Pre-compute analytics, store results for fast dashboards!

Performance Optimization

# 1. Push-down predicates (filter in Cassandra) df = spark.read \ .format("org.apache.spark.sql.cassandra") \ .options(table="orders", keyspace="ecommerce") \ .load() \ .filter("user_id = 'abc-123'") ← Pushed to Cassandra! # 2. Select only needed columns df = spark.read \ .format("org.apache.spark.sql.cassandra") \ .options(table="orders", keyspace="ecommerce") \ .load() \ .select("order_id", "total_amount") ← Less data transferred # 3. Repartition for parallelism df = df.repartition(100) ← More parallel tasks # 4. Cache frequently used DataFrames df.cache() df.count() ← Trigger caching

🚀 Presto/Trino + Cassandra

What is Presto/Trino?

Presto (now Trino) = Distributed SQL query engine for interactive analytics. Query Cassandra using standard SQL with sub-second latency for exploratory analysis.

Presto vs Spark

🔥 Spark

  • Batch: Minutes to hours
  • Use for: ETL, ML, heavy aggregations
  • Language: Python, Scala, SQL
  • Latency: Seconds to minutes
  • Best for: Scheduled jobs

🚀 Presto/Trino

  • Interactive: Seconds to minutes
  • Use for: Ad-hoc queries, dashboards
  • Language: Standard SQL
  • Latency: Sub-second to seconds
  • Best for: Exploratory analysis

Presto Configuration

# 1. Configure Cassandra connector in Presto # File: etc/catalog/cassandra.properties connector.name=cassandra cassandra.contact-points=127.0.0.1 cassandra.load-policy.dc-aware.local-dc=datacenter1 cassandra.username=cassandra cassandra.password=cassandra

Querying with Presto SQL

-- Connect to Presto CLI presto --server localhost:8080 --catalog cassandra --schema ecommerce -- Standard SQL queries! SELECT category, SUM(total_amount) AS revenue, COUNT(*) AS orders, AVG(total_amount) AS avg_order FROM orders WHERE order_date >= DATE '2024-01-01' GROUP BY category ORDER BY revenue DESC; -- Complex analytics with window functions SELECT user_id, order_date, total_amount, SUM(total_amount) OVER ( PARTITION BY user_id ORDER BY order_date ) AS running_total FROM orders; -- Join multiple Cassandra tables (federated query!) SELECT u.username, COUNT(o.order_id) AS order_count, SUM(o.total_amount) AS total_spent FROM users u JOIN orders o ON u.user_id = o.user_id GROUP BY u.username ORDER BY total_spent DESC LIMIT 10;

Performance Warning

Presto joins can be expensive! Cassandra doesn't support joins natively, so Presto fetches data from both tables and joins in-memory. For large datasets, consider pre-joining in Spark and writing back to Cassandra.

📈 Analytics Patterns & Strategies

Pattern 1: Pre-Aggregation (Recommended)

Pre-compute analytics with Spark, store in Cassandra

# Spark job runs nightly daily_revenue = spark.read \ .format("org.apache.spark.sql.cassandra") \ .options(table="orders", keyspace="ecommerce") \ .load() \ .filter("order_date = current_date()") \ .groupBy("category") \ .agg(sum("total_amount").alias("revenue")) # Write results to analytics table daily_revenue.write \ .format("org.apache.spark.sql.cassandra") \ .options(table="daily_revenue", keyspace="analytics") \ .mode("append") \ .save() # Dashboard queries Cassandra directly (fast!) SELECT * FROM analytics.daily_revenue WHERE date = '2024-01-15';

✅ Best Practice: Pre-aggregate in Spark, serve from Cassandra!

Pattern 2: Lambda Architecture

-- Real-time + Batch layers -- Batch Layer (Spark): Processes historical data nightly -- → Writes to: analytics.batch_revenue -- Real-time Layer (Spark Streaming): Processes live data -- → Writes to: analytics.realtime_revenue -- Serving Layer: Merge both for complete view SELECT COALESCE(b.category, r.category) AS category, (b.revenue + r.revenue) AS total_revenue FROM analytics.batch_revenue b FULL OUTER JOIN analytics.realtime_revenue r ON b.category = r.category;

Pattern 3: Data Export (ETL to Data Warehouse)

# Export Cassandra data to Snowflake/Redshift for complex analytics # 1. Read from Cassandra df = spark.read \ .format("org.apache.spark.sql.cassandra") \ .options(table="orders", keyspace="ecommerce") \ .load() # 2. Transform data df_transformed = df.select( "order_id", "user_id", "total_amount", date_format("order_date", "yyyy-MM-dd").alias("order_date") ) # 3. Write to data warehouse (Snowflake, Redshift, BigQuery) df_transformed.write \ .format("snowflake") \ .options( sfURL="account.snowflakecomputing.com", sfUser="user", sfPassword="pass", sfDatabase="analytics", sfSchema="public", sfWarehouse="compute_wh", dbtable="orders" ) \ .mode("overwrite") \ .save() # ✅ Now run complex SQL in Snowflake/Redshift!

✅ Analytics Best Practices

✅ DO These

  • Pre-aggregate in Spark
  • Use push-down predicates
  • Select only needed columns
  • Cache frequently accessed data
  • Schedule batch jobs off-peak
  • Export to data warehouse for complex queries
  • Monitor Cassandra load during analytics
  • Use separate analytics cluster

❌ DON'T Do These

  • Run analytics on production cluster
  • Scan entire tables without filters
  • Join large tables in Presto
  • Run analytics during peak hours
  • Use ALLOW FILTERING for analytics
  • Expect real-time aggregations
  • Query without partition keys
  • Ignore data locality

Recommended Architecture

  • ✅ Production writes: Cassandra (OLTP)
  • ✅ Batch analytics: Spark (nightly/hourly)
  • ✅ Interactive queries: Presto (ad-hoc)
  • ✅ Complex analytics: Export to data warehouse
  • ✅ Dashboards: Query pre-aggregated Cassandra tables

🎯 Analytics Summary

You now understand analytics with Cassandra!

📚 Key Takeaways:

  • 📊 Cassandra = OLTP, not OLAP
  • 🔥 Use Spark for batch analytics
  • 🚀 Use Presto/Trino for interactive queries
  • ✅ Pre-aggregate results in Spark
  • 📈 Store analytics in Cassandra for fast serving
  • 🏢 Export to data warehouse for complex analytics
  • ⚡ Lambda architecture for real-time + batch
  • 🎯 Separate analytics from production cluster

Cassandra + Spark/Presto = Powerful analytics at scale! 📊🚀

Advertisement

Responsive Ad