Apache Spark Was Hard Until I Learned These 30 Concepts! — Transcript
Full transcript
- 0:00Apache Spark is a critical topic that
- 0:03helped me clear interviews and land
- 0:05high-paying opportunities at multiple
- 0:08big tech companies like Apple, Uber,
- 0:11Atlassian, and Databricks. In this
- 0:13video, I am going to simplify the 30
- 0:17most important Spark concepts
- 0:19that transform the way I write,
- 0:23tune, and troubleshoot my Spark jobs, so
- 0:25that you can confidently do the same.
- 0:28So, let's get started. Before the world
- 0:30started using Apache Spark, MapReduce
- 0:33already existed, right? But, there were
- 0:36multiple performance problems with
- 0:38MapReduce. The most pressing problem was
- 0:41that MapReduce writes
- 0:43every stage's intermediate data to HDFS,
- 0:46and that is what makes it super slow.
- 0:49So, let's understand this with an
- 0:51example where we want to count the words
- 0:54in a file. So, let's say we have a file,
- 0:56and the file simply has this content.
- 1:00And this is going to go through the map
- 1:02step. The map step simply transforms
- 1:05each line into key-value pairs. So,
- 1:07basically, what it does is that it
- 1:09converts each of the line into key-value
- 1:12pairs, which looks like word, {comma}
- 1:15one. So, the next step is converting
- 1:17each of the lines into words. And this
- 1:20is what happens. Each of it gets split
- 1:22into word, hello world, hello Spark, and
- 1:24so on. And then,
- 1:26it creates key-value pairs. So, hello is
- 1:29going to be looking like this, hello
- 1:31{comma} one, similarly for world and all
- 1:34of this, right? So, the map step creates
- 1:37key-value pairs, as you see over here,
- 1:39and the output of the map step is going
- 1:42to be written back to the disk, right?
- 1:45So, the result is going to be written
- 1:48back to the disk, and this is what I was
- 1:49talking about earlier.
- 1:51Now, the next step is the shuffle and
- 1:54sort step.
- 1:56And the shuffle and sort step is going
- 1:59to read the data from the disk that was
- 2:02written in the previous step, right?
- 2:04In the shuffle and the sort step, the
- 2:06same keys are going to land at the same
- 2:09place or in the same partition, right?
- 2:11So, we see hello landing at the same
- 2:13partition or the same place along with
- 2:16all of it counts, right? And similarly
- 2:19for all of the other words. And the
- 2:21result of this is again written back to
- 2:23the disk. The reduce step reads the
- 2:27result from the previous step from the
- 2:30disk again, which is over here. And the
- 2:32job of the reduce step is to combine the
- 2:34output. So, it is simply going to sum
- 2:37these outputs or the the the individual
- 2:39counts, right? And then we see that
- 2:42hello has a final count of three, while
- 2:44all the others have a final count of
- 2:47one.
- 2:48And the result is finally
- 2:51again written back to the disk. So, this
- 2:53repeated read and write from the disk is
- 2:57what makes MapReduce super slow. The
- 2:59second problem, as we saw in this
- 3:02example, is that MapReduce has a rigid
- 3:06programming model in the sense that it
- 3:08only has two steps, which is map and
- 3:11reduce. So, if you have complex
- 3:14pipelines involving steps like filter,
- 3:16joins, and aggregation, you would need
- 3:19to break them into multiple MapReduce
- 3:22steps. So, how did Spark solve this
- 3:25problem? Spark introduced number one,
- 3:28in-memory processing, which means that
- 3:30data can be cached or be present in
- 3:34memory across the cluster nodes. So, in
- 3:38this example, what we saw was that after
- 3:40every intermediate step, data was
- 3:43flushed back to the disk, right? Data
- 3:46was flushed back to the disk over here,
- 3:48here, and here,
- 3:50right?
- 3:51Which is very costly, which is very
- 3:53costly. So, instead of this, across all
- 3:57of these steps,
- 3:58data can be present in the memory
- 4:01itself, in the RAM itself, without the
- 4:05need to flush back intermediate results
- 4:08to the disk. And this was the reason why
- 4:10Spark became 10 to 100 times faster.
- 4:13Number two, and instead of the rigid
- 4:16MapReduce workflow,
- 4:18Spark built something called the
- 4:19directed acyclic graph. Right? And it is
- 4:23going to be composed of transformations,
- 4:25actions, and it is going to properly
- 4:28divide it into something called jobs,
- 4:30stages, and tasks. So, it is going to
- 4:33look something like this.
- 4:35This is what exactly it is going to look
- 4:37like. So, don't worry, we will cover
- 4:39this in a lot of detail, but the idea
- 4:42was to tell you that it is going to be
- 4:44converted into something called jobs,
- 4:48and then stages, and then all
- 4:52all of these tasks. Right?
- 4:56So, this programming model simplifies a
- 4:59lot of things from the MapReduce
- 5:02programming model that we had earlier. I
- 5:04want to take a moment to talk about
- 5:06Educative. Educative is a fully
- 5:08interactive and text-based platform with
- 5:11excellent visuals and diagrammatic
- 5:14explanations, where you learn by doing.
- 5:16Their legendary Rocking the System
- 5:18Design Interview is one of the most
- 5:21widely accessed resources used by
- 5:23thousands of learners preparing for
- 5:25interviews at Google, Meta, Amazon, and
- 5:28other big tech companies. So, the course
- 5:30walks you through concepts for building
- 5:32your foundation, helps you understand
- 5:34how to do back-of-the-envelope
- 5:35calculations, and then prepares you
- 5:37through several examples. Right? They
- 5:40also have excellent courses on interview
- 5:43preparation, G&A, cloud platforms, labs,
- 5:46and mock interviews. Some of my personal
- 5:48favorites are Mastering MCP for advanced
- 5:52agentic applications and this one on
- 5:54cursor AI, something that [music] most
- 5:56of us use in our daily life. Right now,
- 5:58they are offering a massive 50% one-time
- 6:01discount and on top of that, you can use
- 6:03the link in the description below to get
- 6:06an additional 10% off. Now, let's get
- 6:09back to the video. Now, before we take a
- 6:10deep look into the Spark job itself,
- 6:14let's understand how and where does the
- 6:17Spark job get executed, right? So, it
- 6:20gets executed on a cluster and a cluster
- 6:24is simply a group of machines connected
- 6:28together through a network, right? So,
- 6:32think of it as multiple computers
- 6:34connected together through a network and
- 6:38we want to use the combined computing
- 6:41power of all the computers within the
- 6:44network, right? So, that is simply a
- 6:47cluster. So, let's assume that this is
- 6:50our cluster and our cluster has seven
- 6:53machines in total. So, you see six
- 6:55workers and one master node.
- 6:58And the configuration of our cluster is
- 7:0232 cores,
- 7:0432 cores for one machine and 120 GB
- 7:09of RAM for one machine, right? So, in
- 7:12total, the capacity of our cluster is 7
- 7:15into 32, which is 224
- 7:19cores
- 7:21and 7 into 120, which is 840
- 7:24GB
- 7:25of RAM.
- 7:27Right? So, this is the configuration of
- 7:30our cluster, right? So, there is going
- 7:33to be one master node, as you see over
- 7:36here, and the rest are worker nodes,
- 7:39right? So, the master node is the place
- 7:42where your cluster manager the cluster
- 7:45manager as we see over here the cluster
- 7:48manager is going to reside right
- 7:52so the cluster manager or the resource
- 7:54manager is the one which is responsible
- 7:58for managing the entire pool of the CPU
- 8:02and the RAM that we've got right yarn
- 8:05and kubernetes are the most widely used
- 8:08cluster manager so in our example we are
- 8:11simply going to stick to yarn
- 8:13right so in order to start the spark
- 8:16application the user is going to submit
- 8:20or the user is going to run a spark
- 8:23submit command along with the code the
- 8:26job configuration and all of that right
- 8:29so this request this request is going to
- 8:32go to yarn
- 8:34so this request from the user goes to
- 8:38yarn the resource manager is a very
- 8:41smart guy is a born leader and it knows
- 8:45how to delegate it task right so yarn
- 8:48the first thing that it is going to do
- 8:49is it is going to create the first
- 8:52container which is called application
- 8:55master so it is going to select any of
- 8:57the nodes the worker nodes at random and
- 9:00create a
- 9:02application master container right so
- 9:04let's say it creates it on worker five
- 9:08right
- 9:10so now just for your context a container
- 9:13is basically an isolated environment
- 9:15within a node with fixed amount of CPUs
- 9:19and RAM it is not strictly not allowed
- 9:23to use any resources on the node beyond
- 9:26what has been allocated to it right so
- 9:28now if we just zoom in into the
- 9:31application master container this is
- 9:34what it is going to have right this is
- 9:38what it is going to have.
- 9:41The application master container is
- 9:43going to start an application driver.
- 9:47That is what you see over here. And an
- 9:49application driver is nothing but a JVM
- 9:53process, right?
- 9:54Its job is to simply execute the main
- 9:58method of our Spark application code.
- 10:00So, you remember we submitted the Spark
- 10:02application code along with the job
- 10:05parameters and all of that, the
- 10:06configuration, right? So, the the
- 10:08purpose of the application driver is to
- 10:12execute the main method of our Spark
- 10:16application,
- 10:17right?
- 10:18Now, all of this is good if we've
- 10:21written our job in
- 10:23Java or Scala because they are JVM
- 10:26languages, right? They're JVM languages,
- 10:28but what if we've written our code in
- 10:30PySpark? JVM does
- 10:33doesn't just understand Python, right?
- 10:37So, in that case, a PySpark driver
- 10:41a PySpark driver is going to spin up and
- 10:44it is going to have its own method, own
- 10:47main method, the only job of which is to
- 10:50call the main method of the application
- 10:54driver, right?
- 10:56So, there is a PySpark driver, there is
- 10:57the application driver, the PySpark
- 10:59driver's main method calls the
- 11:02application driver's main method, right?
- 11:05And the application driver's main
- 11:06method, because it is written in Java or
- 11:09Scala, it simply comfortably runs on the
- 11:12JVM. A very important point to notice
- 11:15that Spark core is written in Scala,
- 11:18right? And to make it available to
- 11:21Python developers, it has a Java wrapper
- 11:24on top of it, which further has a Python
- 11:27wrapper on top of it called PySpark. So,
- 11:31the way it works is that the Python
- 11:33wrapper makes a call to the Java wrapper
- 11:36and then the Java wrapper makes a call
- 11:39to Spark core.
- 11:40Right? And this finally comfortably
- 11:43executes on the JVM. Right?
- 11:46So, going back to our
- 11:48application driver, now the application
- 11:51driver is going to execute the main
- 11:53method of our application.
- 11:55Now, it realizes that, "Okay, I need
- 11:57resources. Right? I need resources in
- 12:00order to be able to execute my Spark
- 12:03job." So, it looks at what were the
- 12:06resources that were requested and it
- 12:09says it understand that 30 GB of memory,
- 12:13four cores, and three executors
- 12:17were requested. Right? Three executors
- 12:20were requested. So, what it simply does
- 12:22is that it goes to YARN,
- 12:26which is your
- 12:28resource manager, and it simply tells it
- 12:31that, "Hey, I
- 12:33I need resources.
- 12:36Right?
- 12:37I need resources and the resources that
- 12:39I need is
- 12:41three executors
- 12:43and each of these three executors should
- 12:45have 30 GB of RAM
- 12:48and four cores. Right? And it should be
- 12:51three of such executors. Now, what YARN
- 12:55does is that it tells the application
- 12:58driver that, "Okay, I'm going to assign
- 13:00you resources for completing your job."
- 13:04And it is simply going to pick up
- 13:07randomly any three workers and create
- 13:10executor containers. Right? So, let's
- 13:13say
- 13:14it
- 13:15picks up worker number one,
- 13:18worker number two,
- 13:20and worker number four.
- 13:24Right? It simply creates
- 13:27executor containers over here. This is
- 13:29number one. This is number two. And this
- 13:33is number three, right? It simply
- 13:35creates them and then hands over the
- 13:38details to the application driver over
- 13:41here, right? So, let me just quickly
- 13:44change the color.
- 13:45So, it simply tells the application
- 13:47driver that, "Hey, these are your three
- 13:49executors, and here are the details.
- 13:52They are present on worker one,
- 13:54worker two, and worker four."
- 13:58Now, what the application driver does is
- 13:59that it then schedules the task across
- 14:03these workers. So, it then schedules the
- 14:05task across this one, this one, and this
- 14:08one. The workers
- 14:10um the executors then execute the task
- 14:12and then finally report the result back
- 14:16to the driver, right? Now, again, a very
- 14:18important point to note is that if we
- 14:20UDF if we've written UDFs or libraries
- 14:24that are not natively available in
- 14:26Spark, each executor container may also
- 14:29spin up a Python process. So, right now
- 14:32we see that it only has a JVM process,
- 14:35right? It only has a JVM process. What
- 14:38it might do is that it might also spin
- 14:40up a Python process,
- 14:44similar to what we see over here, right?
- 14:48So, along with the JVM process, it is
- 14:51going to spin up a Python process to
- 14:54help you execute your UDF code or code
- 14:58written in libraries that are not
- 15:00natively available in PySpark, right?
- 15:03Now, in this architecture, we see that
- 15:05the user or the client submits a request
- 15:09to the YARN resource manager, right? So,
- 15:12the user submits a request to the YARN
- 15:16resource manager, right?
- 15:18The resource manager then starts your
- 15:21application master container and the
- 15:24application master container then starts
- 15:27your application driver process, right?
- 15:31The application master then starts your
- 15:33application driver process. Now, this
- 15:35execution model where the application
- 15:39driver runs on one of the nodes in the
- 15:42cluster. And in our case
- 15:45the node of the cluster is worker number
- 15:49five, right?
- 15:51So, this execution model where the
- 15:53application driver runs on one of the
- 15:55nodes in the cluster is called the
- 15:59cluster mode,
- 16:03right?
- 16:04Now, there is another deployment method
- 16:07where the Spark submit is not routed to
- 16:10the YARN resource manager. Instead, when
- 16:12we run Spark submit, the command starts
- 16:15the application driver on the user or
- 16:18the client machine itself, right? Now,
- 16:20instead
- 16:21in the client mode, the application
- 16:24driver is started on the user or the
- 16:28client machine itself,
- 16:32right?
- 16:32And it is the application driver which
- 16:35is then going to coordinate with the
- 16:37YARN resource manager for the allocation
- 16:40of resources, creation of executor
- 16:42containers, and all of that, right? So,
- 16:44this execution model where the
- 16:47application driver runs on the client or
- 16:51the user machine instead of running on
- 16:55the worker nodes on the cluster is
- 16:57called the client mode.
- 17:02So, let's assume that our deployment has
- 17:04now happened using the cluster mode and
- 17:07our Spark application contains a mix of
- 17:11transformations and actions, right? So,
- 17:15what exactly are these transformations
- 17:18and action? And we'll refer back to the
- 17:21same code the word count example that we
- 17:24were referring to earlier.
- 17:26Transformations are simply operations
- 17:29that produce a new data set from an
- 17:31existing one, right? So, they are lazy
- 17:34operations simply meaning that no matter
- 17:37how many operations you perform, they
- 17:39are not executed immediately. So, all of
- 17:43these steps, right? The reading of the
- 17:45file
- 17:47the reading of the file
- 17:49splitting of each of the line into
- 17:51words, exploding each of the line into
- 17:53words, and then actually doing a group
- 17:56by and the count
- 17:58all of this, right? They are going to be
- 18:01appended to the list of operations
- 18:03needed to be executed. It's not actually
- 18:06going to be executed. That is what it
- 18:09means when I say it's lazy in nature,
- 18:12right? So, getting back to
- 18:13transformation example Examples of
- 18:15transformation include, let's say, a
- 18:17select operation or a
- 18:20filter operation. Let's say I had a
- 18:23filter operation over here saying that,
- 18:24"Okay, filter out all of the words which
- 18:27start with the letter W." Right? Or
- 18:31let's say you want to add a column using
- 18:33with column in Spark. All of these are
- 18:38transformations, right? So, they keep on
- 18:40get add getting added to the list of
- 18:42operations until an action is called.
- 18:47Right? Now, what is an action? An action
- 18:50is an operation that triggers actual
- 18:53computation, right? So, some examples of
- 18:55it is collect or a show or a count or a
- 18:59save or let's say you write to a parquet
- 19:01file. All of this, they are actions. So,
- 19:04let's assume that I would call a dot
- 19:07show after
- 19:10computing the count, this is going to be
- 19:13an action.
- 19:15Yeah?
- 19:16So, it's important to understand that
- 19:19again, transformations can be classified
- 19:21into narrow and wide. In narrow
- 19:24transformations, each output partition
- 19:26depends on a single input partition and
- 19:30there is no shuffle involved. The
- 19:32keyword here is there is no shuffle
- 19:34involved. So, we are going to understand
- 19:36what exactly shuffle is, but
- 19:39in wide transformation, a shuffle is
- 19:42involved. So, transformation they are of
- 19:43two types, narrow and wide. Narrow
- 19:46doesn't require shuffle. One input
- 19:48partition gives one output partition,
- 19:50but in a wide transformation, it
- 19:53requires a shuffle. And shuffle is
- 19:56another interesting activity where your
- 19:58data is redistributed across the
- 20:01clusters. And to put it simply, the same
- 20:04keys land on the same partition in a
- 20:07shuffle. So, let's understand that with
- 20:09an example over here.
- 20:11The same word count example. And don't
- 20:15worry about the stages,
- 20:17the exchange and all of this over here
- 20:20for now. We are going to go into a lot
- 20:22of details on the job stages, tasks,
- 20:24shuffle and all of this, right? So, we
- 20:26read in this file, it got split up into
- 20:28words, and then we got the key value
- 20:32pairs, right? Similar to what we were
- 20:34doing in MapReduce. So, let's exactly
- 20:36replicate the same step. Now, I want to
- 20:39do a group by, right? So, this means
- 20:42that I need to count the hellos, I need
- 20:45to count all of the other letters,
- 20:47right?
- 20:48So, in a shuffle, the same keys land on
- 20:52the same partition. So, you see all the
- 20:54hellos are going to land in the same
- 20:56partition, right? And then you are going
- 20:59to have the word
- 21:01Spark and MapReduce, right?
- 21:03Then finally, the count over here. So,
- 21:06shuffle is going to redistribute your
- 21:09data to make your keys. In this case,
- 21:12the keys is this one over here because
- 21:14we are doing a group by word, word
- 21:16becomes a key. It is going to make sure
- 21:18that the same keys land in the same
- 21:21partition, right? Now, this is becoming
- 21:24a little over simplified. What would
- 21:26generally happen is there would be a
- 21:28rule which would say, let's say key
- 21:30modulo some number, right? And let's say
- 21:33there are three partitions.
- 21:361 2 and
- 21:383, right? So, I would say key modulo 3
- 21:42and it is always going to give me a
- 21:43number between 0 1 and 2 and whatever
- 21:46number I get
- 21:48it is going to decide which
- 21:50partition the key is going to go into.
- 21:52So, let's say for hello,
- 21:54I'm going to say simply hello mod 2.
- 21:57Now, because this is not a number, let's
- 21:59say I'm going to do a hash of hello, it
- 22:02is going to give me an integer and this
- 22:05integer mod 3 is going to give me a
- 22:07number between 0 1 and 2
- 22:10and this is going to decide which
- 22:12partition it is
- 22:13going to land into. So, this makes it
- 22:15obvious that all the hellos are going to
- 22:18land in the same partition because
- 22:21this function is going to give the same
- 22:24values for all the hellos, right? So, I
- 22:27hope this makes sense. Now, every time a
- 22:30user calls an action, the driver is a
- 22:34very smart and organized guy. So, before
- 22:36it sends the task to the executors, it
- 22:39creates a plan.
- 22:42It is going to take up the code that
- 22:44we've written in the form of SQL or data
- 22:47frames and it checks whether the syntax
- 22:50of the code is correct or not, right?
- 22:53Once validated, it creates what is known
- 22:56as the unresolved logical plan.
- 22:59Yeah?
- 23:00Internally, it maintains something
- 23:03called a catalog. So, it has a structure
- 23:06called a catalog.
- 23:08And it has details about the table,
- 23:10databases, data types, and all of it,
- 23:13right? So, let's say you're writing a
- 23:14query which looks like something like
- 23:16select a {comma} b from
- 23:21table a.
- 23:24Right? So, let's say you wrote this
- 23:25query. What it is going to do is it
- 23:28verifies whether table a is present
- 23:32within the Spark ecosystem or not,
- 23:34right? And whether column a and b are a
- 23:38part of table a, right? So, this
- 23:40basically makes sure that your query is
- 23:43semantically correct. Now, once this is
- 23:46verified, it creates the logical plan.
- 23:50So, the unresolved logical plan goes
- 23:52through the catalog, and then finally it
- 23:54creates the logical plan. The logical
- 23:58plan then goes through the Catalyst
- 24:00Optimizer,
- 24:02and finally produces the optimized
- 24:04logical plan. Now, you can ask me that
- 24:07what exactly are these optimizations,
- 24:11right?
- 24:12What exactly are these optimizations?
- 24:15So, some examples could be filter push
- 24:17down,
- 24:18projection push down, right? Some
- 24:20examples can be filter push down,
- 24:25and projection push down.
- 24:28Filter push down simply means that the
- 24:30filtering logic is directly executed at
- 24:34the data source, right? So, let's say
- 24:36you wrote a filter
- 24:39city equals Boston.
- 24:43So, this filtering logic is directly
- 24:46executed at the data source without
- 24:48having to bring back the data to the
- 24:51executor, and then removing the unwanted
- 24:54record, right? Because then you've lost
- 24:56a lot of energy and bandwidth in
- 24:59bringing all of the chunk of records to
- 25:02the executors and then the executors
- 25:04doing the heavy lifting in removing the
- 25:06unwanted rows, right? So, that is where
- 25:08filter push down can be very, very
- 25:11helpful, yeah? And projection push down
- 25:14reduces the number of columns that you
- 25:17need to read that needs to be read from
- 25:19the data source by fetching only the
- 25:22columns that are written within the
- 25:23query, right? So, let's say if you have
- 25:25a query which looks like select name,
- 25:28{comma} age and then inside of this
- 25:33inside of this you have a query which
- 25:35says select star
- 25:37from whatever
- 25:40instead of bringing all of the data
- 25:43because we mentioned a select star, it
- 25:45is specifically going to bring the name
- 25:47and the age because it knows that the
- 25:50output only requires name and age,
- 25:53right? So, this is where also projection
- 25:56push down reduces the amount of data
- 25:59that we are pulling and executing on our
- 26:02executors. So, now after the
- 26:04optimization that I said and they have
- 26:06been applied, it creates an optimized
- 26:09logical plan and this optimized logical
- 26:12plan is converted into several physical
- 26:16plans, right? But before choosing which
- 26:19physical plan is going to run on the
- 26:20cluster all of these generated physical
- 26:24plans all of these generated physical
- 26:26plans go through something which is the
- 26:29cost model, right? It goes through the
- 26:32cost model
- 26:34which is going to decide which is the
- 26:36most optimal plan that can be run on the
- 26:39cluster, right? So, once it goes through
- 26:41the cost model, the final physical plan
- 26:44comes out and now the DAG scheduler
- 26:46kicks in.
- 26:48Right? So, the DAG scheduler kicks in
- 26:52and it creates a job. It takes up this
- 26:55final physical plan and it converts it
- 26:58into something called a DAG of stages
- 27:02and task and finally executes it on the
- 27:05cluster. Now, what exactly are the jobs,
- 27:08stages, and tasks? So, let's understand
- 27:11this with this example. Let's assume
- 27:13that we submitted this piece of code.
- 27:16So,
- 27:17what exactly is happening here is we are
- 27:19reading a parquet file and then we are
- 27:21applying a filter uh where the amount is
- 27:24greater than 1,500. And then we do a
- 27:28repartition of three, meaning that if we
- 27:30read in one partition, it is now going
- 27:33to be split into three partition, right?
- 27:35Then we select the region and amount and
- 27:39then we add a new column which is tax
- 27:42and finally we do a group by by the
- 27:45region and sum the sale amount and the
- 27:49tax amount, right?
- 27:50The sale amount over here and the tax
- 27:52amount. And lastly, we finally add what
- 27:55exactly was the timestamp where this
- 27:57computation happened. And we trigger an
- 27:59action which is collect. So, let's first
- 28:02break this down into different steps,
- 28:04right?
- 28:05So, the first step over here is the
- 28:08reading of the file which is spark.read.
- 28:11Number one. And then the second step
- 28:13over here is a filter. The third step
- 28:17over here is a repartition,
- 28:19right? The fourth step over here is a
- 28:22select. The fifth step over here is a
- 28:26withColumn where we add in the tax
- 28:30column.
- 28:31The next step over here is group by
- 28:33which is sixth step.
- 28:35Then we do an aggregation
- 28:37which is the seventh step and then the
- 28:40last step over here before
- 28:42the action is a withColumn and then
- 28:45finally we do a collect, which is an
- 28:49action, right? Now, the the
- 28:53the ones which are marked in gray, these
- 28:56are narrow transformations, right?
- 28:59These are narrow transformations.
- 29:04Because they don't require any shuffle,
- 29:06right? Repartition number three and
- 29:08group by, these are wide transformations
- 29:12because they are going to require a
- 29:14shuffle, right? So, a very important
- 29:16point to note is
- 29:19of course, we understood about narrow
- 29:21transformations.
- 29:23Narrow {slash} wide transformations.
- 29:30Number one, number two, what exactly are
- 29:33the job stages and tasks, right? So,
- 29:35whenever we see an action being called,
- 29:38for example, {dot} collect, right?
- 29:40Uh this is an action being invoked by
- 29:43the user. Spark is going to create job
- 29:46in order to get the result of this
- 29:49action, right? So, this is going to
- 29:51create a job where all of these eight
- 29:54steps are going to be executed to
- 29:56produce the output and give you back the
- 29:59result, right?
- 30:00Now, that is a job. Now, stages
- 30:03are created at the boundary of a
- 30:06shuffle. So, this is where a stage is
- 30:08going to be created, right? Stage one is
- 30:11going to be created. So, this is all a
- 30:13part of stage one.
- 30:15And then a stage is going to be created
- 30:17over here again.
- 30:19Stage two. This is going to be a part of
- 30:23stage two. And this is going to be stage
- 30:26three.
- 30:28So, if you have n shuffles
- 30:31or n wide transformations,
- 30:34n shuffles or wide transformations,
- 30:40you are going to have n plus one
- 30:45stages.
- 30:46So, in this case, we have two wide
- 30:48transformations.
- 30:50We are going to have three stages. So,
- 30:53stage one over here
- 30:56is going to look like this.
- 31:00This one over here. Stage two, what we
- 31:03discussed over here is going to look
- 31:05like this.
- 31:06And the final stage three
- 31:10is
- 31:11the one over here, right? So, let's go
- 31:14through this example
- 31:16through all of the items in the stages,
- 31:19right? The first step was a read, and
- 31:21then a filter, and then a repartition,
- 31:23right? Let's say we read in this file,
- 31:26which is the first step, and then we do
- 31:28a filter, where filter amount was
- 31:30greater than 1,500, and that is what it
- 31:33leaves you with over here. And then we
- 31:35did a repartition.
- 31:38So, you see that we read in one
- 31:40partition over here,
- 31:42but now we repartition it into three
- 31:44partitions. So, this is what number one,
- 31:47partition number two, and number three.
- 31:49Three partitions, right? And this get
- 31:52written to the shuffle write exchange.
- 31:55So, this is basically a buffer storage,
- 31:57which is where your data gets written
- 32:00when shuffling is happening, right? So,
- 32:02that's the first step, first stage. And
- 32:06then we go to the second stage, where
- 32:08all of this is read in, and we are going
- 32:11to do a select. The select was
- 32:15to select the region and the amount.
- 32:18And then we added a tax column, right?
- 32:20You remember that the tax the the
- 32:22formula for the tax
- 32:24tax column was 0.8 into the amount, and
- 32:27that is what we are going to apply over
- 32:28here.
- 32:31So, we simply select the relevant
- 32:33column, which is these two columns, and
- 32:35then we add a new column called tags,
- 32:38and this is going to apply for all of
- 32:40the three partitions.
- 32:42Right? These two steps, number one and
- 32:44number two. Right? And then a group by
- 32:47is going to happen where we are going to
- 32:49write to the shuffle write exchange.
- 32:52Now, a very important thing to note is
- 32:55these steps, right?
- 32:57These steps are executed on each of the
- 33:00partition. Right? This is going to be
- 33:02executed on this partition, on this
- 33:04partition, on and on this partition.
- 33:06Right? So, when it's executed on one
- 33:09partition, it is called a task. Right?
- 33:12And one task is operated on one
- 33:15partition
- 33:17by one core.
- 33:20Right? To keep it To keep things simple,
- 33:22by one core. Right? So, we have three
- 33:24partitions over here, number one, two,
- 33:27and three. So, there are going to be
- 33:29three tasks.
- 33:31So, remember that if you have n
- 33:34partitions,
- 33:37there are going to be n parallel
- 33:40tasks.
- 33:42Right?
- 33:42So, these can be executed parallelly
- 33:44because there's no dependency of one on
- 33:46the other. The select and with column on
- 33:49one partition can completely go on in
- 33:52parallel with this partition and this
- 33:54partition. Right?
- 33:56So, then we do a group by, and this data
- 33:59is written back to the shuffle write
- 34:01exchange, and this is written in as
- 34:04three partitions.
- 34:06Yeah?
- 34:06So, when we do a group by, the first
- 34:08thing that happens is a shuffling, and
- 34:10in shuffling, the same keys come to the
- 34:12same partition, and you see that Kolkata
- 34:16has come to the same partition. Earlier,
- 34:17it was over here and over here. Now, it
- 34:19has come to the same partition. Right?
- 34:21It gets written to the shuffle write
- 34:23exchange, and then it is read up
- 34:27in the shuffle read exchange, which is
- 34:29stage number three over here. Right?
- 34:32Now, once this is read up, we do a sum
- 34:35and it is simply going to sum up the
- 34:37total amount
- 34:38and the tax over here. And this is going
- 34:41to happen for all of the three
- 34:43partition, right? And the last
- 34:46step is a width column. The width column
- 34:48is simply going to add a timestamp and
- 34:51this gives us the final
- 34:54result.
- 34:56Right? So, what we've seen here is, to
- 34:59summarize, we called an action as a
- 35:01result of it and that action was this
- 35:04collect action over here. As a result of
- 35:06it, a job was created and when a job was
- 35:09created, we figured out the narrow and
- 35:12the wide dependencies,
- 35:14the wide transformation, and we figured
- 35:16out that there are going to be three
- 35:20stages.
- 35:21And this is how the stages are created.
- 35:23Number one, number two, and number
- 35:26three.
- 35:27Right? And we created three partitions,
- 35:30because of which we had three parallel
- 35:33task in place. So, now I believe you
- 35:35understand how job stages and task are
- 35:39created. So, finally, you must be
- 35:41thinking that Spark does all of this
- 35:44magic in memory and what exactly goes
- 35:47behind the scenes in memory to enable
- 35:51Spark to be able to do all of this
- 35:53magic. So, we've seen YARN spin up
- 35:56executor containers on the worker nodes,
- 35:59right? Now, let's have a deeper look
- 36:01into the memory section of the executor
- 36:05container, yeah? So, the container has
- 36:09three important regions, right? The
- 36:11first one is the on-heap memory, right?
- 36:15This is the most important section. This
- 36:17is the main area which is managed by the
- 36:21JVM and this is where most of the Spark
- 36:24operations run, right? And this is
- 36:26divided into four sections, which as you
- 36:29see over here is execution memory,
- 36:31storage memory, user memory, and
- 36:34reserved memory. So, execution memory is
- 36:37the place where your joins, your sorts,
- 36:40your shuffle, and your aggregation, some
- 36:43of the things that we do very frequently
- 36:45in our code, all of this happen in the
- 36:48execution memory, right? The storage
- 36:51memory is the place where caching, let's
- 36:54say we've cached our data frame or we've
- 36:56created broadcast variable, we are doing
- 36:58broadcast joins, right? So, this is the
- 37:01place where all of those entities live,
- 37:04where our cached variable,
- 37:06where our cached data frame, where our
- 37:08broadcast variables live, right? The
- 37:10third one is user memory, yeah? User
- 37:14memory is the one which is used by our
- 37:17own program's objects, right? So, let's
- 37:19say we create a lot of list, variables,
- 37:22user-defined functions, and all of that.
- 37:25All of that is going to live in the user
- 37:28memory. And lastly, is the reserved
- 37:31memory. Reserved memory is a very small
- 37:34slice. So, Spark keeps it for its own
- 37:38internal housekeeping and working. The
- 37:40next important section is the off-heap
- 37:43memory,
- 37:45which we see over here. This is
- 37:47generally disabled by default. It is
- 37:50used when you enable
- 37:52this setting using
- 37:53spark.memory.offHeapEnabled,
- 37:56and actually give it a size by using
- 37:58spark.memory.offHeap.size, right? So,
- 38:00this this section of the memory
- 38:05is useful when Spark is handling large
- 38:08data sets or doing heavy shuffle, and
- 38:11frequent garbage collection is
- 38:13happening, right? Frequent garbage
- 38:14collection generally happens when
- 38:17there's a lot of aggravation of objects
- 38:20within the memory, and Spark needs to
- 38:23free up that memory in order to be able
- 38:25to use it for subsequent steps. That is
- 38:28when garbage collection happens, so that
- 38:31it removes those unwanted objects,
- 38:34right? And garbage collection is a
- 38:36costly operation because it stops your
- 38:38program and then does the cleanup and
- 38:40then resumes your program, right? So, in
- 38:44such a case where frequent GC or garbage
- 38:47collection is happening, the data can be
- 38:49kept outside the JVM on the off-heap
- 38:53memory because the off-heap memory is
- 38:55generally immune to garbage collection.
- 38:58But, then it's very important to note
- 39:00then that you would have to do them the
- 39:03the cleaning of the unwanted objects in
- 39:05order to avoid memory leaks, right?
- 39:08If you want to completely avoid garbage
- 39:10collection, right? So, in this case
- 39:12specifically, you can use the off-heap
- 39:15memory.
- 39:16The third and the last part is the
- 39:19overhead memory.
- 39:20The overhead memory is the extra space
- 39:23that Spark requests from YARN
- 39:27in order to handle its internal
- 39:29system-level operations, right? Now, a
- 39:31very interesting thing to note is this
- 39:33execution plus uni- plus the storage
- 39:37memory is together called the unified
- 39:41memory
- 39:43because the slider that you see, right?
- 39:46The slider that you see is actually
- 39:48movable. So, they both can borrow memory
- 39:51from each other, but more preference is
- 39:54given to the execution memory because,
- 39:56of course, most of the important
- 39:58operation within our Spark lifecycle
- 40:01happen in the execution memory, right?
- 40:04So, let me take a quick example and
- 40:06explain you how each of the sections of
- 40:08this memory is calculated. So, let's say
- 40:10when you submit your Spark job, right?
- 40:13So, you say exec {hyphen} {hyphen}
- 40:15executor
- 40:16{hyphen} memory
- 40:18is let's say 10 GB.
- 40:21Right? Let's say 10 GB. So, this whole
- 40:24section for the on heap memory from here
- 40:26to here is defined by
- 40:29spark.executor.memory,
- 40:31right? So, this is going to be 10 GB.
- 40:35What we submitted for
- 40:37the executor memory over here. This is
- 40:39the on heap memory. Now, this portion
- 40:41this portion of the unified memory from
- 40:43here to here execution and storage
- 40:46memory, it is defined by
- 40:48spark.memory.fraction,
- 40:50and this is 0.6 of your
- 40:54spark.executor.memory,
- 40:55right? So, this is going to be 0.6 into
- 40:5710 GB, which is 6 GB, and this is going
- 41:01to be equally divided between execution
- 41:04and storage. And remember that this is
- 41:08movable, both of them can borrow memory
- 41:10from each other. Yeah? Reserved memory,
- 41:12as we've already seen, is 300 MB.
- 41:15Therefore, this portion of the memory is
- 41:18going to be 0.4 because this was 0.6,
- 41:21right? 0.4 into
- 41:25into 10 GB minus 300 MB,
- 41:29right?
- 41:30So, this is going to be 0.4
- 41:33into 10 into 1024
- 41:36minus 300 MB,
- 41:39which is simply going to give me 37
- 41:4396 MB,
- 41:44which is 3.
- 41:46close to 3.8 GB.
- 41:49Right? So, this is how different
- 41:50sections of the memory will be
- 41:52calculated. So, that's all about the 30
- 41:55most important Apache Spark concepts. If
- 41:57you found this video insightful, I
- 42:00believe you'll also love the Spark
- 42:02performance tuning playlist and the
- 42:046-hour long Delta Lake master class on
- 42:07my YouTube channel. If you're preparing
- 42:09for interviews or upskilling in data
- 42:12engineering, these videos might be very
- 42:14helpful for you. Thank you for watching
- 42:17and I will see you in the next video.
About this transcript
This page contains the full transcript of Apache Spark Was Hard Until I Learned These 30 Concepts! by Afaque Ahmad, generated from the public captions YouTube serves with the video. The transcript has 6,045 words across 961 segments, with the original timestamps preserved so you can click any line to jump to that moment in the embedded player.
What you can do with it
Use the transcript to take notes, quote the speaker, build a study guide, generate a summary with ChatGPT or Claude via the YouTube Summary tool, or export it as a timed subtitle file with YouTube to SRT. You can also re-open it in the transcriber to translate the transcript into 100+ languages.
Free YouTube transcript tool
YouTube2Text is a free YouTube transcript generator — no signup, no daily limit. Paste any YouTube link and get the full transcript instantly, with timestamps, click-to-jump, translation to 100+ languages, AI prompts for ChatGPT, Claude, and Gemini, and exports to TXT, SRT, VTT, or Markdown.