Big Data and AnalyticsUnit 68 min read

Apache Spark: Architecture, Components, and Applications

Unit 6 of Big Data and Analytics explores Apache Spark, its architecture, core components (RDDs, DataFrames, Spark SQL), execution model, and real-world applications in analytics, machine learning, and stream processing. This note covers how Spark differs from Hadoop MapReduce, its distributed computing model, and hand

Apache Spark: The Engine for Big Data Processing

What is Apache Spark?

Apache Spark is an open-source, distributed computing framework designed for large-scale data processing. It provides an interface for programming entire clusters with implicit data parallelism and fault tolerance. Unlike Hadoop MapReduce, Spark keeps data in memory, making it 100x faster for iterative algorithms and interactive queries.

Key Features

  • In-memory processing: Reduces disk I/O bottlenecks.
  • Unified engine: Supports SQL, streaming, machine learning, and graph processing.
  • Fault tolerance: Recovers lost tasks via lineage (DAG-based execution).
  • Multi-language support: Scala, Python (PySpark), Java, R, and SQL.

Core Components of Spark

Spark’s architecture consists of:

  1. Spark Core: The foundation for distributed computing (RDDs, scheduling, memory management).
  2. Spark SQL: Structured data processing (DataFrames, Datasets).
  3. Spark Streaming: Real-time data processing.
  4. MLlib: Machine learning library.
  5. GraphX: Graph processing.

Visual: Spark Architecture

graph TD
    subgraph "Driver Node"
        DP["Driver Program"]
        SS["Spark Context / Session"]
    end
    subgraph "Cluster Manager"
        CM["Standalone / YARN / Mesos / K8s"]
    end
    subgraph "Worker Nodes"
        WN1["Worker Node 1"]
        WN2["Worker Node 2"]
        WN3["Worker Node N"]
    end
    subgraph "Executors"
        E1["Executor 1"]
        E2["Executor 2"]
        E3["Executor N"]
    end
    DP --> SS
    SS --> CM
    CM --> WN1
    CM --> WN2
    CM --> WN3
    WN1 --> E1
    WN2 --> E2
    WN3 --> E3
    E1 -->|"Shuffle Data"| E2
    E2 -->|"Shuffle Data"| E3

1. Spark Core: RDDs (Resilient Distributed Datasets)

RDDs are immutable, partitioned collections of objects that can be processed in parallel. They are the fundamental data structure in Spark.

Data BlocksPartition 1Data BlocksPartition 2Data BlocksPartition NRDD (Resilient Distributed Dataset)
Structure of an RDD: A logical collection of partitions distributed across the cluster.

How RDDs Work

  • Partitioning: Data is split into chunks for parallel processing.
  • Lineage: Tracks transformations to rebuild lost data (fault tolerance).
  • Lazy evaluation: Operations are only executed when an action (e.g., collect(), count()) is called.

Example: Word Count in Spark (vs. MapReduce)

# PySpark example: Word count on a text file
from pyspark import SparkContext

sc = SparkContext("local", "WordCount")
text_file = sc.textFile("book.txt")
words = text_file.flatMap(lambda line: line.split())
word_counts = words.countByValue()
print(word_counts)

Why Spark is faster?

  • MapReduce reads/writes data to disk after each stage.
  • Spark caches intermediate RDDs in memory.

2. Spark SQL and DataFrames

Spark SQL enables structured data processing using SQL or DataFrame APIs. DataFrames are optimized execution plans (like RDDs but with schema).

02.557.510RDD (Unstructured)1DataFrame (Structured)10
Relative Performance: DataFrames are significantly faster than RDDs due to Catalyst Optimizer and Tungsten Execution Engine.

DataFrame vs. RDD

Feature RDD DataFrame
Schema No (raw data) Yes (columnar format)
Performance Slower (Java serialization) Faster (Tungsten engine)
Use Case Low-level transformations SQL queries, analytics

Example: Filtering Data in a DataFrame

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("SQLExample").getOrCreate()
df = spark.read.csv("sales_data.csv", header=True, inferSchema=True)
filtered_df = df.filter(df["amount"] > 1000).groupBy("customer_id").count()
filtered_df.show()

3. Spark Streaming: Real-Time Processing

Spark Streaming processes live data streams (e.g., IoT sensors, clickstreams) in micro-batches (DStreams).

Example: Fraud Detection in eSewa Transactions

sequenceDiagram
    participant User
    participant eSewaAPI
    participant SparkStreaming
    participant MLModel

    User->>eSewaAPI: Transaction (e.g., Rs. 50,000 in 1 sec)
    eSewaAPI->>SparkStreaming: Stream data to Spark
    SparkStreaming->>MLModel: Check for anomalies
    MLModel-->>SparkStreaming: Flag as fraud (95% confidence)
    SparkStreaming->>eSewaAPI: Block transaction

Real-World Use Case:

  • eSewa uses Spark Streaming to detect unusual transaction patterns (e.g., multiple high-value payments in seconds) and block fraudulent activities in real time.

4. MLlib: Machine Learning on Big Data

MLlib provides scalable machine learning algorithms (classification, regression, clustering).

Example: Customer Segmentation in Daraz

from pyspark.ml.clustering import KMeans

# Load data
df = spark.read.csv("customer_purchases.csv", header=True)

# Train K-Means model
kmeans = KMeans().setK(5).setSeed(1)
model = kmeans.fit(df)
predictions = model.transform(df)
predictions.show()

Output: Daraz can group customers by spending habits (e.g., "High-value buyers," "Occasional shoppers") for targeted promotions.


5. GraphX: Graph Processing

GraphX is Spark’s library for graph algorithms (PageRank, community detection).

Example: Traffic Route Optimization in Pathao

graph TD
    A["Pathao Driver"] -->|"Requests route"| B["GraphX"]
    B -->|"Calculates shortest path"| C["Traffic Data"]
    C -->|"Real-time updates"| B
    B -->|"Optimized route"| A

How it works:

  • GraphX processes road networks as graphs (nodes = intersections, edges = roads).
  • Pathao uses real-time traffic data (from NTC or GPS) to reroute drivers dynamically.

In the Real World

  1. eSewa (Fraud Detection)

    • Uses Spark Streaming + MLlib to analyze transaction patterns in real time.
    • Example: If a user suddenly transfers Rs. 100,000 to 5 different accounts in 2 minutes, Spark flags it as suspicious and blocks the transaction.
  2. Pathao (Dynamic Routing)

    • GraphX processes live traffic data to suggest the fastest route.
    • Example: During Kathmandu’s busy hours, Pathao reroutes drivers via less congested streets (e.g., avoiding Thapathali to Mahaboudha).
  3. Ncell (Churn Prediction)

    • Uses Spark MLlib to predict which customers might cancel their plans.
    • Example: If a customer reduces call duration by 30% in a week, Ncell offers a discount to retain them.

Advantages and Disadvantages of Spark

Advantages Disadvantages
100x faster than Hadoop MapReduce Requires more memory than disk-based systems
Unified engine for SQL, ML, streaming Smaller community than Hadoop
Fault-tolerant via lineage Steeper learning curve for beginners
Supports multiple languages Not ideal for batch-only workloads

Exam Tip

  1. Compare Spark vs. Hadoop MapReduce:

    • Spark: In-memory, faster iterations, unified API.
    • MapReduce: Disk-based, batch-only, slower for iterative tasks.
  2. Key Terms to Remember:

    • RDD: Resilient Distributed Dataset (immutable, fault-tolerant).
    • DataFrame: Schema-enforced, optimized for SQL.
    • DAG: Directed Acyclic Graph (Spark’s execution plan).
    • Micro-batch: Spark Streaming’s approach to real-time processing.
  3. Worked Example Focus:

    • Always show how Spark processes data differently (e.g., caching vs. disk I/O).
    • Relate to Nepalese companies: eSewa (fraud), Pathao (routing), Ncell (churn).
  4. Practical Questions:

    • "How would you optimize a word count job in Spark vs. MapReduce?" Answer: Use persist() to cache intermediate RDDs in Spark.
    • "Explain how Spark Streaming detects fraud in eSewa." Answer: Uses DStreams + MLlib to analyze transaction velocity and amount anomalies.

Based on the TU BIM syllabus for Big Data and Analytics (IT278), unit 6.

Discussion

Loading…