Spark - Part Infinity
What is Apache Spark?
Spark is an distributed computing engine for large data processing. Spark provide in memory storage for intermediate computation, making it 10x faster than Hadoop's Map Reduce.
It is written in Scala, run on JVM.
Spark vs Hadoop MapReduce
- Spark process data in memory where as MapReduce uses disk-based processing (reads/writes HDFS between stages).
- Spark uses DAG (Directed Acyclic Graph) execution model that optimizes multi-stage pipeline. In contrast, MapReduce is a rigid map-then-reduce execution model.
- Spark provides iterative algorithms natively (e.g.: ML workload) but in MapReduce to achieve iterative processing chaining multiple jobs is required.
- Spark supports real-time processing, MapReduce is batch processing only.
Spark Ecosystem
- Spark Core: RDDs, task scheduling, memory management, fault tolerance
- Spark SQL: Dataframes, Datasets, Catalyst Optimizer, Hive Integration
- Spark Streaming / Structured Streaming: Real-time data processing
- MLib: Machine Learning library
- GraphX: Graph processing and computation
Spark Architecture & Internals
Spark follows a master-slave architecture, where driver is the master and executor are the slaves:
- Driver → Acts as the central coordinator. It creates the Spark application, builds the execution plan, and schedules tasks.
- Executors → Run on worker nodes and execute the tasks assigned by the Driver. They also store data in memory/disk when needed.
- Cluster Manager → Allocates resour
- ces (CPU and memory) and manages where the Driver and Executors run.
Driver → Cluster Manager → Executors → Tasks
Driver = coordinates | Executors = execute | Cluster Manager = manages resources
Let's go over in detail for each component
1. The Driver (Master Node)
The Driver is the process that runs your Spark application's main program. It acts as the control center of the application. It:
- Creates the
SparkSession/SparkContext - Builds the execution plan (DAG)
- Converts the plan into stages and tasks
- Schedules tasks on Executors
- Coordinates the overall execution
2. The Cluster Manager
SparkSession vs SparkContext
SparkContext was the entry point to connect spark application to the cluster in Spark 1.x, from Spark 2.x a new unified entry point for SparkContext, HiveContext, SQLContext was introduced that is called SparkSession.
Internally SparkSession is a wrapper for SparkContext.
Code for initializing SparkSession (in pyspark)
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName('Spark Application') \
.master('local[*]') \
.getOrCreate()
RDDs vs Dataframes vs Datasets
- RDDs are the low level abstraction whereas Datasets and Dataframes are high level abstraction
- RDDs do not have any schema. Dataframes has schema from
StructTypeobject and Datasets has schema as well compile-time safety. - RDDs are not optimized by Catalyst and Tungsten, while Dataframes and Datasets leverage these optimizations for improved performance.
RDD transformations are often defined using lambda functions, making them opaque to Spark. As Spark cannot understand the transformation logic, it cannot apply Catalyst and Tungsten optimizations. Dataframes and Datasets, however, provide a structured API that Spark can optimize
- RDDs provide only a functional API (
map,filter). Dataframes provide a SQL-like declarative API. Datasets support both a type-safe functional API (map,filter,flatMap) and a SQL-like API, combining RDD flexibility with Dataframe optimizations. - RDDs are preferred when working with unstructured data, implementing complex custom transformations, requiring fine-grained control over partitioning and execution, or maintaining legacy Spark applications. For most ETL and analytical workloads, Dataframes are preferred because they leverage Catalyst and Tungsten optimizations, resulting in significantly better performance
In Pyspark, Dataframe is the primary abstraction of (Dataset[Row])
Transformations vs Actions
Transformations are lazy operations that define how data should be changed or processed. Examples include select(), where(), filter(), groupBy(), join(), and aggregations such as sum(), count(), and avg().
They do not execute immediately. Instead, Spark records these operations and builds a logical execution plan (DAG). This allows Spark to optimize the plan before actually processing the data.
Actions are operations that trigger the execution of the DAG and produce a result or write data to storage. Examples include collect(), count(), show(), save(), and write operations such as df.write.parquet(...) or df.write.save(...).
Partitions
Partitions are the smallest unit of parallelism in Spark. One task work on one partition & on one core. When reading a file, Spark creates 1 partition per 128MB block (HDFS block size). For example if we have 50GB file i.e., 50000MB / 128MB = ~ 391 paritions, so total tasks would be 391.
Having too few partitions will result in under utilization of compute, large task, OOM risks. In contrast, having too many partitions will cause scheduling overhead, many small files.
A rule of thumb is to target 128-256MB per partition.
Section 2: Spark Architecture and Internals
The Architecture of Spark

Cluster Components
- Driver Program: It runs the main application, creates the SparkSession, converts the code into DAG, schedules tasks and collect results.
- Cluster Manager: It manages and allocates resources for the cluster of nodes (executors) on which Spark Application runs.
- Executors: Executors or workers processes on cluster nodes that execute tasks, store cached data, report back to driver.
Note: Number of stages = Number of shuffle boundaries + 1
1. While working with large or varying workload use Dynamic Resource Allocation.
spark.dynamicAllocation.enabled true
spark.dynamicAllocation.minExecutors 2
spark.dynamicAllocation.schedulerBacklogTimeout 1m
spark.dynamicAllocation.maxExecutors 20
spark.dynamicAllocation.executorIdleTimeout 2min
By default spark.dynamicAllocation.enabled is set to false. When enabled with the settings shown here, the Spark driver will request that the cluster manager create two executors to start with, as a minimum (spark.dynamicAllocation.minExecutors). As the task queue backlog increases, new executors will be requested each time the backlog timeout (spark.dynamicAllocation.schedulerBacklogTimeout) is exceeded. In this case, whenever there are pending tasks that have not been scheduled for over 1 minute, the driver will request that a new executor be launched to schedule backlogged tasks, up to a maximum of 20 (spark.dynamicAllocation.maxExecutors). By contrast, if an executor finishes a task and is idle for 2 minutes
(spark.dynamicAllocation.executorIdleTimeout), the Spark driver will terminate
it.
Data Types In Pyspark
Structured Data Type - Basic

Structured Data Type - Complex

Window Function In Spark
from pyspark.sql.function import Window, col, desc
window_fn = Window.partitionBy(col("user_id")).orderBy(desc("login_time"))
df2 = df.withColumn("rn", rank().over(window_fn))
Four Phases of Execution Plan = https://oreil.ly/jMDOi
Logical Optimization
├─ Rule Based Optimizations
│ - Predicate Pushdown
│ - Column Pruning
│ - Constant Folding
│
└─ Statistics/CBO Assisted Decisions
- Join Reordering
↓
Optimized Logical Plan
↓
Physical Planning
Candidate Physical Plans:
- Broadcast Hash Join
- Sort Merge Join
- Shuffle Hash Join
↓
Choose Best Physical Plan
↓
Execution
Lineage Graph