Big Data and AnalyticsUnit 68 min read
Apache Spark: Architecture, Components & Big Data Processing
Unit 6 of Big Data and Analytics covers Apache Spark’s architecture, core components (Spark Core, Spark SQL, MLlib, GraphX, and Spark Streaming), its in-memory processing model, and how it outperforms Hadoop MapReduce for iterative and real-time analytics. Includes real-world use cases, performance comparisons, and han
What is Apache Spark?
Apache Spark is an open-source, distributed processing framework designed for large-scale data processing across clusters. Unlike Hadoop MapReduce, which relies on disk I/O for intermediate results, Spark uses in-memory computation, making it 100x faster for iterative algorithms and interactive queries.
Why Spark?
- Speed: In-memory processing reduces disk I/O bottlenecks.
- Ease of Use: Supports scala, Python, R, and Java APIs.
- Versatility: Handles batch processing, real-time analytics, machine learning, and graph processing.
- Scalability: Runs on Hadoop, Kubernetes, or standalone clusters.
Core Components of Apache Spark
Spark’s modular architecture includes:
| Component | Purpose | Key Features |
|---|---|---|
| Spark Core | Distributed task scheduling and memory management. | RDDs (Resilient Distributed Datasets), DAG execution engine. |
| Spark SQL | Structured data processing (SQL, DataFrames). | Schema enforcement, Catalyst optimizer, Pandas-like APIs. |
| MLlib | Machine learning library. | Algorithms for classification, clustering, regression (e.g., logistic regression). |
| GraphX | Graph processing (e.g., social networks, fraud detection). | Graph algorithms (PageRank, connected components). |
| Spark Streaming | Real-time data processing (micro-batch). | DStreams (Discretized Streams), Kafka integration. |
How Spark Works: In-Memory Processing
1. Resilient Distributed Datasets (RDDs)
- Immutable, partitioned collections of data stored in memory.
- Two key properties:
- Partition Tolerance: Data split across nodes.
- Fault Tolerance: Lineage (history of transformations) enables recovery.
- Example:
# Creating an RDD from a text file rdd = sc.textFile("data.txt") # Transformations (lazy evaluation) filtered_rdd = rdd.filter(lambda line: "error" in line) # Actions (trigger computation) filtered_rdd.count()
2. DAG (Directed Acyclic Graph) Execution Engine
- Spark breaks jobs into stages (wide vs. narrow transformations).
- Wide transformations (e.g.,
join,groupBy) require shuffling data across nodes. - Narrow transformations (e.g.,
map,filter) can be pipelined.
Spark vs. Hadoop MapReduce
| Feature | Apache Spark | Hadoop MapReduce |
|---|---|---|
| Processing Model | In-memory (RAM) | Disk-based (HDFS) |
| Speed | 10–100x faster for iterative tasks | Slower due to disk I/O |
| Use Case | Real-time analytics, ML, interactive SQL | Batch processing (ETL, logs) |
| Ease of Use | APIs in Scala/Python/R | Java-centric, complex APIs |
| Fault Tolerance | RDD lineage (recomputes lost data) | Task retries (slower recovery) |
Real-World Applications of Spark
1. Fraud Detection in Banks (e.g., Nabil Bank, Global IME)
- How Spark is used:
- Spark MLlib trains models on transaction data to detect anomalies in real time.
- Spark Streaming processes live transactions (e.g., sudden large withdrawals).
- Example:
A bank uses Spark to flag suspicious transactions in <100ms using:
from pyspark.ml.classification import RandomForestClassifier model = RandomForestClassifier(featuresCol="features", labelCol="is_fraud") model.fit(training_data)
2. Traffic Optimization in Kathmandu (NTC, Pathao)
- How Spark is used:
- Spark SQL analyzes GPS data from Pathao rides to predict congestion.
- GraphX models road networks to optimize routes dynamically.
- Example: NTC uses Spark to process 1M+ GPS points/day to reroute buses during peak hours.
3. Recommendation Systems (Daraz, Amazon)
- How Spark is used:
- Spark MLlib builds collaborative filtering models for product recommendations.
- Spark Streaming updates recommendations in real time based on user clicks.
Hands-On Example: Analyzing Kathmandu Traffic Data
Scenario: NTC wants to predict traffic jams using historical GPS data from Pathao rides.
Step 1: Load Data into Spark
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("TrafficAnalysis").getOrCreate()
df = spark.read.parquet("pathao_gps_data.parquet")
Step 2: Preprocess Data (Spark SQL)
-- Filter data for peak hours (7–9 AM)
SELECT *
FROM df
WHERE hour >= 7 AND hour <= 9;
Step 3: Train a Model (MLlib)
from pyspark.ml.feature import VectorAssembler
from pyspark.ml.regression import LinearRegression
# Feature engineering
assembler = VectorAssembler(inputCols=["speed", "distance"], outputCol="features")
df_features = assembler.transform(df)
# Train model
lr = LinearRegression(featuresCol="features", labelCol="congestion_level")
model = lr.fit(df_features)
Step 4: Predict Congestion
predictions = model.transform(df_features)
predictions.filter(predictions.congestion_level > 0.8).show()
Spark Architecture: Cluster Deployment
Spark runs on a master-slave architecture with:
- Driver Program: Coordinates tasks (submits jobs to the cluster).
- Cluster Manager: Allocates resources (YARN, Mesos, or Spark’s standalone).
- Executor: Runs tasks on worker nodes (caches data in memory).
Advantages and Disadvantages of Spark
| Advantages | Disadvantages |
|---|---|
| Speed: In-memory processing. | Memory Overhead: Requires large RAM. |
| Ease of Use: APIs for multiple languages. | Smaller Data Handling: Struggles with <100MB datasets (use Pandas instead). |
| Unified Engine: Supports SQL, ML, streaming. | Learning Curve: Steeper than Pandas. |
| Fault Tolerance: RDD lineage. | Not a Database: Needs HDFS/S3 for storage. |
Exam Tip
- Focus on Key Concepts:
- Explain RDDs, DAG execution, and Spark’s in-memory model clearly.
- Compare Spark with Hadoop MapReduce (speed, use cases).
- Practical Questions:
- Expect code snippets (PySpark/Scala) for RDD transformations or MLlib usage.
- Be ready to design a Spark job for a given scenario (e.g., fraud detection).
- Real-World Links:
- Connect Spark to Nepali examples (e.g., NTC traffic analysis, bank fraud detection).
- Mention tools like Zeppelin, Databricks, or Jupyter for Spark development.
Based on the TU BITM syllabus for Big Data and Analytics (IT278), unit 6.
Discussion
Loading…