Chapter 4 - Cloud Storage Systems (Spark?)

Updated 4 Oct 2026

Project Goals

  • Extend the MapReduce model to better support two common classes of analytics apps:
    • Iterative algorithms (machine learning, graphs)
    • Interactive data mining
  • Enhance programmability:
    • Integrate into Scala programming language
    • Allow interactive use from Scala interpreter

Motivation

The Problem with Traditional Cluster Programming

  • Most current cluster programming models are based on acyclic data flow from stable storage to stable storage
  • Data flows: Input → Map → Reduce → Output

Analogy: Think of traditional MapReduce like a factory assembly line where each product must go through every station sequentially, and the final product is stored in a warehouse. Every time you want to modify the product, you have to retrieve it from the warehouse and send it through the entire assembly line again.

Benefits of Data Flow Model

  • Runtime flexibility: Runtime can decide where to run tasks
  • Automatic fault recovery: Can automatically recover from failures

The Inefficiency Problem

Acyclic data flow is inefficient for applications that repeatedly reuse a working set of data:

  • Iterative algorithms (machine learning, graphs)
  • Interactive data mining tools (R, Excel, Python)

Current issue: With frameworks like Hadoop MapReduce, apps reload data from stable storage on each query

Analogy: Imagine reading a textbook where you have to go back to the library, check out the book, read one chapter, return it, then repeat this process for every chapter. Spark is like keeping the book on your desk so you can reference it quickly whenever needed.


Solution: Resilient Distributed Datasets (RDDs)

Core Concept

RDDs allow apps to keep working sets in memory for efficient reuse

Key Properties Retained from MapReduce

  • Fault tolerance
  • Data locality
  • Scalability
  • Support for wide range of applications

RDD Concept & Properties

1. Fault Tolerance

  • RDDs are fault-tolerant
  • Spark does NOT replicate data by default
  • Instead, it remembers how the data was created (lineage)
  • If a partition is lost (e.g., node failure), Spark recomputes it automatically

2. Partitioning

  • An RDD is split into partitions and stored across multiple nodes
  • Each partition can be processed in parallel
  • Enables scalable, high-performance computation

3. Lineage Tracking

  • Spark tracks the DAG (Directed Acyclic Graph) of transformations
  • Example: File → map → filter → reduce
  • If a partition is lost → Spark re-executes only the missing part
  • No need to recompute everything

Analogy: Think of RDD lineage like a recipe. If your cake gets destroyed, you don't need someone to give you another cake - you just follow the recipe again to recreate it. The recipe (lineage) is much smaller to store than the actual cake (data).


Hadoop MapReduce vs Spark RDD

Hadoop MapReduceSpark RDD
Disk-basedIn-memory
Slow iterationFast iterative computation
Limited APIsRich transformations
Poor fault modelLineage-based recovery

Real-world impact: This difference means Spark can be 100x faster for in-memory computations compared to Hadoop MapReduce.


Spark Programming Model

Resilient Distributed Datasets (RDDs)

  • Immutable, partitioned collections of objects
  • Created through parallel transformations (map, filter, groupBy, join, …) on data in stable storage
  • Can be cached for efficient reuse

Operations on RDDs

1. Transformations

These operations are applied to create a new RDD:

  • map - Apply function to each element
  • filter - Select elements matching a condition
  • groupBy - Group elements by key
  • join - Join two RDDs
  • And more...

Key property: Transformations are lazy - they are not executed until an action is called

2. Actions

These operations trigger computation and return results:

  • count - Count elements
  • reduce - Aggregate elements
  • collect - Retrieve all elements
  • save - Save to storage
  • And more...

Example of Lazy Evaluation

rdd2 = rdd1.map(lambda x: x*2)   # NOT executed yet (transformation)
rdd2.count()                     # Execution starts here (action)

Analogy: Transformations are like writing a shopping list (you're planning what to buy), while actions are like actually going to the store and making the purchases.


Example 1: Log Mining

Use Case

Load error messages from a log into memory, then interactively search for various patterns

Code Example

lines = spark.textFile("hdfs://...")
errors = lines.filter(_.startsWith("ERROR"))
messages = errors.map(_.split('\t')(2))
cachedMsgs = messages.cache()
 
cachedMsgs.filter(_.contains("foo")).count
cachedMsgs.filter(_.contains("bar")).count
...

Architecture

[Driver] --tasks--> [Worker 1: Block 1, Cache 1]
         --tasks--> [Worker 2: Block 2, Cache 2]
         --tasks--> [Worker 3: Block 3, Cache 3]
         <--results--

Performance Results

  • Full-text search of Wikipedia:
    • In-memory: < 1 sec
    • On-disk: 20 sec
  • Scaled to 1 TB data:
    • In-memory: 5-7 sec
    • On-disk: 170 sec

Real-world usage: Netflix uses similar Spark patterns to analyze billions of events from their streaming platform in real-time, caching frequently accessed user behavior data to provide instant recommendations.


RDD Fault Tolerance

Lineage Tracking

RDDs maintain lineage information that can be used to reconstruct lost partitions

Example

messages = textFile(...).filter(_.startsWith("ERROR"))
                       .map(_.split('\t')(2))

Lineage graph:

[HDFS File] --filter--> [Filtered RDD] --map--> [Mapped RDD]
            (func = _.contains(...))   (func = _.split(...))

If any partition is lost, Spark can recreate it by re-executing only the necessary transformations.

Analogy: It's like having a trail of breadcrumbs showing how you got somewhere. If you get lost, you can follow the breadcrumbs back and retrace your steps.


Example 2: Logistic Regression

Goal

Find the best line separating two sets of points (binary classification)

Visualization

  • Start with random initial line
  • Iteratively adjust to find optimal separation
  • Target: achieve best classification boundary

Code Implementation

val data = spark.textFile(...).map(readPoint).cache()
var w = Vector.random(D)
 
for (i <- 1 to ITERATIONS) {
  val gradient = data.map(p =>
    (1 / (1 + exp(-p.y*(w dot p.x))) - 1) * p.y * p.x
  ).reduce(_ + _)
  w -= gradient
}
 
println("Final w: " + w)

Mathematical Formula

gradient=∑p∈data(11+e−y⋅(w⋅x)−1)⋅y⋅x\text{gradient} = \sum_{p \in \text{data}} \left(\frac{1}{1 + e^{-y \cdot (w \cdot x)}} - 1\right) \cdot y \cdot x

Performance Comparison

IterationsHadoop TimeSpark Time
1~174s~174s
5~635s~174s
10~1270s~174s
20~2540s~174s
30~3810s~174s

Key insight:

  • Hadoop: 127 seconds per iteration (reads from disk each time)
  • Spark first iteration: 174 seconds (initial data loading)
  • Spark subsequent iterations: ~6 seconds (data cached in memory)

Speedup: ~20-25x for iterative algorithms

Real-world usage: Uber uses Spark's machine learning capabilities to train models that predict rider demand, driver availability, and optimize pricing in real-time across millions of rides daily.


Apache Spark Overview

Definition

Apache Spark is an open-source, distributed computing system designed for big data processing and analytics. It provides fast, scalable, and fault-tolerant computations by enabling parallel processing across clusters.

Key Features

Lightning-Fast

  • 100x faster than Hadoop MapReduce for in-memory computations
  • Optimized for iterative algorithms

Scalability

  • Can run on clusters of thousands of machines
  • Handles petabytes of data

Fault Tolerance

  • Automatically recovers lost computations using lineage information
  • No need for data replication

Multi-Language Support

  • Works with Scala, Python (PySpark), Java, and R
  • Each language has full API support

Unified Analytics

  • Supports batch processing, streaming, machine learning, and graph analytics
  • Single platform for multiple workloads

Core Components of Apache Spark

1. Spark Core

  • Provides the foundation for distributed execution and memory management
  • Implements RDD (Resilient Distributed Dataset), the fundamental data structure in Spark

2. Spark SQL

  • Allows structured data processing using SQL queries
  • Supports DataFrames and Datasets, enabling efficient optimizations
  • Can query data from various sources (Parquet, JSON, Hive, etc.)

Real-world usage: Airbnb uses Spark SQL to run thousands of SQL queries daily on petabytes of data to analyze booking patterns, pricing strategies, and user behavior.

3. Spark Streaming

  • Enables real-time processing of data streams
  • Supports data sources like Kafka, Flume, and Amazon Kinesis
  • Micro-batch processing for low-latency analytics

Real-world usage: Twitter uses Spark Streaming to analyze trending topics in real-time, processing millions of tweets per second to identify viral content and emerging trends.

4. MLlib (Machine Learning Library)

  • A scalable machine learning library for:
    • Classification (e.g., spam detection)
    • Clustering (e.g., customer segmentation)
    • Regression (e.g., price prediction)
    • Collaborative filtering (e.g., recommendations)

Real-world usage: Spotify uses MLlib to power their recommendation engine, analyzing billions of song plays to suggest personalized playlists to users.

5. GraphX (Graph Processing)

  • Used for graph-based computations (e.g., social network analysis)
  • Supports algorithms like:
    • PageRank - Ranking importance of nodes
    • Connected Components - Finding clusters
    • Shortest Paths - Finding optimal routes

Real-world usage: LinkedIn uses GraphX to analyze their professional network graph, suggesting connections, job opportunities, and content based on relationship patterns.


Deployment Modes

Spark can be deployed in multiple ways:

1. Standalone Mode

  • Runs on a single machine or cluster
  • Simple setup for testing and development

2. YARN (Hadoop Cluster Mode)

  • Uses Hadoop's resource manager
  • Integrates with existing Hadoop infrastructure

3. Mesos

  • Integrates with Apache Mesos
  • Fine-grained resource sharing

4. Kubernetes

  • Cloud-based orchestration
  • Modern containerized deployments
  • Auto-scaling capabilities

5. Cloud (AWS, Azure, GCP)

  • Fully managed Spark services
  • Examples:
    • AWS EMR (Elastic MapReduce)
    • Azure HDInsight
    • Google Cloud Dataproc
    • Databricks (cloud-native Spark platform)

Real-world usage:

  • Netflix runs Spark on AWS EMR to process over 1 trillion events per day
  • Apple uses Databricks on Azure for machine learning workloads
  • Lyft runs Spark on Kubernetes for their ride-sharing analytics

Real-World Applications

1. In-memory data mining on Hive data (Conviva)

Use case: Aggregations on many keys with same WHERE clause

Performance:

  • Hive: 20 hours
  • Spark: 0.5 hours

40× speedup comes from:

  • Not re-reading unused columns or filtered records
  • Avoiding repeated decompression
  • In-memory storage of deserialized objects

2. Predictive Analytics (Quantifind)

  • Real-time anomaly detection
  • Financial risk modeling

3. City Traffic Prediction (Mobile Millennium)

  • Analyzing GPS data from smartphones
  • Predicting traffic patterns in real-time

4. Twitter Spam Classification (Monarch)

  • Machine learning on streaming tweets
  • Real-time spam detection

5. Collaborative Filtering via Matrix Factorization

  • Recommendation systems
  • User-item preference prediction

Frameworks Built on Spark

Bagel (Pregel on Spark)

  • Google message passing model for graph computation
  • Only 200 lines of code
  • Demonstrates Spark's flexibility

Shark (Hive on Spark)

  • 3000 lines of code
  • Compatible with Apache Hive
  • ML operators in Scala
  • Much faster than Hive on MapReduce

Analogy: Think of Spark Core as a powerful engine, and these frameworks (Bagel, Shark) as different vehicles built on that engine - each optimized for specific use cases.


Implementation Details

Architecture

┌─────────────────────────────────────────┐
│  Spark    │  Hadoop   │   MPI    │ ... │
├─────────────────────────────────────────┤
│            Apache Mesos                  │
├─────────────────────────────────────────┤
│  Node  │  Node  │  Node  │  Node  │...│
└─────────────────────────────────────────┘

Key Implementation Features

  • Runs on Apache Mesos to share resources with Hadoop & other apps
  • Can read from any Hadoop input source (e.g., HDFS)
  • No changes to Scala compiler required

Spark Scheduler

Features

1. Dryad-like DAGs

  • Creates execution plans as Directed Acyclic Graphs
  • Optimizes the execution order

2. Pipelines functions within a stage

  • Combines multiple transformations
  • Reduces intermediate data materialization

3. Cache-aware work reuse & locality

  • Reuses cached data when possible
  • Schedules tasks close to data

4. Partitioning-aware

  • Avoids unnecessary shuffles
  • Optimizes join operations

Example DAG

    Stage 1          Stage 2          Stage 3
    
A: ─┐                              
    ├─► join ──┐                   
B: ─┘          │                   
               ├─► union ─► groupBy ──► map
C: ─┐          │                   
    ├─► join ──┘                   
D: ─┘                              

E: ────────────────────────────────────────►

F: ────────────────────────────────────────►

G: (cached) ═══════════════════════════════►

Legend:
─── = RDD lineage
═══ = cached partition

Interactive Spark

Scala Interpreter Integration

Spark provides a modified Scala interpreter to allow Spark to be used interactively from the command line.

Required Changes

  1. Modified wrapper code generation
    • Each line typed has references to objects for its dependencies
  2. Distribute generated classes over the network
    • Classes are sent to worker nodes automatically

Usage

$ spark-shell
Welcome to Spark Shell!
 
scala> val data = sc.textFile("hdfs://...")
scala> val filtered = data.filter(_.contains("error"))
scala> filtered.count()
res0: Long = 12345

Analogy: Think of Interactive Spark like a Python Jupyter notebook, but for big data - you can experiment with massive datasets interactively, seeing results immediately.


Behavior with Limited RAM

Performance vs Memory

Experiment: Running iterative algorithm with different amounts of cached data

Memory CachedIteration Time (s)
Cache disabled68.8
25% cached58.1
50% cached40.7
75% cached29.7
Fully cached11.5

Key Insight

Even partial caching provides significant speedup:

  • 25% cached: 15% faster than no cache
  • 50% cached: 41% faster than no cache
  • Fully cached: 83% faster than no cache

Important: Spark gracefully degrades when memory is insufficient - it doesn't fail, just runs slower by recomputing or spilling to disk.


Fault Recovery Results

Experiment Setup

  • Run 10 iterations of an algorithm
  • Inject a failure in iteration 6
  • Measure recovery time

Results

IterationNo Failure (s)With Failure (s)
1119119
25757
35656
45858
55858
65781
75757
85959
95757
105959

Analysis

  • Iteration 6 with failure: 81 seconds (24 seconds overhead)
  • Subsequent iterations: Return to normal speed
  • Only lost partition is recomputed - not entire dataset
  • Automatic recovery without manual intervention

Real-world importance: In production clusters with thousands of nodes, failures are common. Spark's fast recovery ensures applications continue running with minimal disruption.


Spark Operations Reference

Transformations (Define a new RDD)

OperationDescription
mapApply function to each element
filterSelect elements matching predicate
sampleRandom sample of elements
groupByKeyGroup values by key
reduceByKeyReduce values by key
sortByKeySort RDD by key
flatMapMap then flatten results
unionCombine two RDDs
joinJoin two RDDs by key
cogroupGroup two RDDs by key
crossCartesian product
mapValuesApply function to values only

Actions (Return a result to driver)

OperationDescription
collectReturn all elements
reduceAggregate elements
countCount elements
saveSave to storage
lookupKeyReturn value for key

Practical Examples

1. Log Processing (Real-time Log Analysis)

Use Case: Analyze server logs to find the most common error messages

from pyspark.sql import SparkSession
 
# Initialize Spark session
spark = SparkSession.builder.appName("LogAnalysis").getOrCreate()
sc = spark.sparkContext
 
# Load log file into RDD
logs_rdd = sc.textFile("server_logs.txt")
 
# Filter error messages
error_logs = logs_rdd.filter(lambda line: "ERROR" in line)
 
# Count occurrences of each error type
error_counts = error_logs.map(lambda line: (line.split(" ")[1], 1)) \
                        .reduceByKey(lambda a, b: a + b)
 
# Collect and display results
print(error_counts.collect())

Output Example:

[("ConnectionTimeout", 120), ("DatabaseDown", 45), ("DiskFull", 30)]

Real-world usage: Datadog and Splunk use similar Spark patterns to analyze billions of log events per day for their customers, providing real-time alerting and analytics.


2. Sentiment Analysis (Text Processing)

Use Case: Analyze tweets and count positive vs. negative words

# Sample word lists
positive_words = {"good", "happy", "love", "awesome", "excellent"}
negative_words = {"bad", "sad", "hate", "terrible", "worst"}
 
# Load tweets into RDD
tweets_rdd = sc.textFile("tweets.txt")
 
# Analyze sentiment
sentiment_counts = tweets_rdd.flatMap(lambda line: line.split()) \
    .map(lambda word: ("positive", 1) if word.lower() in positive_words 
         else ("negative", 1) if word.lower() in negative_words else None) \
    .filter(lambda x: x is not None) \
    .reduceByKey(lambda a, b: a + b)
 
print(sentiment_counts.collect())

Output Example:

[("positive", 300), ("negative", 120)]

Real-world usage: Brand monitoring companies like Brandwatch use Spark to analyze millions of social media posts in real-time, helping companies understand public sentiment about their products and services.


Advantages of Apache Spark

1. ⚡ Speed

  • Processes petabytes of data 100x faster than Hadoop
  • In-memory processing eliminates disk I/O bottlenecks

2. 🔄 Flexibility

  • Works with batch, real-time, and machine learning workloads
  • Unified platform reduces complexity

3. 👨‍💻 Ease of Use

  • Supports SQL, Python, Scala, and Java
  • High-level APIs simplify development

4. 📊 Scalability

  • Runs on thousands of nodes
  • Handles growing data volumes seamlessly

Summary: RDD Applications Across Domains

RDDs are widely used in various domains:

🔍 Big Data Analytics

  • Log Processing
  • Sentiment Analysis
  • Clickstream Analysis

💰 Finance & Security

  • Fraud Detection
  • Risk Modeling
  • Network Traffic Analysis
  • Anomaly Detection

🏥 Healthcare & Bioinformatics

  • DNA Sequence Analysis
  • Medical Image Processing
  • IoT Health Monitoring

🎯 Recommendation Systems & AI

  • Collaborative Filtering
  • Content Recommendations
  • Personalization Engines

🌐 Real-Time Applications

  • Stream Processing
  • Live Dashboard Updates
  • Real-time Alerting

DryadLINQ, FlumeJava

  • Similar "distributed collection" API
  • Cannot reuse datasets efficiently across queries (Spark's advantage)

Relational Databases

  • Use lineage/provenance, logical logging, materialized views
  • Spark extends these concepts to big data scale

GraphLab, Piccolo, BigTable, RAMCloud

  • Fine-grained writes similar to distributed shared memory
  • Different approach than Spark's coarse-grained transformations

Iterative MapReduce (e.g., Twister, HaLoop)

  • Implicit data sharing for a fixed computation pattern
  • Less flexible than Spark's explicit caching

Caching Systems (e.g., Nectar)

  • Store data in files
  • No explicit control over what is cached (Spark provides explicit cache control)

Key Takeaways

  1. RDDs enable efficient reuse of data across multiple operations
  2. Lineage-based fault tolerance eliminates need for data replication
  3. In-memory computing provides orders of magnitude speedup
  4. Lazy evaluation optimizes execution plans
  5. Unified platform supports multiple workloads (batch, streaming, ML, graph)
  6. Production-ready - used by thousands of companies worldwide

Conclusion

Spark provides a simple, efficient, and powerful programming model for a wide range of applications:

  • ✅ Fast iterative algorithms
  • ✅ Interactive data exploration
  • ✅ Real-time stream processing
  • ✅ Machine learning at scale
  • ✅ Graph analytics

Download open source release: www.spark-project.org
Contact: matei@berkeley.edu


Additional Resources

Cloud Services Using Spark

  1. Databricks (databricks.com)

    • Unified analytics platform built on Spark
    • Used by 5,000+ organizations including Shell, Comcast, H&M
  2. AWS EMR (Amazon Elastic MapReduce)

    • Managed Spark clusters on AWS
    • Used by Netflix, Airbnb, Samsung
  3. Azure HDInsight

    • Microsoft's cloud Spark offering
    • Integrated with Azure ecosystem
  4. Google Cloud Dataproc

    • Fast, easy-to-use managed Spark service
    • Auto-scaling and integrated with GCP
  5. Cloudera Data Platform

    • Enterprise data platform with Spark
    • On-premises and cloud deployment

Production Use Cases by Company

  • Netflix: 1 trillion events/day for recommendations and content optimization
  • Uber: Demand prediction, driver matching, pricing optimization
  • Airbnb: Pricing algorithms, search ranking, fraud detection
  • Apple: Siri improvements, App Store analytics
  • LinkedIn: Profile recommendations, job matching, spam detection
  • Pinterest: Content recommendations, visual search
  • Spotify: Music recommendations, playlist generation
  • Twitter: Real-time trend analysis, spam detection

Formula Reference

Logistic Regression Gradient

gradient=∑p∈data(11+e−y⋅(w⋅x)−1)⋅y⋅x\boxed{\text{gradient} = \sum_{p \in \text{data}} \left(\frac{1}{1 + e^{-y \cdot (w \cdot x)}} - 1\right) \cdot y \cdot x}

Where:

  • ww = weight vector
  • xx = feature vector
  • yy = label (+1 or -1)
  • The term 11+e−y⋅(w⋅x)\frac{1}{1 + e^{-y \cdot (w \cdot x)}} is the sigmoid function

End of Chapter 4 Notes