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:
- Spark Core: The foundation for distributed computing (RDDs, scheduling, memory management).
- Spark SQL: Structured data processing (DataFrames, Datasets).
- Spark Streaming: Real-time data processing.
- MLlib: Machine learning library.
- 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"| E31. 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.
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).
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 transactionReal-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"| AHow 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
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.
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).
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
Compare Spark vs. Hadoop MapReduce:
- Spark: In-memory, faster iterations, unified API.
- MapReduce: Disk-based, batch-only, slower for iterative tasks.
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.
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).
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.
- "How would you optimize a word count job in Spark vs. MapReduce?"
Answer: Use
Based on the TU BIM syllabus for Big Data and Analytics (IT278), unit 6.
Discussion
Loading…