Spark’s Internal - Core Concepts
Spark’s architecture, execution flow, and core concepts like jobs, stages, tasks, partitions, and DAG

What is Spark: Spark is a distributed data processing framework at large scale.
Backbone: Spark is built around the concepts of Resilient Distributed Datasets (RDD) and Direct Acyclic Graph (DAG) representing transformations and dependencies between them.

Spark Timeline
· 2007- Dryad paper (published by Microsoft)
· 2009- Spark is founded at UC Berkeley
· 2010- Spark is open sourced
· 2013- Spark goes Apache
· 2015- spark 1.4
· 2024-spark 3.5
Why do we need RDD in Spark?
DSM (Distributed Shared Memory) is a very general abstraction, but this generality makes it harder to implement in an efficient and fault tolerant manner on commodity clusters. Here the need of RDD comes into the picture.
DSM is a mechanism that manages memory across multiple nodes and makes inter-process communications transparent to end-users.

The key motivations behind the concept of RDD are-
· Iterative algorithms.
· Interactive data mining tools.
· DSM (Distributed Shared Memory) is a very general abstraction, but this generality makes it harder to implement in an efficient and fault tolerant manner on commodity clusters. Here the need of RDD comes into the picture.
· In distributed computing system data is stored in intermediate stable distributed store such as HDFS or Amazon S3. This makes the computation of job slower since it involves many IO operations, replications, and serializations in the process.
RDD – Resilient Distributed Datasets
RDD is defined as a “collection of elements partitioned across the nodes of the cluster that can be operated on in parallel”.
How would it be visualized in distributed env?


It is the fundamental data structure of Apache Spark. RDDs are huge collections of records with following properties –
Partitioned
Persisted
Fault tolerant
Lazily evaluated
Immutable
Created by coarse grained operations
PPFLIC
Let’s try to understand these characteristics –
Immutability and partitioning
RDDs composed of collection of records which are partitioned. Partition is basic unit of parallelism in a RDD, and each partition is one logical division of data which is immutable and created through some transformations on existing partitions. Immutability helps to achieve consistency in computations.
Users can define their own criteria for partitioning based on keys on which they want to join multiple datasets if needed.
Coarse grained operations: Coarse grained operations are operations which are applied to all elements in datasets. For example – a map, or filter or groupBy operation which will be performed on all elements in a partition of RDD.
Transformations and actions
RDDs can only be created by reading data from a stable storage such as HDFS or by transformations on existing RDDs. All computations on RDDs are either transformations or actions.
Fault Tolerance
Since RDDs are created over a set of transformations, it logs those transformations, rather than actual data. Graph of these transformations to produce one RDD is called as Lineage Graph.
For example –
firstRDD=spark.textFile("hdfs://...") secondRDD=firstRDD.filter(someFunction); thirdRDD = secondRDD.map(someFunction); result = thirdRDD.count()

Spark RDD Lineage Graph
In case of we lose some partition of RDD, we can replay the transformation on that partition in lineage to achieve the same computation, rather than doing data replication across multiple nodes. This characteristic is biggest benefit of RDD, because it saves a lot of efforts in data management and replication and thus achieves faster computations.
Lazy evaluations
Spark computes RDDs lazily the first time they are used in an action, so that it can pipeline transformations. So, in above example RDD will be evaluated only when count() action is invoked.
Persistence
Users can indicate which RDDs they will reuse and choose a storage strategy for them (e.g., in-memory storage or on Disk etc.)
These properties of RDDs make them useful for fast computations.
Dataframes: similar to table in mysql in relation DB.
Tuning with SparkConf -->
A class which configures and tunes spark jobs
Key/values pairs
Can be coded in the driver or at the command line.
Execution workflow recap
Here’s a quick recap on the execution workflow before digging deeper into details: user code containing RDD transformations forms Direct Acyclic Graph which is then split into stages of tasks by DAGScheduler. Stages combine tasks that don’t require shuffling/repartitioning of the data. Tasks run on workers and results then returned to the client.

Splitting DAG into Stages
Spark stages are created by breaking the RDD graph at shuffle boundaries
Machine Learning with MLlib-->
Native Machine learning framework
Classification, Regression, Clustering
Chi-Square, Correlation, Summary Stats.
Automatic Algorithm parallelization.
Pipeline API
Apache Spark:
Apache Spark is a distributed and highly scalable in-memory data analytics system, providing the ability to develop applications in Java, Scala and Python.
Apache Spark is a distributed, in-memory, parallel processing system, which needs an associated storage mechanism.
Apache Spark provides four main submodules, which are SQL, MLlib, GraphX, and Streaming.
The following figure explains how this book will address Apache Spark and its modules.

Apache Spark cluster manager in terms of the master, slave (worker), executor, and Spark client applications

Spark Context can connect to several types of cluster managers (either Spark’s own standalone cluster manager or Mesos/YARN), which allocate resources across applications. Spark cluster manager allocates resources across the worker nodes for the application. The cluster manager allocates executors across the cluster worker nodes. It copies the application jar file to the workers, and finally it allocates tasks
If the Spark master value is set as yarn-cluster, then the application can be submitted to the cluster, and then terminated. The cluster will take care of allocating resources and running tasks.
However, if the application master is submitted as yarn-client, then the application stays alive during the life cycle of processing, and requests resources from YARN.
Apache Mesos is an open source system for resource sharing across a cluster. It allows multiple frameworks to share a cluster by managing and scheduling resources.
RDD (Resilient Distributed Dataset). The main approach to work with unstructured data. Pretty similar to a distributed collection that is not always typed.
Datasets. The main approach to work with semi-structured and structured data. Typed distributed collection, type-safety at a compile time, strong typing, lambda functions.
Dataframes. It is the Dataset organized into named columns. It is conceptually equivalent to a table in a relational database or a data frame in R/Python, but with richer optimizations under the hood. Think about it as a table in a relational database.
Why Spark is fast:
Because of reducing the number of the reading/write cycle to disk and storing intermediate data in-memory Spark makes it possible
Apache DAG:
DAG is a much-recognized term in Spark. It describes all the steps through which our data is being operated. In one line “DAG is a graph denoting the sequence of operations that are being performed on the target RDD”. DAG in Apache Spark is a set of Vertices and Edges, where vertices represent the RDDs and the edges represent the Operation to be applied on RDD.
When an action is called on Spark RDD at a high level, DAG is created and is submitted to the DAG scheduler.
DAGScheduler
But before that let’s understand DAGScheduler’s two fundamental concepts: Jobs and Stages.
Jobs
A job is a top-level work item (computation). When an action is called the processing gets started and a Job is created which is then submitted to DAGScheduler to be computed.
Stages
A stage is a physical unit of execution. It is a step in a physical execution plan.
A stage is a set of parallel tasks — one task per partition (of an RDD that computes partial results of a function executed as part of a Spark job).
In other words, a Spark job is a computation with that computation sliced into stages.
There are two types of Stages:
1. ResultStage
A ResultStage is the final stage in a job that applies a function to one or many partitions of the target RDD to compute the result of an action.
2. ShuffleMapStage
ShuffleMapStage is an intermediate stage in the physical execution DAG that corresponds to a ShuffleDependency. A ShuffleMapStage may contain multiple pipelined operations, e.g. map and filter, before shuffle operation.
Workflow
Whenever an action is called over an RDD, it is submitted as an event of type DAGSchedulerEvent by the Spark Context to DAGScheduler. It is submitted as a JobSubmitted case.
The first thing done by DAGScheduler is to create a ResultStage which will provide the result of the spark job which is submitted. Now to execute the submitted job, we need to find out on which operation our RDD is based on. So backtracking begins.
In backtracking, when we find that the current operation is dependent of a shuffle operation (this is called a shuffle dependency) a new stage is created (shuffleMapStage) which is placed before the current stage (Result stage).
This new stage’s output will be the input to our ResultStage. And the backtracking continues.
If we find another shuffle operation happening then again a new shuffleMapStage will be created and will be placed before the current stage (also a ShuffleMapStage) and the newly created shuffleMapStage will provide an input to the current shuffleMapStage.
Hence all the intermediate stages will be ShuffleMapStages and the last one will always be a ResultStage.
In other words, there will only be one ResultStage and can be any number of ShuffleMapStages in between (from 0 to n depending on the number of shuffle operation in an RDD)

Why a new stage is formed when there is shuffling of data?
DAGScheduler splits up a job into a collection of stages.
Each stage contains a sequence of narrow transformations (that can be completed without shuffling the entire data set) separated at shuffle boundaries, i.e. where shuffle occurs.
Stages are thus a result of breaking the RDD graph at shuffle boundaries.
Shuffle boundaries introduce a barrier where stages/tasks must wait for the previous stage to finish before they fetch map outputs.
There are two advantages of breaking tasks into stages:
After every shuffle operation, a new stage is created so that whenever data is lost due to shuffle(network I/O) only the previous stage will be calculated for fault tolerance.
For executing operations in one go Spark groups the operation which doesn’t need to share data between executors (when one partition requires the data from another partition to complete some operation like groupBy).
So, after DAGScheduler has done its work of converting this job into stages, it hands over the stage to TaskScheduler for its execution which will do the rest of the computation.
A broadcast variable. Broadcast variables allow the programmer to keep a read-only variable cached on each machine rather than shipping a copy of it with tasks. They can be used, for example, to give every node a copy of a large input dataset in an efficient manner. Spark also attempts to distribute broadcast variables using efficient broadcast algorithms to reduce communication cost.
scala> val broadcastVar = sc.broadcast(Array(1, 2, 3))
What are accumulators?
Accumulators are variables that are used for aggregating information across the executors. For example, this information can pertain to data or API diagnosis like how many records are corrupted or how many times a particular library API was called.
accumulators à to aggregate information broadcast variables à to efficiently distribute large values.
final Accumulator<Integer> blankLines = sc.accumulator(0);
People familiar with Hadoop Map-Reduce will notice that Spark’s accumulators are similar to Hadoop’s Map-Reduce counters
When using accumulators there are some caveats that we as programmers need to be aware of,
1. Computations inside transformations are evaluated lazily, so unless an action happens on an RDD the transformations are not executed. As a result of this, accumulators used inside functions like map() or filter() won’t get executed unless some action happen on the RDD.
Spark guarantees to update accumulators inside actions only once. So even if a task is restarted and the lineage is recomputed, the accumulators will be updated only once.
Spark does not guarantee this for transformations. So if a task is restarted and the lineage is recomputed, there are chances of undesirable side effects when the accumulators will be updated more than once.
Storage in Spark:
BlockManager manages the storage for blocks (chunks of data) that can be stored in memory and on disk.

BlockManager runs as part of the driver and executor processes. BlockManager provides interface for uploading and fetching blocks both locally and remotely using various stores (i.e. memory, disk, and off-heap).
BlockManager manages the lifecycle of a ShuffleClient:
The ShuffleClient can be an ExternalShuffleClient or the given BlockTransferService based on spark.shuffle.service.enabled configuration property. When enabled, BlockManager uses the ExternalShuffleClient. The ShuffleClient is available to other Spark services (using shuffleClient value) and is used when BlockStoreShuffleReader is requested to read combined key-value records for a reduce task.
When External Shuffle Service is enabled, BlockManager uses ExternalShuffleClient to read shuffle files (of other executors).
BlockManager uses the ShuffleManager for the following:
· Retrieving a block data (for shuffle blocks)
· Retrieving a non-shuffle block data (for shuffle blocks anyway)
· Registering an executor with a local external shuffle service (when initialized on an executor with externalShuffleServiceEnabled)
BlockManager creates a DiskBlockManager when created.
reference: https://books.japila.pl/apache-spark-internals/storage/DiskBlockManager/
Serialization in Spark
A serialization framework helps you convert objects into a stream of bytes and vice versa in new computing environment. This is very helpful when you try to save objects to disk or send them through networks. Those situations happen in Spark when things are shuffled around. RDDs can be stored in serialized form, to decrease memory usage, reduce network bottleneck and performance tuning.
Java serialization
Kryo serialization
Store RDD as serialized Java objects (one byte array per partition). This is generally more space-efficient than deserialized objects, especially when using a fast serializer, but more CPU-intensive to read. By default, Java serialization is used. To enable Kryo, initialize the job with a SparkConf and set spark.serializer to org.apache.spark.serializer.KryoSerializer.
set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.set("spark.kryoserializer.buffer.mb","24")Understanding How Apache Spark Optimization Works?
In order to understand how Apache Spark optimization works, you need to understand its architecture first and in the subsequent section, we will elaborate the same.
The architecture of Apache Spark
The Run-time architecture of Spark consists of three parts –
1. Spark Driver (Master Process)
The Spark Driver converts the programs into tasks and schedules the tasks for Executors. The Task Scheduler is the part of the Driver and helps to distribute tasks to Executors.
2. Spark Cluster Manager
A cluster manager is the core in Spark that allows launching executors, and sometimes drivers can be launched by it also. Spark Scheduler schedules the actions and jobs in Spark Application in FIFO way on cluster manager itself. You should also read about Apache Airflow.
3. Executors (Slave Processes)
Slave processes or Executors are the individual entities on which the individual task of the job runs. They will always run till the lifecycle of a spark Application once they are launched. Failed executors don’t stop the execution of spark job.
4. RDD (Resilient Distributed Datasets)
An RDD is a distributed collection of immutable datasets on distributed nodes of the cluster. An RDD is partitioned into one or many partitions. RDD is the core of Spark as its distribution among various cluster nodes leverages data locality. To achieve parallelism inside the application, Partitions are the units for it. Repartition or coalesce transformations can help to maintain the number of partitions. Data access is optimized utilizing RDD shuffling. As Spark is close to data, it sends data across various nodes through it and creates required partitions as needed.
5. DAG (Directed Acyclic Graph)
Spark tends to generate an operator graph when we enter our code to the Spark console. When an action is triggered to Spark RDD, Spark submits that graph to the DAGScheduler. It then divides those operator graphs into stages of the task inside the DAGScheduler. Every step may contain jobs based on several partitions of the incoming data. The DAGScheduler pipelines those individual operator graphs together. For Instance, Map operator graphs schedule for a single stage, and these stages pass on to the. Task Scheduler in cluster manager for their execution. This is the task of Work or Executors to execute these tasks on the slave.
6. Distributed processing using partitions efficiently
Increasing the number of Executors on clusters also increases parallelism in processing Spark Job. But for this, one must have adequate information about how that data would be distributed among those executors via partitioning. RDD is helpful for this case with negligible traffic for data shuffling across these executors. One can customize the partitioning for pair RDD (RDD with key-value Pairs). Spark assures that set of keys will always appear together in the same node because there is no explicit control in this case.
Mistakes to avoid while Optimizing Apache Spark
1. reduceByKey or groupByKey
Both groupByKey and reduceByKey produce the same answer, but the concept to produce results is different. reduceByKey is best suitable for large datasets because, in Spark, it combines output with a shared key for each partition before shuffling of data. While on the other side, groupByKey shuffles all the key-value pairs. GroupByKey causes unnecessary shuffles and transfer of data over the network.
2. Maintain the required size of the shuffle blocks
By default, the Spark shuffle block cannot exceed 2GB. The better use is to increase partitions and reduce its capacity to ~128MB per partition that will reduce the shuffle block size. We can use repartition or coalesce in regular applications. Large partitions make the process slow due to a limit of 2GB, and few partitions don’t allow to scale the job and achieve parallelism.
3. File Formats and Delimiters
Choosing the right File formats for each data-related specification is a headache. One must choose wisely the data format for Ingestion types, Intermediate type, and Final output type. We can also Classify the data file formats for each type in several ways, such as we can use the AVRO file format for storing Media data as Avro is best optimized for binary data than Parquet. Parquet can be used for storing metadata information as it is highly compressed.
4. Small Data Files
Broadcasting is a technique to load small data files or datasets into Blocks of memory so that they can be joined with more massive data sets with less overhead of shuffling data. For Instance, We can store Small data files into n number of Blocks, and Large data files can be joined to these data Blocks in the future as Large data files can be distributed among these blocks in a parallel fashion.
5. No Monitoring of Job Stages
DAG is a data structure used in Spark that describes various stages of tasks in Graph format. Most of the developers write and execute the code, but monitoring of Job tasks is essential. This monitoring is best achieved by managing DAG and reducing the stages. The job with 20 steps is prolonged as compared to a job with 3-4 Stages.
6. ByKey, repartition or any other operations which trigger shuffles
Most of the time, we need to avoid shuffles as much as we can as data shuffles across many, and sometimes it becomes very complex to obtain Scalability out of those shuffles. GroupByKey can be a valuable asset, but its need must be described first.
7. Reinforcement Learning
Reinforcement Learning is not only the concept to obtain a better Machine learning environment but also to process decisions in a better way. One must apply deep reinforcement learning in Spark if the transition model and reward model are built correctly on data sets and also agents are capable enough to estimate the results.
Apache Parquet, on the other hand, is a free and open-source column-oriented data storage format of the Apache Hadoop ecosystem. It is similar to the other columnar-storage file formats available in Hadoop namely RCFile and ORC. It provides efficient data compression and encoding schemes with enhanced performance to handle complex data in bulk.
Avro:
Avro is an open source data serialization system that helps with data exchange between systems, programming languages, and processing frameworks. Avro helps define a binary format for your data,
It has schema
{
"namespace": "example.avro",
"type": "record",
"name": "User",
"fields": [
{"name": "name", "type": "string"},
{"name": "favorite_number", "type": ["null", "int"]},
{"name": "favorite_color", "type": ["null", "string"]}
]
}
AVRO file format for storing Media data as Avro is best optimized for binary data than Parquet. Parquet can be used for storing metadata information as it is highly compressed.

Executor memory breakdown
Prior to spark 1.6, mechanism of memory management was different, this article describes about memory management in spark version 1.6 and above. Earlier memory management was done using StaticMemoryManager class, in spark 1.6 and above, memory management is done by UnifiedMemoryManager class.
Following diagram shows how the memory is segregated in Spark executor

Spark Executor Memory Distribution
For the sake of understanding, we will take an example of 4GB Memory allocated to an executor and leave the default configuration and see how much memory each segment gets.
Spark Partitioning
By default, Spark creates partitions that are equal to the number of CPU cores in the machine. Data of each partition resides in a single machine. Spark creates a task for each partition. Spark Shuffle operations move the data from one partition to other partitions. Partitioning is an expensive operation as it creates a data shuffle (Data could move between the nodes).
By default, Dataframe shuffle operations create 200 partitions.
Partition is a logical division of the data. This idea is derived from Map Reduce (split). Logical data is specifically derived to process the data. Small chunks of data also it can support scalability and speed up the process. Input data, intermediate data, and output data everything is partitioned RDD.

How does spark partition the data?
Spark uses the map-reduce API to do the partition the data. In input format, we can create a number of partitions. By default, HDFS block size is partition size (for best performance), but it’s possible to change partition size like split.
Factor to consider while partitioning the data
Avoid having too big or small files. Having bigger partitioned data will lead to some of the executor doing the heavy load work, while others are just sitting idle. We need to ensure that no executors in the cluster is sitting idle due to skewed workload distribution across the executors. This will lead to increased data processing time because of weak utilization of the cluster. On the other hand, having too many small files may require lots of shuffling data on disk space, taking a lot of your network compute.
The recommendation is to keep your partition file size ranging 256 MB to 1GB.
How to decide the partition Key(s)?
⦁ Do not partition by columns having high cardinality. For example, don’t use your partition key such as roll no, employee id etc, instead your state code, country code etc.
⦁ Partition data by specific columns that will be mostly used during filter and group by operations.
Partition in memory
You can partition or repartition the Dataframe by calling repartition() or coalesce() method transformations.
· Coalesce – The coalesce method reduces the number of partitions in Dataframe. It avoids full shuffle, instead of creating new partitions, it shuffles the data using the default Hash Partitioner, and adjusts into existing partitions, this means it can only decrease the number of partitions.
· Repartition – The repartition method can be used to either increase or decrease the number of partitions in a Dataframe. Repartition is a full shuffle operation, where whole data is taken out from existing partitions and equally distributed into newly formed partitions
· Partition By – Spark partition By() is a function of pyspark.sql.DataFrameWriter class which is used to partition based on one or multiple column values while writing Data frame to Disk/File system. When you write Spark Data frame to disk by calling partition By(), Pyspark splits the records based on the partition column and stores each partition data into a sub-directory.
Print has taken till here
Block Manager:
It is crucial component that handles the storage and retrieval of data blocks.
Crucial for efficient data processing.
Duties:
1) Block storage and management: keep track of all blocks in local node and their status.
2) Data serialization/deserialization: ser/dev in reading/writing.
3) Data replication: block replica for fault tolerant
4) Block transfer: required when shuffling operation. Data redistribution in cluster.
5) Memory management: works with spark memory management to evict block when needed for other task, implementing strategies like LRU.
6) Integration with storage system: intract with external storage, S3, HDFS
7) Metadata management: maintains metadata of blocks(size, location, status in memory or disk)
8) Handling faults: help to recover the lost blocks using replicas.
Apache Spark: workflow internal

When Spark reads a file from HDFS, it creates a single partition for a single input split. Input split is set by the Hadoop InputFormat used to read this file. For instance, if you use textFile() it would be TextInputFormat in Hadoop, which would return you a single partition for a single block of HDFS (but the split between partitions would be done on line split, not the exact block split), unless you have a compressed text file. In case of compressed file you would get a single partition for a single file (as compressed text files are not splittable).
When you call rdd.repartition(x) it would perform a shuffle of the data from N partitions you have in rdd to x partitions you want to have, partitioning would be done on round robin basis.
If you have a 30GB uncompressed text file stored on HDFS, then with the default HDFS block size setting (128MB) it would be stored in 235 blocks, which means that the RDD you read from this file would have 235 partitions. When you call repartition (1000) your RDD would be marked as to be repartitioned, but in fact it would be shuffled to 1000 partitions only when you will execute an action on top of this RDD (lazy execution concept)
Promote this post
Share & amplify
LinkedIn post generator
Generate a polished, professional LinkedIn post from this article. Edit before posting.

Comments (0)