Shuffle Partition Spark Optimization: 10x Faster! — Transcript
Full transcript
- 0:00[Music]
- 0:03hey everyone welcome back in this video
- 0:06we are going to talk about chuffle
- 0:08partition but before we talk about
- 0:10Shuffle partition let's first talk about
- 0:14shuffling which is where Shuffle
- 0:16partitions come into picture so
- 0:19shuffling is a very popular term that we
- 0:21hear whenever we are working with spark
- 0:23right so shuffling basically happens
- 0:26whenever you do a wide transformation
- 0:29and those transformation can be anything
- 0:31like a group buy or a join so the idea
- 0:35behind shuffling is that spark tries to
- 0:39bring together all of the data that is
- 0:42related but it resides across different
- 0:45notes right so it brings together all of
- 0:48the data that is related but it
- 0:50currently resides across different nodes
- 0:53in the cluster right so this is the idea
- 0:56and the g behind shuffling and let's
- 0:58understand it with an example okay so
- 1:01let's have a look at this
- 1:03example and before we go into details
- 1:07let me first explain what the data set
- 1:10is about so here you basically have a
- 1:14simple data set which contains the store
- 1:17ID and it contains the sale
- 1:22amount the sale that happened at that
- 1:25particular store now let's ignore the
- 1:27dates and all of that for now right and
- 1:30the objective that we have over here is
- 1:33that we want to find out the total sales
- 1:37per store right so the objective is that
- 1:41we want to find out the total
- 1:44sales per
- 1:46store and you've been given this these
- 1:50two column right so the first thing that
- 1:55might come to your mind
- 1:57and for that matter it could be
- 1:59something like this you could simply do
- 2:01a DF do group by and you can simply do a
- 2:06store ID and then you can simply
- 2:09Aggregate and do the sum over here right
- 2:13sum of the sale amount right so over
- 2:15here we see that we want to do a group
- 2:18by right we want to do a group bu
- 2:22operation so let's try to understand how
- 2:25this entire process is going to be
- 2:27executed right from load in the files
- 2:30into the data frame till the point we to
- 2:33the group by and then aggregate the
- 2:35final values so let's say you had
- 2:38certain files and we basically read
- 2:42these files into a data frame now the
- 2:45moment we read it it is it was read into
- 2:49partition and this is the partition that
- 2:51you see over here P1 P2 P3 and P4 so we
- 2:55see that P1 basically contains data for
- 2:59quite a few stores it contains data for
- 3:01S1 S3 S2 and S4 as well right similarly
- 3:06you have partitions P2 P3 and
- 3:10P4 yeah now you have all of these data
- 3:14all of the data in all of these three
- 3:17Ford partition right all of the Ford
- 3:20partition now what we want to do from
- 3:22here is we want to group the same stores
- 3:25and then aggregate the final values
- 3:28right so in order to do that we first
- 3:30have to bring the data of all the same
- 3:33stores in the same partition right and
- 3:37that process is basically called
- 3:38shuffling so what happens is we do a
- 3:42shuffling over here and all of the data
- 3:45for S1 ends up in this
- 3:48partition this S1 goes over here S1 goes
- 3:53over here so you see P1 basically
- 3:56containing after Shuffle P1 containing
- 3:59all of the data for the store S1
- 4:03similarly for P2 it contain data for S2
- 4:07P3 contain data for S3 and P4 dat
- 4:10contain data for S4 so going back to the
- 4:13definition again shuffling basically
- 4:17intends to bring together the data that
- 4:20is related right and initially it may be
- 4:23residing across different nodes across
- 4:25the cluster so this is what we've
- 4:27achieved and we brought together all of
- 4:31the data that is related now it's very
- 4:33simple to be able to do a group by we
- 4:37simply just take all of these values
- 4:39over here and we do a sum so we sum up
- 4:43all these values all these sale
- 4:46amount so the group buy after the group
- 4:49buy we simply do an aggregation and this
- 4:52aggregation gives you the final sum now
- 4:55this partition that you see over here
- 4:57after shuffling these
- 5:00partition P1 P2 P3 and P4 these are
- 5:05nothing but Shuffle partition right so
- 5:09now a very important thing to understand
- 5:11is why do Shuffle partitions matter why
- 5:15is it even important right so let's
- 5:18imagine you have a 1,000 core cluster
- 5:22right you have
- 5:25a, core
- 5:27cluster and these are your core
- 5:36so let's imagine a case where you have
- 5:39set your default actually you don't need
- 5:42to set but your default Shuffle
- 5:48partition your default Shuffle partition
- 5:51equals 200 right and basically this
- 5:55number of Shuffle partition is being
- 5:57employed to run your join or a group by
- 6:03operation wherever a y transformation is
- 6:06involved right so we do know the fact
- 6:10that one
- 6:12partition is acted upon by one core
- 6:16right one partition is acted upon by one
- 6:19core now if you have 200 Shuffle
- 6:23partition whenever the shuffle is going
- 6:25to happen during the join or the group
- 6:27by operation these number of partitions
- 6:31are going to be
- 6:33occupied right let's say these are 200
- 6:37cor so there are 200 Shuffle partitions
- 6:41so 200 CES are going to be occupied the
- 6:44remaining
- 6:46800 the remaining 800 CES are going to
- 6:50sit idle right so you you realize the
- 6:53fact that these 800 codes are a huge
- 6:56number and these are going to sit idle
- 6:59so the repercussion for this is slow
- 7:02completion time of your
- 7:05job slow completion
- 7:08time that is the first one the second
- 7:11one is under
- 7:16utilization under utilization of the
- 7:20cluster so these are two of the most
- 7:23important things that is going to affect
- 7:26your spark
- 7:28jobs if if you do not manage your
- 7:30Shuffle partitions well right so your
- 7:34jobs are going to take larger even if
- 7:36you have a very good cluster with huge
- 7:38amount of resources right so it is very
- 7:41important to be careful about the number
- 7:43of shule partition that you say so let's
- 7:45have a look at a few scenario based
- 7:47questions which is going to help us
- 7:49understand how to tune and set the
- 7:52property spark. SQL do Shuffle partition
- 7:56right so basically how to set shuffle
- 7:59partition partion so the first scenario
- 8:01over here that we have is the data per
- 8:05Shuffle partition is
- 8:07large data per Shuffle
- 8:11partition is large and you would see how
- 8:15so let's first have a look at the
- 8:17parameters that are already available to
- 8:19us the first one is that we have five
- 8:23CES sorry we have five executors and
- 8:26each of them has four cores yeah each of
- 8:29them has four code we are going to use
- 8:33the default value the default spark.
- 8:35Shuffle SQL
- 8:37partition this
- 8:39is this is been set to 200 by default
- 8:43and the amount of data that is being
- 8:45shuffled is 300 gab so if you were to
- 8:50have a look at the spark UI you would
- 8:51see
- 8:54shuffle.
- 8:56right this data would be 300 megab so
- 8:59the data that is being shuffled is 300
- 9:02megab and just to give you an
- 9:05example the shuffle right would look
- 9:07something like this on the right hand
- 9:10side you see a column which is shuffled
- 9:12right so I'm referring to this column on
- 9:14The Spar Qi so the data that is being
- 9:17shuffled is 300 GB now let's do a few
- 9:20calculation the Total
- 9:23Core the Total Core that you have is 5
- 9:26into 4 five is the number of executors
- 9:30and four is the core per executor which
- 9:33is 20
- 9:35here now the shuffle partition is 200
- 9:39Shuffle partition is 200 which is by
- 9:42default we've not changed that
- 9:44number the data that is being shuffled
- 9:48is 300
- 9:51gab yeah now if we want to find out what
- 9:55is the size of data per Shuffle
- 9:57partition it is simply going to
- 10:00be size per Shuffle partition it is
- 10:05simply going to be the total data
- 10:08size let me just put this down for
- 10:11Simplicity the total data size by the
- 10:15number of Shuffle partitions here this
- 10:17is simply going to
- 10:19be 1.5
- 10:22gab yeah so this simply means that each
- 10:26of these cores that you see over here
- 10:29all of these cores these guys are
- 10:32handling 1.5 GB of data now remember
- 10:37that the optimal partition size optimal
- 10:40Shuffle partition size should be
- 10:43somewhere between 1 to 200
- 10:46megabyte it should always fall in this
- 10:49range now this is a very huge number now
- 10:53we need to tune the number of Shuffle
- 10:56partition in order to make sure that
- 10:59each score is handling an adequate
- 11:01amount of data yeah so the way we would
- 11:04do that is we would simply change the
- 11:07number of Shuffle partition so we would
- 11:09say the number of Shuffle partition is
- 11:12simply is going to be 300
- 11:19gab and I want each Shuffle partition to
- 11:23be of 200 megab yeah so I simply put
- 11:27this over here so I put the total data
- 11:29size by the optimal data size and this
- 11:33is going to give me the number of
- 11:34Shuffle partition so this is simply
- 11:37going to be
- 11:391,500 Shuffle partition so if we set
- 11:43this number to be
- 11:451,500 this number to be
- 11:481,500 what essentially we are going to
- 11:50get is
- 11:52that each of this each of this core over
- 11:56here is going to get only 200 MB of data
- 12:00and this is very much adequate that each
- 12:03core can handle right so what we are
- 12:06essentially ensuring is that the
- 12:09utilization for each of the core is
- 12:11adequate and each of the core is being
- 12:13given an adequate workload yeah so this
- 12:17is the first case where you see that the
- 12:19data per Shuffle partition was very
- 12:21large and we tuned it to a reasonable
- 12:24number by changing the number of Shuffle
- 12:27partition yeah okay so now let's have a
- 12:30look at the second scenario where the
- 12:33data per Shuffle partition is very small
- 12:37yeah so scenario 2 and the data per
- 12:41Shuffle
- 12:43partition is very
- 12:48small so let's have a look at the
- 12:50parameters that have already been given
- 12:52to us so we have three executors and
- 12:55four cores that means a total total of 3
- 13:00into 4 which is 12 cores and the data
- 13:04that is being shuffled is 50 megabytes
- 13:08so the
- 13:09shuffle right data is 50
- 13:14megab and we are not changing the
- 13:17shuffle partition over here the number
- 13:18of Shuffle partition it is been set to
- 13:21200 yeah so the data that is being
- 13:24shuffled is 50
- 13:27mbes
- 13:29the number of Shuffle partition is 200
- 13:32now if I were to calculate the data per
- 13:35Shuffle
- 13:37partition data per Shuffle partition it
- 13:40is going to be 50 megabyte by
- 13:44200 which is simply 250 KB and this is a
- 13:50very small number this is a very small
- 13:54part Shuffle partition size the optimal
- 13:57recommendation is
- 13:59always to have a range somewhere between
- 14:021 to 200
- 14:04mgab this is the optimal size of a
- 14:07shuffle partition and this is a very
- 14:10small number yeah so we have two options
- 14:12over here the first one is that we
- 14:15change the shuffle partition size over
- 14:18here yeah so the first option is we
- 14:20change the number of Shuffle partition
- 14:24we say that the data size is 50 mb and
- 14:27we choose any number between 1 to 200
- 14:30and we say that that is going to be the
- 14:32optimal Shuffle partition s so we say we
- 14:36choose let's say 10 megab if we say that
- 14:40we want each Shuffle partition to be of
- 14:4310
- 14:45megab so this is going to give me five
- 14:48Shuffle partition right so this simply
- 14:51gives me five Shuffle partition so what
- 14:53I'm going to do is that I'm just going
- 14:54to set this value to five and this is is
- 14:58going to ensure that each
- 15:02of each of the CES over here these guys
- 15:06they going to be five cores and each is
- 15:09going to process 10 megab of data which
- 15:13is a good number right so here I have
- 15:16ensured that the shuffle partitions are
- 15:19not very small and it falls within the
- 15:22optimal range that you see over here
- 15:25yeah but there is a problem the problem
- 15:27is that all of the other
- 15:30cores all of the other cores that you
- 15:33see over here they are sitting
- 15:37idle these guys are sitting idle so the
- 15:40other option that you have is that you
- 15:43could utilize all of the cores in your
- 15:46cluster so we've seen here that we have
- 15:49a total of 12 cores right we have a
- 15:52total of 12 cores and we know that one
- 15:56partition is processed by one core so
- 16:00what we could do instead is that we can
- 16:02say
- 16:03alternatively number of Shuffle
- 16:05partitions could be
- 16:07somewhere between
- 16:1050 50 megab by 12 so I'm going to say
- 16:14that 50 megab is my total size and I
- 16:17have 12 cores so I'm going to give some
- 16:20amount of data to each core and that
- 16:23amount of data is simply going to be
- 16:24this value which is going to be 4
- 16:27something right let's assume to be
- 16:3242 4.2 megabyte right so now what I've
- 16:37essentially ensured is that each of my
- 16:40cores is going to get 4.2 megabytes of
- 16:44data they're going to get 4.2 megab of
- 16:47data and all of them are going to be
- 16:52utilized and this would simply ensure
- 16:56that my job are completed much faster
- 17:00because there are more people who are
- 17:02working on a smaller data set right so
- 17:05this is again another way uh at which
- 17:08you can look at how your data set is
- 17:11structured uh what is the size of your
- 17:13data set what is the size of your
- 17:16cluster and then according to that you
- 17:18can tune the number of Shuffle partition
- 17:21you would have to set this number to 12
- 17:24for the
- 17:26second approach yeah
- 17:29okay while still keeping the two
- 17:30scenarios in mind even after adjusting
- 17:33the number of Shuffle partition there
- 17:36may be cases where your job is still
- 17:39running very slow right and in those
- 17:42cases you might need to look outside of
- 17:45Shuffle partitioning and think of maybe
- 17:47issues like data
- 17:50Q because in cases of data skew what
- 17:53happens is let's say you're doing a join
- 17:56operation if there is a if if the
- 17:59operation is queued on a particular
- 18:01value most of the keys are going to go
- 18:04to a particular partition or a set of
- 18:07partitions right and those partitions
- 18:09are going to be heavily loaded so a few
- 18:13cores are going to process those heavily
- 18:16loaded partition and it is of course
- 18:18going to take a lot more time because
- 18:20all the other cores are sitting idle
- 18:22right so you can think of solving these
- 18:24issues using aqe by enabling aqe or by
- 18:28using salting and you can refer to these
- 18:32in my other videos so yeah in these
- 18:35cases you would have to think of issues
- 18:37like this and this may not be 100%
- 18:41solved by tuning or adjusting the number
- 18:44of Shuffle partition yeah so keep all of
- 18:47these in mind and I hope this video give
- 18:50you a complete overview of What shuffle
- 18:52partitioning is it's great to see that
- 18:54you've completed the full video if you
- 18:56found value please don't forget to like
- 18:59share and subscribe thank you again for
- 19:02watching
About this transcript
This page contains the full transcript of Shuffle Partition Spark Optimization: 10x Faster! by Afaque Ahmad, generated from the public captions YouTube serves with the video. The transcript has 2,601 words across 377 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.