Big Data Technologies
Spark Architecture, RDDs, DataFrames, Spark SQL and Machine Learning
PGCP-BDA
Apache Spark
Apache Spark is a distributed execution engine that builds task graphs for batch, SQL, machine-learning and streaming workloads across a cluster.
What is Spark?
Apache Spark is a fast, general-purpose distributed computing framework that:
- Is 100x faster than Hadoop MapReduce (in-memory processing)
- Supports batch, streaming, SQL, ML and graph processing
- Written in Scala; APIs in Python (PySpark), Java, Scala, R, SQL
Spark Architecture
+--------------------------------------------------+
| SPARK CLUSTER |
| |
| Driver Program Worker Nodes |
| +----------+ +---------+ |
| | SparkContext| ──────> | Executor| → Tasks |
| | (your app) | +---------+ |
| +----------+ +---------+ |
| ──> | Executor| → Tasks |
| Cluster Manager +---------+ |
| (YARN, K8s, Mesos) ...N workers... |
+--------------------------------------------------+
RDD (Resilient Distributed Dataset)
The fundamental data structure in Spark:
- Resilient — fault-tolerant (recomputed from lineage if lost)
- Distributed — partitioned across cluster
- Dataset — collection of data elements
Spark Transformations (lazy): map, filter, flatMap, groupByKey, reduceByKey, join
Spark Actions (trigger execution): collect, count, save, reduce, take
PySpark Examples
from pyspark.sql import SparkSession
# Create Spark session
spark = SparkSession.builder \
.appName("Sales Analysis") \
.getOrCreate()
# Read data
df = spark.read.csv("hdfs:///data/sales.csv", header=True, inferSchema=True)
# DataFrame operations (lazy)
result = df \
.filter(df.amount > 1000) \
.groupBy("product") \
.agg({"amount": "sum"}) \
.orderBy("sum(amount)", ascending=False) \
.limit(10)
# Action: actually execute the computation
result.show()
# SQL style
df.createOrReplaceTempView("sales")
spark.sql("""
SELECT product, SUM(amount) as total_sales
FROM sales
WHERE amount > 1000
GROUP BY product
ORDER BY total_sales DESC
LIMIT 10
""").show()
# Spark Streaming
from pyspark.streaming import StreamingContext
ssc = StreamingContext(spark.sparkContext, 1) # 1 second batch interval
driver executor and cluster manager
The Spark driver plans work, the cluster manager allocates resources and executors run tasks and retain application data.
RDD
A resilient distributed dataset is an immutable partitioned collection described by lineage so lost partitions can be recomputed.
lineage and lazy evaluation
Spark records transformations as lineage and delays execution until an action requires a result.
transformation and action
A Spark transformation lazily defines a new dataset, while an action triggers execution and returns or writes a result.
narrow and wide transformation
A narrow transformation reads each output partition from few input partitions; a wide transformation requires a shuffle across partitions.
DataFrame and Catalyst
A Spark DataFrame is a distributed table with named typed columns; Catalyst analyzes and optimizes its logical plan before physical execution.
Spark SQL
Spark SQL analyzes SQL or DataFrame expressions into optimized logical and physical plans executed as distributed stages.
partition and cache
Partitions define Spark task parallelism and data placement, while caching retains reused partitions to avoid recomputation.
Spark ML pipeline
A Spark ML pipeline chains estimators and transformers so fitting produces a reproducible model pipeline that applies identical feature steps during.
Spark Execution Model
A Spark application contains a driver and executors. The driver builds a logical plan and schedules jobs. Executors run tasks and retain cached partitions. An action triggers execution of lazy transformations. Narrow transformations can process each output partition from one input partition while wide transformations require a shuffle across the cluster. Stages are separated at shuffle boundaries.
An RDD is an immutable partitioned collection with lineage used for recomputation. A DataFrame adds named columns and schema. Spark SQL can optimize DataFrame and SQL expressions through logical and physical planning. DataFrames normally outperform opaque row-by-row functions because the engine can prune columns, push filters and choose join strategies. Caching helps only when reused data justifies memory and recomputation cost.
Spark SQL and Machine Learning
Spark SQL reads tables and file formats through a common structured interface. Temporary views belong to a session or application while catalog tables preserve metadata. Partitioning and bucketing affect file layout. Broadcast joins avoid shuffling a small relation. Skewed keys and many tiny partitions produce slow tasks even when aggregate capacity is high.
Spark ML pipelines connect transformers and estimators. A transformer converts one DataFrame to another. An estimator learns a model through fit and produces a transformer. Feature preparation, training and evaluation must use separate data to avoid leakage. Parameters and fitted stages can be assembled into a repeatable pipeline. Distributed execution helps large feature sets but does not correct poor sampling, invalid labels or an unsuitable metric.
Continue learning
Related notes
Put this topic into timed practice
Open mock tests when you want full-exam pacing, or keep drilling in practice mode.