Apache Spark Executor Tuning | Executor Cores & Memory — Transcript
Full transcript
- 0:03Even if you've written the bestin-class
- 0:06Spark code, your jobs may take forever
- 0:09to complete if you've not done the
- 0:11allocation of CPU and memory resources
- 0:15correctly. Hey everyone, welcome back.
- 0:17In this video, we are going to talk
- 0:19about executor tuning. basically how do
- 0:23you decide the number of executors that
- 0:26you should create and the amount of
- 0:29memory and the number of cores that you
- 0:32should be allotting to those executors.
- 0:35So let's first have a look at how
- 0:38executors are created and how executors
- 0:41would look like inside a node.
- 0:46So here I've taken one node
- 0:49and one node is basically one machine in
- 0:52a cluster and you could have several of
- 0:55these machines in the cluster and the
- 0:58configuration of this machine is that
- 1:02it has 17 cores and it has 20 GB of RAM.
- 1:08Now let's say you want to do a spark
- 1:11submit.
- 1:13And when we do a spark submit, we
- 1:16basically specify the number of
- 1:19executors we want to create, the amount
- 1:22of codes, and the amount of memory that
- 1:25we want to allot to those executors.
- 1:28Right? So let's let's go ahead and
- 1:31create three executors.
- 1:35Let's go ahead and create three
- 1:36executors. And let's say that I have
- 1:40assigned five cores
- 1:43and 6 GB of RAM to each of these
- 1:48executors. So in total you would end up
- 1:52using 5 into 3 which is 15 15 cores and
- 1:576 into 3 which is 18 GB of RAM from the
- 2:03whole cluster. Yeah. And this is how
- 2:05your executors would look like. So this
- 2:08executor, you would have three
- 2:10executors.
- 2:12You would end up having three executors.
- 2:14And each of them would have five cores
- 2:18and 6 GB of RAM.
- 2:21Five cores and 6 GB of RAM. Yeah. So
- 2:26here we saw that inside a node, this is
- 2:29how you're going to create executors
- 2:32during Spark submit. Yeah, but again
- 2:36whatever number I assigned over here,
- 2:39right, this was based on my wishful
- 2:41thinking, right? There are of course
- 2:43certain logic and certain rules that you
- 2:46would need to follow in order to
- 2:48optimally size these executor
- 2:52in order to optimally decide what are
- 2:55the number of executors you should
- 2:57create and how much memory and how much
- 3:00course you should allot to those
- 3:02executors. Yeah. So let's imagine you
- 3:05had to create and run a spark job. How
- 3:09would you decide what are the number of
- 3:11executors?
- 3:13What are the number of executors and the
- 3:16cores
- 3:18and memory
- 3:21that you would assign to those
- 3:24executors, right? A very critical
- 3:26question. So let's go ahead and take a
- 3:29few examples in order to understand how
- 3:31we would do that. Okay? So let's say you
- 3:34have this cluster and you have the
- 3:36following configuration. you have the
- 3:39following configuration list. So let's
- 3:41say you have five nodes which is five
- 3:43different machines and each of those
- 3:46machines have this configuration. They
- 3:50have 12 cores and 48 GB of RAM. Yeah. So
- 3:55basically you have five machines in a
- 3:58cluster you have five machines and each
- 4:00of those machines have 12 cores and 48
- 4:05GB of RAM. Now the question is how do
- 4:08you decide the number of executors,
- 4:13the number of executors
- 4:16that you should create when running a
- 4:18spark job, the course per executor and
- 4:23the memory per executor.
- 4:26Yeah. So there is a way to think about
- 4:30this. So you would naturally have three
- 4:33options as we see over here. The first
- 4:36option would be thin executors. The
- 4:39other one would be fat executors. And
- 4:43the last one would be optimally sized
- 4:46executors. Of course, by looking at
- 4:48this, we would always want to go ahead
- 4:51with optimally sized executors, right?
- 4:54But all three of them have advantages
- 4:57and disadvantages of their own. So let's
- 5:01first have a look at fat executors. If
- 5:05you were to create fat executors, how
- 5:08would the configuration look like and
- 5:11what kind of benefits or dis benefits
- 5:14and disadvantages that it would have?
- 5:16Yeah. So,
- 5:19let me just move my camera over here
- 5:23and
- 5:25let's go ahead and perform the
- 5:26calculation. Right? So, first of all,
- 5:29what are fat executors? Right? What are
- 5:32fat executors? So fat executors are
- 5:35those executors which occupy a large
- 5:39portion of the resources on a node.
- 5:42Yeah. So simply to keep it simply fat
- 5:45executors are those executors which
- 5:48occupy a large portion of the resources
- 5:52on a node. Yeah. So here we have
- 5:5612 core
- 5:5812 core and 48 GB of RAM. So fat
- 6:03executors are going to occupy a good
- 6:06portion of these resources. Yeah. So now
- 6:10let's go ahead and do some calculation.
- 6:12So in order to calculate the number of
- 6:14executors and the codes and all of that,
- 6:17we are first going to leave out one core
- 6:21and 1 GB of RAM for operating system
- 6:25Hadoop and YAN and other processes.
- 6:27Right? So per node we are going to leave
- 6:30out one core and 1 GB of RAM. So that is
- 6:34the first thing that we should do. So if
- 6:37we had if we had 12 cores and 48 GB of
- 6:40RAM we just subtract one. We just
- 6:43subtract one
- 6:45and what we are going to get is the
- 6:47number that we have over here. So this
- 6:49is per node
- 6:52you're going to have 11 cores and 47 GB
- 6:56of RAM. Yeah. So now what we said that
- 7:01fat executors are going to occupy a
- 7:03large portion of the resources. Yeah. So
- 7:05what we are going to say is one
- 7:07executor, one fat executor will take all
- 7:11of 11 cores
- 7:14and it is going to take all of 47 GB of
- 7:17RAM. So this simply means that one node
- 7:22is simply going to have one executor. It
- 7:25is going to have one executor which
- 7:27takes up all of the 11 cores and 47 GB
- 7:31of RAM and
- 7:33this much amount of space is left for
- 7:36the operating system and other
- 7:37processes. Yeah. So this would simply
- 7:41mean now that one node
- 7:44has one executor
- 7:48and our cluster has five nodes as you've
- 7:51seen over here.
- 7:54So one one one cluster has five nodes.
- 7:57So five nodes are going to simply have
- 8:00five executors.
- 8:03We are going to have five executors. Now
- 8:09the number of executors
- 8:13the number of executors is simply going
- 8:15to be five. And executor codes
- 8:21how much is it going to be? Take a
- 8:23guess. It is going to be 11. We already
- 8:26decided that over here. Yeah. So it is
- 8:28going to be 11. And executor memory
- 8:34executor memory is going to be 47
- 8:38GB over here. So this configuration is
- 8:41the configuration that we supply when we
- 8:44do a spark submit when we are creating a
- 8:47job.
- 8:50Yeah. So basically if we were to use fat
- 8:54executor we are going to have five fat
- 8:58executors. Yeah one present on each of
- 9:01the nodes and you see that these are
- 9:04very strong very powerful executors
- 9:07because each of them has 11 cores and 47
- 9:12GB of RAM. Yeah. So you're going to end
- 9:16up with five executors each having 11
- 9:19cores and 47 GB of RAM. So this is how
- 9:23the scenario would look like for fat
- 9:25executors. Now let's have a look at thin
- 9:28executors.
- 9:32So first of all, what are thin
- 9:33executors? Thin executors are just the
- 9:35opposite of fat executors. Thin
- 9:38executors occupy minimal resources from
- 9:42the node. Yeah. So they would occupy
- 9:45minimal
- 9:49resources from the node. So let's
- 9:52quickly do a few calculations.
- 9:54[clears throat] So here we've seen that
- 9:56we've already left out one core and 1 GB
- 9:59of RAM. So per node again we are left
- 10:03with 11 cores and 47 GB of RAM. Yeah. So
- 10:08one node has 11 cores and 47 GB of RAM.
- 10:12Now what we decide is that one executor
- 10:16because it contains minimal it takes up
- 10:19minimal resources. I'm going to only
- 10:21give it one core. One executor is only
- 10:24going to take up one core and we have 11
- 10:27cores. So this simply means that one
- 10:29node is going to end up with 11
- 10:33executors.
- 10:35Yeah, simple unitary method. We
- 10:40have one core per exeutor and there are
- 10:4411 cores in total. So there are going to
- 10:48be 11 executors. Yeah. Now what is the
- 10:54memory per executor?
- 10:58What is the memory per executor?
- 11:02So we know that we have 11 executors in
- 11:05total and 47 GB of RAM is all that I
- 11:10have on node. So 47 is my total RAM and
- 11:15in one node I have 11 executors. So this
- 11:17is somewhere going to be 4 GB.
- 11:21Yeah. So what we've essentially found
- 11:24out here is one executor
- 11:27is going to contain one core and 4 GBTE
- 11:32of RAM.
- 11:35Yeah. Now the last thing that we need to
- 11:37find out is what is the total number of
- 11:41cores and that is very simple to find
- 11:43out. So one node has
- 11:4711 executors as we've seen over here and
- 11:50we have a total of five nodes. So this
- 11:53simply means we are going to end up with
- 11:5555 executors.
- 11:5811 into 5 is 55. So again the
- 12:02configuration is going to look something
- 12:04like this. Num executors
- 12:09is going to be 55.
- 12:11The executor course
- 12:15the executor core is going to be 1 and
- 12:20the executor
- 12:23memory
- 12:26is going to be close to 4 GB. Yeah. So
- 12:29this is how the configuration is going
- 12:31to look like for a thin and a fat
- 12:34executor. Now let's understand the
- 12:39differences, the advantages and the
- 12:41disadvantages for each of them and then
- 12:45we would go ahead and understand how an
- 12:48optimally sized executor would look
- 12:50like. Okay. So now let's have a look at
- 12:52the advantages of a fat executor. So the
- 12:56first advantage of it is increase
- 12:58parallelism. So we've seen that fat
- 13:01executors are pretty powerful. they have
- 13:04a lot of cores and a lot of memory. So
- 13:07they'll be able to crunch a lot of data,
- 13:10right? So if you have a lot of cores,
- 13:13that means that you'll be able to run a
- 13:17lot of task, right? And because many
- 13:20codes are present on one executor, many
- 13:24tasks can be run together, thereby
- 13:26increasing the parallelism. Right? Now
- 13:30it is of course very beneficial because
- 13:35it allows you to load data which
- 13:37requires significant amount of memory.
- 13:40So tasks that require significant amount
- 13:44of memory can be easily consumed by
- 13:48consumed and processed by fat executors.
- 13:52Another advantage of it is if managing a
- 13:56lot of executors is a concern in any
- 13:58case. Right? We've seen that one node
- 14:03one node only has one executor
- 14:08and similarly all the other nodes would
- 14:11only have one executor either one or
- 14:14minimal executors. So in cases where
- 14:17managing executors is a concern you
- 14:20could think of fat executors.
- 14:23Now the other advantage is enhanced data
- 14:26locality. Now given that the executor
- 14:31already has a large memory,
- 14:37it is going to be able to fit a lot of
- 14:41partitions lot of partitions in this
- 14:43memory. Right? So that would mean that
- 14:46the data it wants to process is already
- 14:49local to it. It's already loaded in
- 14:52memory. Right? So it wouldn't need to
- 14:55shuffle data from the other nodes on uh
- 14:57in the clusters. Right? So it is for
- 15:00that reason that the data locality is
- 15:03enhanced and this overall reduces the
- 15:07network traffic and the overall
- 15:09application piece. First of all, because
- 15:11you don't need to move data across the
- 15:13cluster, right? And because of this, it
- 15:17improves the overall application speed.
- 15:21Now coming over to the disadvantages.
- 15:24Of course we have a lot of resources
- 15:26within an executor. If we don't fully
- 15:30utilize it, we are going to end up pay
- 15:33for resources which are sitting idle for
- 15:36resources which are not which are which
- 15:37you're not even using. Right? The second
- 15:41one is fault tolerance. So let's imagine
- 15:44that you have
- 15:47two executors and these decide on node
- 15:50one and node two and these two executors
- 15:53are crunching a lot of data. So let's
- 15:56say they are crunching 32 GB and 32 GB
- 15:59of data. Right? Now let's say for
- 16:02whatever reason this executor failed,
- 16:07something happened and this executor
- 16:09crashed.
- 16:11The amount of computation, the amount of
- 16:14effort that needs to be done in order to
- 16:16recomputee this is going to be large
- 16:19because it was processing a huge amount
- 16:21of data. It was processing 32 GB of
- 16:24data. So there is going to be a good
- 16:26amount of time in terms of computation
- 16:29that is going to be lost in order to
- 16:34recomputee this size of data right and
- 16:36this is going to reduce the application
- 16:39reliability.
- 16:41Now the last one is HDFS throughput. So
- 16:45for those of you using HDFS, HDFS
- 16:48throughput simply mean the rate at which
- 16:51you can write data to HDFS or the rate
- 16:56at which you can read data from HDFS.
- 16:59Right? Write data to HDFS or read data
- 17:02from HDFS. And the rate at which you can
- 17:05do these two is called HDFS throughput.
- 17:08Now if you use a lot of codes
- 17:12basically more than three to five cores
- 17:16actually more than five cores
- 17:19then this is going to cause a lot of
- 17:22garbage collection and I've discussed
- 17:23about garbage collection in my previous
- 17:25video on spark memory management. If
- 17:28you've not watched it please go ahead
- 17:29and watch it. So it is going to cause a
- 17:32lot of garbage collection and garbage
- 17:34collection is basically a process within
- 17:36the JVM
- 17:38in which if in which it cleans up your
- 17:42memory for unwanted objects right so if
- 17:45your memory is full it basically cleans
- 17:48up your memory of the unwanted objects
- 17:51and during this time it pauses your
- 17:53program. So it pauses your program,
- 17:56cleans up the memory for unwanted
- 17:57objects and then resumes back your
- 18:00program. Now imagine if this is going to
- 18:02happen, this pause is going to happen
- 18:04again and again and again. It is going
- 18:07to take a performance toll on your
- 18:10program. Right? So that is one of the
- 18:14reasons why you're recommended to have
- 18:16something between 3 to five executor. So
- 18:20these are some of the advantages and
- 18:23disadvantages of fat executors. Yeah.
- 18:27Okay. So now let's talk about the
- 18:29advantages and disadvantages of thin
- 18:32executors. So the first advantage is
- 18:36increased parallelism again. And this
- 18:38might confuse you a little bit because
- 18:39you saw increased parallelism for fat
- 18:42executors as well. But this is a little
- 18:44different, right? in the sense that you
- 18:48would have a lot of executors, right?
- 18:52You would have a lot of executors.
- 18:57So this is basically executor level
- 19:00parallelism. The last one was task level
- 19:03parallelism. In one executor, you had a
- 19:06lot of cores, right?
- 19:10So each core is capable of processing
- 19:13one task and you can ex you can process
- 19:16many of such task in parallel. So it was
- 19:19task level parallelism.
- 19:22Task level parallelism but this is
- 19:25executor level parallelism and that's
- 19:27how it's different. So the advantage of
- 19:29it is that you would still be able to
- 19:33process a lot of things parallelly but
- 19:36these jobs need to be lightweight. The
- 19:40amount of work that each executor is
- 19:43doing needs to be lightweight because
- 19:44you cannot stuff in a lot of memory
- 19:47because these guys have very small
- 19:49memory. So you can still do lightweight
- 19:53jobs. Yeah, lightweight task.
- 19:56Now the second advantage of it is fault
- 19:59tolerance. So we've seen earlier that
- 20:02one executor was processing a huge
- 20:04amount of data and then that worker
- 20:07crashed for some reason and because of
- 20:10that you lost a huge amount of data on
- 20:14which already computation had happened.
- 20:16But in this case your executors are very
- 20:20small right?
- 20:22So even if you lose an executor, it is
- 20:26easy to be able to recomputee
- 20:30whatever has been done over here
- 20:32already. Yeah. So the fall tolerance is
- 20:35pretty good when comparing it to fat
- 20:39executor. Now the disadvantage is high
- 20:43network traffic. Now because this
- 20:46executor has a very small memory, there
- 20:49are chances that the data it needs might
- 20:52not be fully present on this executor.
- 20:55So it needs to move data across the
- 20:58cluster in order to bring the relevant
- 21:01data into this executor. Yeah. And this
- 21:06is going to happen to all the executors
- 21:08within the cluster. So that is one
- 21:10reason why it is going to increase the
- 21:13network traffic. Now the last
- 21:14disadvantage is reduce data locality.
- 21:18Now we know that each of these executors
- 21:21have a small amount of memory. Right? So
- 21:24the amount of partitions it is going to
- 21:27be able to load in this memory is also
- 21:30going to be small. So the number of
- 21:32partitions which are local to this
- 21:34executor will therefore be small. So
- 21:36because this this memory is small, the
- 21:40amount of partitions that could be
- 21:42loaded into this memory is also going to
- 21:44be small. Let's say it need P10.
- 21:47It would need to load this in this
- 21:50memory. And P10 is not local to this
- 21:53executor, right? And the reason why it's
- 21:56not local because it couldn't be loaded
- 21:58into this memory because the memory
- 22:00itself was small. So this leads to a
- 22:04reduced data locality. Yeah. So these
- 22:07are the overall advantages and
- 22:10disadvantage of both thin and fat
- 22:13executors. Yeah. Okay. Now that we've
- 22:17seen the advantages and disadvantage of
- 22:20a fat and thin executor, let's try to
- 22:23understand how do we size or how do we
- 22:27create an optimal executor. Yeah. And
- 22:31there are a few rules that we should
- 22:34always keep in mind when trying to size
- 22:37an optimal executor. And there are four
- 22:40of them over here. The first one is
- 22:43leave out one core and 1 GB of RAM for
- 22:48Hadoop Yan and operating system
- 22:51processes. Right. So you always leave
- 22:54out one core and 1 GB of RAM.
- 22:59Yeah.
- 23:01This is the first one. Now the second
- 23:04one is the yan application master. Now
- 23:07yan application master is basically the
- 23:11one which is responsible for negotiating
- 23:14resources to the resource manager. So
- 23:16the application master basically
- 23:19negotiates for resources
- 23:22from the resource manager. So basically
- 23:25when you say that I want to create an
- 23:28executor with 11 cores and 47 GB of RAM,
- 23:33it is this guy who is going to go and
- 23:35ask the resource manager that I need
- 23:38these resources in order to be able to
- 23:40provision
- 23:41an executor. Yeah. So we also need to
- 23:45leave out something for this guy to
- 23:48function properly. So you can either
- 23:51leave out one executor
- 23:54or you can leave out one core and one GB
- 23:59of RAM. The reason why we have two
- 24:02options over here is because application
- 24:06master generally works quite well with
- 24:08one core and 1 GB of RAM. Yeah, but just
- 24:12for simplicity, you may just want to
- 24:14remove out one executor.
- 24:16But this may not be very suitable for
- 24:19cases where you have a fat executor.
- 24:21Yeah. So in fat executor you have
- 24:24configurations like 11 cores and 47 GB
- 24:28of RAM. You wouldn't want to give away
- 24:30such a big executor for an application
- 24:33master which just needs one core or and
- 24:361 GB of RAM. So if your executor is
- 24:39small just subtract one executor when
- 24:42you define the num executors right when
- 24:45you define the num executor just
- 24:47subtract one executor. We are going to
- 24:49look at these example. So either you can
- 24:52just subtract one executor or you can
- 24:56spare out one core and 1 GB of RAM. So
- 25:00that's the second important thing to
- 25:03remember. The third one is 3 to five
- 25:08task per executor. And we mentioned
- 25:11we've discussed earlier that the HDFS
- 25:13throughput deteriorates if we have more
- 25:16than five ex uh cores per executor. It
- 25:20leads to a lot of garbage collection.
- 25:23So a general rule of thumb and this is a
- 25:26general practice rather than just saying
- 25:28HDFS throughput. This is a general good
- 25:31practice to have three to five cores per
- 25:34executor.
- 25:38Yeah. And the last one is when you
- 25:42define your executor memory, right?
- 25:47When you define your executor memory,
- 25:49this executor memory should exclude
- 25:53the memory overhead. So the overhead
- 25:55memory as we discussed in one of my last
- 25:58videos on park memory management this is
- 26:01basically used for internal system
- 26:04system processes right. So we need to
- 26:07spare out some memory. So the actual
- 26:10executor memory should always exclude
- 26:16the overhead memory.
- 26:21And we are going to look at an example
- 26:23taking into consideration all of these
- 26:26rules. So don't worry if this this this
- 26:29sounds this sounds intimidating right
- 26:31now. Yeah. Okay. So now let's go ahead
- 26:34and try to understand how would we size
- 26:37an optimal executor. Right. So let me
- 26:41quickly iterate. We have a five node
- 26:44cluster and each of those nodes have 12
- 26:48cores and 48 GB of RAM. So let's go
- 26:52through the rules that we just walk
- 26:54through and follow each of them. Yeah.
- 26:58So the first one that we saw was we
- 27:01leave out one core and 1 GB of RAM for
- 27:05Hadoop, Pan and other operating system
- 27:08processes, right? So we would do the
- 27:12same over here. We would leave out one
- 27:14core and 1 GB of RAM for Hadoop and
- 27:18another operating system demon. Right?
- 27:20[snorts] And this is basically done per
- 27:22node level. So it's quite important to
- 27:25understand that this is per node. Per
- 27:28node we have to leave out one core and 1
- 27:31GB of RAM. So per node
- 27:36after subtracting this we are left with
- 27:3911 cores and 47 GB of RAM. Yeah. So this
- 27:45is the amount of resources that we would
- 27:48be left with. Now let's go ahead and
- 27:52follow the f the second rule. The second
- 27:55rule was that we leave out either one
- 27:59executor
- 28:01or we leave out one core and 1 GB of RAM
- 28:09for the application master and this is
- 28:13at the cluster level.
- 28:17So it's important to know that this is
- 28:19at the cluster level. This was at a node
- 28:23level. So let's go ahead and do those
- 28:26calculation. So before doing that
- 28:27calculation, let's first calculate
- 28:31the total
- 28:33memory that we have.
- 28:35The total memory that we have is per
- 28:38node we have 47 GB and we have five
- 28:42nodes in total. So this is going to be
- 28:4547 into 5 which is going to be 235.
- 28:50The total core
- 28:53the total cores is going to be 11 core
- 28:56in one node and we have five nodes. So
- 29:00it is going to be 55
- 29:02cores. Yeah. Now this is the resources
- 29:07that we have at a cluster level.
- 29:12So now we would simply go ahead and
- 29:15follow this rule which is the
- 29:18subtracting out either one core 1 GB of
- 29:20RAM or one exeutor. So we'll go ahead
- 29:22with one core and 1 GB of RAM. Yeah. So
- 29:26we subtract out 1 GB of RAM which gives
- 29:29me 234 GB
- 29:32and
- 29:35I subtract out one core which gives me
- 29:3854 cores.
- 29:41So this is the configuration of my
- 29:44cluster. Yeah. So now what I need to do
- 29:48is I need to find out how many executors
- 29:52I need to create and for each of those
- 29:55executors what are the cores the number
- 29:58of cores and the memory
- 30:02that I should be allotting per executor.
- 30:06Yeah. So the first thing that is very
- 30:08simplified is we can find out the number
- 30:11of cores by simply following one of the
- 30:13rules.
- 30:15So we already have these rules which
- 30:18would say that we should assign three to
- 30:21five ex 3 to five coursees per executor.
- 30:25So this should be cores. Yeah, it's
- 30:29interchangeable course or task because
- 30:31one core would process one task. So
- 30:33quotes per exeutor.
- 30:37So I would simply say that for each
- 30:41executor I am going to assign five cores
- 30:46and this is going to help me find out
- 30:49the total number of executor.
- 30:51So the total executors
- 30:56is simply going to be
- 31:02the total course which is 54 cores by
- 31:05the course per exeutor
- 31:08and this is simply going to be giving me
- 31:11some number which is close to 10. So I'm
- 31:15going to have 10 executors
- 31:18and each executor is going to have five
- 31:21cores. Yeah. So the last calculation
- 31:24that I now need to do is to find out the
- 31:27memory per executor.
- 31:30So in order to find out the memory per
- 31:33executor,
- 31:37I simply need to take the total memory
- 31:39which is 234 GB and divided by the
- 31:43number of executors which gives me
- 31:45somewhere close to 23.4
- 31:49GB.
- 31:51So each executor is going to have 23.4
- 31:54GB. But remember this rule. Remember the
- 31:58last rule that we are yet to follow was
- 32:00that we need to subtract the overhead
- 32:03memory from the executor memory and that
- 32:07will be the actual executor memory. And
- 32:09this is the calculation for it. So it's
- 32:12maximum of 384 MB
- 32:16or 10% of whatever executor memory
- 32:21we have. Yeah.
- 32:24So let's go ahead and do the
- 32:25calculation. So the actual
- 32:29memory per executor
- 32:34per exeutor is simply going to be
- 32:3923. So let me just round this off to 23
- 32:43GB. 23 minus maximum of 384 MB or 10% of
- 32:5123GB is 2.3GB.
- 32:54Yeah. So 23 minus 2.3GB
- 33:00and let's let's also round this again to
- 33:0420GB.
- 33:07Yeah. So we've rounded it down to 20 GB.
- 33:10So this is the actual memory per
- 33:13executor. Yeah. So we finally found out
- 33:17the number of executors.
- 33:22The number of executors the number of
- 33:24executors are 10. And let me just
- 33:26highlight this quickly. So the number of
- 33:30executors oops
- 33:33the number of executors are 10.
- 33:37The cores per executor is five cores and
- 33:41the memory per executor is 20 GB.
- 33:47So the number of executor is 10.
- 33:50The executor
- 33:54cores
- 33:56is five and the executor memory
- 34:02is 20 GB.
- 34:06Yeah. Now a possible question that you
- 34:09may have is why are we not talking about
- 34:12the size of data?
- 34:16The size of data.
- 34:19Why are we not talking about let's say
- 34:21we we are processing 10 GB of data over
- 34:24here or 100 GB of data over here.
- 34:27Wouldn't the sizes of data affect all of
- 34:30these calculation that we are doing over
- 34:32here? Yeah. So the reason the reason why
- 34:35I'm not talking about the sizes of data
- 34:38particularly is because
- 34:41I want you to focus on one very
- 34:43important metric which is
- 34:48the memory per core
- 34:52which is the memory per core and let's
- 34:54calculate what is the memory per core.
- 34:56So we here we have cores. The executor
- 35:00core is five and five cores are getting
- 35:0420 GB of memory. That mean one core is
- 35:08going to get 20 by 5 which is close to 4
- 35:12GB of memory. Yeah. So that means one
- 35:16core is simply going to get 4 GB of
- 35:19memory and one core can process one
- 35:23partition. Right.
- 35:26That means as long as my partition is
- 35:29less than equal to 4 GB
- 35:32the processing should happen seamlessly
- 35:34without any issues. Right? So whether is
- 35:3610 GB or 100 GB
- 35:40we need to talk about this at a very
- 35:43granular level. What is the size of my
- 35:46partition? Right? If my partition if my
- 35:49data is partitioned into one into sizes
- 35:51of 128 MB into partitioned of 128 MB
- 35:55this configuration that we've created
- 35:57over here is very good because each core
- 36:01can live up to 4 GB of data. Yeah. So I
- 36:08believe that is how you should be
- 36:09looking at this problem. Try to
- 36:12understand what is your partition size
- 36:14and then try to figure out that the
- 36:17configuration that you've created what
- 36:20is the partition that it can process
- 36:22what is the partition size that it can
- 36:24process right what is the memory that
- 36:27has been allocated to each core yeah
- 36:31because that memory would be the amount
- 36:35of data that that it would be able to
- 36:37process right for one partition so That
- 36:41was one of the most important reason but
- 36:43I hope all of this makes sense and it
- 36:44made um things clear for you right okay
- 36:47so to quickly summarize the benefits of
- 36:49an optimal executor we see that we've
- 36:52tried to maintain a balance between the
- 36:54thin executor and the fat executor right
- 36:57so here we see that we have assigned we
- 37:01have assigned five cores and 20 GB of
- 37:05RAM in this example right which is not
- 37:08too low and not too high as we've seen
- 37:12in a thin and a fat executor. So this is
- 37:15a good configuration for good
- 37:17parallelism right. We also avoid falling
- 37:22into issues with the HDFS throughput.
- 37:26HDFS throughput right because one of the
- 37:29consequence of this is that it is going
- 37:31to lead to large GC cycles which is
- 37:34going to pause your program. So we have
- 37:36assigned five codes. So this should be
- 37:39good with HDFS throughput. The last one
- 37:42is of course data locality and we see
- 37:45that the amount of memory is 20 GB of
- 37:49RAM. So the number of partitions this
- 37:52executor can hold it should be a good
- 37:55number, right? Should be a good number
- 37:56of partitions. So these partitions are
- 38:00going to be local to that executor. So
- 38:03data locality is still preserved is
- 38:06still enhanced. Right? So these are some
- 38:09of the advantages uh of an optimal size
- 38:13executor, right? Okay. So let's solve
- 38:16one more example to size an optimal
- 38:20executor, right? So this is the
- 38:24configuration of your cluster. Here you
- 38:27have three nodes. Each of those nodes is
- 38:3016 core and 48 GB of RAM. Yeah. So now
- 38:36let's again follow those those four
- 38:39rules right we are first going to leave
- 38:42out
- 38:44one core and 1 GB of RAM per node.
- 38:48So once we leave out one core and 1 GB
- 38:51of RAM per node we are going to be left
- 38:53with 15 cores and 47 GB of RAM and this
- 38:59is per node. Yeah. So now let's
- 39:03calculate the total memory and the total
- 39:06code that we are going to have in the
- 39:07cluster that we can use. So total
- 39:12cores is going to be 15 into 3 which is
- 39:1645 cores.
- 39:19And then total memory is going to be
- 39:2447 into 3
- 39:28which is going to be 73.21
- 39:33141. Yeah. So this is going to be 141 GB
- 39:37of memory. Now let's go ahead and follow
- 39:42the second rule which in which we are
- 39:44going to leave out 1 GB of
- 39:48RAM and one core for the application
- 39:52master. Yeah. So let's go ahead and do
- 39:55that. So we are going to leave out one
- 39:57core which gives me 44 cores and 1 GB of
- 40:02RAM which gives me 140 GP of RAM. Yeah.
- 40:06So this is now going to be
- 40:10the resources in my cluster. Yeah. So
- 40:14this is the cluster resources
- 40:17that I can use.
- 40:20So now let's go ahead and find out how
- 40:23many executors
- 40:25how many executors we want to create and
- 40:30what is the number of course that we
- 40:34should give to these executor and what
- 40:36is the memory.
- 40:39Yeah. So let's first find out the
- 40:42course. Right. So again we are going to
- 40:43follow this.
- 40:45We are going to follow this and we are
- 40:47going to give five executors and we also
- 40:50see that if we give five executor it is
- 40:52going to okay it's not going to fully
- 40:55use 44
- 40:5744 uh cores right so instead let's say
- 41:01we give four execute uh four cores so if
- 41:05we give four coursees per executor we
- 41:08are going to end up with 11 executors
- 41:11so the number of executor this is simply
- 41:13going to be total course by the core per
- 41:17executor.
- 41:20So this is going to give you 11
- 41:22executors.
- 41:24So I'm going to have 11 executors in
- 41:27total. Each of them is going to have
- 41:30four cores. And now let's see what is
- 41:32the amount of memory. So I have a total
- 41:36of 140 GB and I have 11 executor. So
- 41:41this is simply going to give me
- 41:43something very close to 1
- 41:4812 something GB right so let let's
- 41:51assume it to be 12 GB so it's going to
- 41:55be give me giving me some number which
- 41:58is very close to 12 GB so this would
- 42:00simply mean that the total memory that I
- 42:05would have is 12 GB now if you remember
- 42:08I also need to subtract out the memory
- 42:13overhead. Yeah. So the memory overhead
- 42:16is again going to be max of 384 MB or
- 42:1910% of executor memory. So the 10% is
- 42:23going to be 1.2 GB, right? 10% of 12 GB
- 42:28is going to be 1.2 GB. So overhead
- 42:33is going to be max of 384 MB or 10% of
- 42:4012 GB and the answer to this is going to
- 42:43be 1.2 GB. Yeah. So let let me just
- 42:48assume again this is I'm doing this for
- 42:51to make the calculations very simple. So
- 42:53I would just assume it to be 1 GB. Yeah.
- 42:57So now the actual memory is going to be
- 43:0312 minus overhead memory which is 1 GB
- 43:06it is going to be 11 GB. Yeah. So if we
- 43:12were to narrow down the final
- 43:14calculations it is going to be something
- 43:16like this. Num executors
- 43:21num executors is going to be 11 as we
- 43:24calculated over here.
- 43:27The executor core
- 43:31is going to be four as we calculated
- 43:34over here. And
- 43:37the executor
- 43:40memory
- 43:42is going to be
- 43:4611 GB.
- 43:48So this is how your calculation is going
- 43:50to look like. Yeah, it's great to see
- 43:53that you've reached the end of the
- 43:55video. So to summarize, we've learned
- 43:58how to size an executor and what are the
- 44:02approaches that you could follow in
- 44:04order to size these executors optimally.
- 44:07Right? So we've seen various options
- 44:09like what are fat executors, what are
- 44:11thin executors, what are the advantages
- 44:14and disadvantage of both of them and
- 44:16then finally how do you size an optimal
- 44:20executor? what are the rules that are
- 44:23going to govern the sizing of that
- 44:25optimal executor. So, I hope all of that
- 44:28made sense and if it did, please don't
- 44:30forget to like this video, share this
- 44:33video. Thank you so much for watching.
About this transcript
This page contains the full transcript of Apache Spark Executor Tuning | Executor Cores & Memory by Afaque Ahmad, generated from the public captions YouTube serves with the video. The transcript has 5,743 words across 841 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.