YouTube2Text

Shuffle Partition Spark Optimization: 10x Faster! — Transcript

by Afaque Ahmad · 2,601 words · 377 segments · language en · Watch on YouTube

Full transcript

  1. 0:00[Music]
  2. 0:03hey everyone welcome back in this video
  3. 0:06we are going to talk about chuffle
  4. 0:08partition but before we talk about
  5. 0:10Shuffle partition let's first talk about
  6. 0:14shuffling which is where Shuffle
  7. 0:16partitions come into picture so
  8. 0:19shuffling is a very popular term that we
  9. 0:21hear whenever we are working with spark
  10. 0:23right so shuffling basically happens
  11. 0:26whenever you do a wide transformation
  12. 0:29and those transformation can be anything
  13. 0:31like a group buy or a join so the idea
  14. 0:35behind shuffling is that spark tries to
  15. 0:39bring together all of the data that is
  16. 0:42related but it resides across different
  17. 0:45notes right so it brings together all of
  18. 0:48the data that is related but it
  19. 0:50currently resides across different nodes
  20. 0:53in the cluster right so this is the idea
  21. 0:56and the g behind shuffling and let's
  22. 0:58understand it with an example okay so
  23. 1:01let's have a look at this
  24. 1:03example and before we go into details
  25. 1:07let me first explain what the data set
  26. 1:10is about so here you basically have a
  27. 1:14simple data set which contains the store
  28. 1:17ID and it contains the sale
  29. 1:22amount the sale that happened at that
  30. 1:25particular store now let's ignore the
  31. 1:27dates and all of that for now right and
  32. 1:30the objective that we have over here is
  33. 1:33that we want to find out the total sales
  34. 1:37per store right so the objective is that
  35. 1:41we want to find out the total
  36. 1:44sales per
  37. 1:46store and you've been given this these
  38. 1:50two column right so the first thing that
  39. 1:55might come to your mind
  40. 1:57and for that matter it could be
  41. 1:59something like this you could simply do
  42. 2:01a DF do group by and you can simply do a
  43. 2:06store ID and then you can simply
  44. 2:09Aggregate and do the sum over here right
  45. 2:13sum of the sale amount right so over
  46. 2:15here we see that we want to do a group
  47. 2:18by right we want to do a group bu
  48. 2:22operation so let's try to understand how
  49. 2:25this entire process is going to be
  50. 2:27executed right from load in the files
  51. 2:30into the data frame till the point we to
  52. 2:33the group by and then aggregate the
  53. 2:35final values so let's say you had
  54. 2:38certain files and we basically read
  55. 2:42these files into a data frame now the
  56. 2:45moment we read it it is it was read into
  57. 2:49partition and this is the partition that
  58. 2:51you see over here P1 P2 P3 and P4 so we
  59. 2:55see that P1 basically contains data for
  60. 2:59quite a few stores it contains data for
  61. 3:01S1 S3 S2 and S4 as well right similarly
  62. 3:06you have partitions P2 P3 and
  63. 3:10P4 yeah now you have all of these data
  64. 3:14all of the data in all of these three
  65. 3:17Ford partition right all of the Ford
  66. 3:20partition now what we want to do from
  67. 3:22here is we want to group the same stores
  68. 3:25and then aggregate the final values
  69. 3:28right so in order to do that we first
  70. 3:30have to bring the data of all the same
  71. 3:33stores in the same partition right and
  72. 3:37that process is basically called
  73. 3:38shuffling so what happens is we do a
  74. 3:42shuffling over here and all of the data
  75. 3:45for S1 ends up in this
  76. 3:48partition this S1 goes over here S1 goes
  77. 3:53over here so you see P1 basically
  78. 3:56containing after Shuffle P1 containing
  79. 3:59all of the data for the store S1
  80. 4:03similarly for P2 it contain data for S2
  81. 4:07P3 contain data for S3 and P4 dat
  82. 4:10contain data for S4 so going back to the
  83. 4:13definition again shuffling basically
  84. 4:17intends to bring together the data that
  85. 4:20is related right and initially it may be
  86. 4:23residing across different nodes across
  87. 4:25the cluster so this is what we've
  88. 4:27achieved and we brought together all of
  89. 4:31the data that is related now it's very
  90. 4:33simple to be able to do a group by we
  91. 4:37simply just take all of these values
  92. 4:39over here and we do a sum so we sum up
  93. 4:43all these values all these sale
  94. 4:46amount so the group buy after the group
  95. 4:49buy we simply do an aggregation and this
  96. 4:52aggregation gives you the final sum now
  97. 4:55this partition that you see over here
  98. 4:57after shuffling these
  99. 5:00partition P1 P2 P3 and P4 these are
  100. 5:05nothing but Shuffle partition right so
  101. 5:09now a very important thing to understand
  102. 5:11is why do Shuffle partitions matter why
  103. 5:15is it even important right so let's
  104. 5:18imagine you have a 1,000 core cluster
  105. 5:22right you have
  106. 5:25a, core
  107. 5:27cluster and these are your core
  108. 5:36so let's imagine a case where you have
  109. 5:39set your default actually you don't need
  110. 5:42to set but your default Shuffle
  111. 5:48partition your default Shuffle partition
  112. 5:51equals 200 right and basically this
  113. 5:55number of Shuffle partition is being
  114. 5:57employed to run your join or a group by
  115. 6:03operation wherever a y transformation is
  116. 6:06involved right so we do know the fact
  117. 6:10that one
  118. 6:12partition is acted upon by one core
  119. 6:16right one partition is acted upon by one
  120. 6:19core now if you have 200 Shuffle
  121. 6:23partition whenever the shuffle is going
  122. 6:25to happen during the join or the group
  123. 6:27by operation these number of partitions
  124. 6:31are going to be
  125. 6:33occupied right let's say these are 200
  126. 6:37cor so there are 200 Shuffle partitions
  127. 6:41so 200 CES are going to be occupied the
  128. 6:44remaining
  129. 6:46800 the remaining 800 CES are going to
  130. 6:50sit idle right so you you realize the
  131. 6:53fact that these 800 codes are a huge
  132. 6:56number and these are going to sit idle
  133. 6:59so the repercussion for this is slow
  134. 7:02completion time of your
  135. 7:05job slow completion
  136. 7:08time that is the first one the second
  137. 7:11one is under
  138. 7:16utilization under utilization of the
  139. 7:20cluster so these are two of the most
  140. 7:23important things that is going to affect
  141. 7:26your spark
  142. 7:28jobs if if you do not manage your
  143. 7:30Shuffle partitions well right so your
  144. 7:34jobs are going to take larger even if
  145. 7:36you have a very good cluster with huge
  146. 7:38amount of resources right so it is very
  147. 7:41important to be careful about the number
  148. 7:43of shule partition that you say so let's
  149. 7:45have a look at a few scenario based
  150. 7:47questions which is going to help us
  151. 7:49understand how to tune and set the
  152. 7:52property spark. SQL do Shuffle partition
  153. 7:56right so basically how to set shuffle
  154. 7:59partition partion so the first scenario
  155. 8:01over here that we have is the data per
  156. 8:05Shuffle partition is
  157. 8:07large data per Shuffle
  158. 8:11partition is large and you would see how
  159. 8:15so let's first have a look at the
  160. 8:17parameters that are already available to
  161. 8:19us the first one is that we have five
  162. 8:23CES sorry we have five executors and
  163. 8:26each of them has four cores yeah each of
  164. 8:29them has four code we are going to use
  165. 8:33the default value the default spark.
  166. 8:35Shuffle SQL
  167. 8:37partition this
  168. 8:39is this is been set to 200 by default
  169. 8:43and the amount of data that is being
  170. 8:45shuffled is 300 gab so if you were to
  171. 8:50have a look at the spark UI you would
  172. 8:51see
  173. 8:54shuffle.
  174. 8:56right this data would be 300 megab so
  175. 8:59the data that is being shuffled is 300
  176. 9:02megab and just to give you an
  177. 9:05example the shuffle right would look
  178. 9:07something like this on the right hand
  179. 9:10side you see a column which is shuffled
  180. 9:12right so I'm referring to this column on
  181. 9:14The Spar Qi so the data that is being
  182. 9:17shuffled is 300 GB now let's do a few
  183. 9:20calculation the Total
  184. 9:23Core the Total Core that you have is 5
  185. 9:26into 4 five is the number of executors
  186. 9:30and four is the core per executor which
  187. 9:33is 20
  188. 9:35here now the shuffle partition is 200
  189. 9:39Shuffle partition is 200 which is by
  190. 9:42default we've not changed that
  191. 9:44number the data that is being shuffled
  192. 9:48is 300
  193. 9:51gab yeah now if we want to find out what
  194. 9:55is the size of data per Shuffle
  195. 9:57partition it is simply going to
  196. 10:00be size per Shuffle partition it is
  197. 10:05simply going to be the total data
  198. 10:08size let me just put this down for
  199. 10:11Simplicity the total data size by the
  200. 10:15number of Shuffle partitions here this
  201. 10:17is simply going to
  202. 10:19be 1.5
  203. 10:22gab yeah so this simply means that each
  204. 10:26of these cores that you see over here
  205. 10:29all of these cores these guys are
  206. 10:32handling 1.5 GB of data now remember
  207. 10:37that the optimal partition size optimal
  208. 10:40Shuffle partition size should be
  209. 10:43somewhere between 1 to 200
  210. 10:46megabyte it should always fall in this
  211. 10:49range now this is a very huge number now
  212. 10:53we need to tune the number of Shuffle
  213. 10:56partition in order to make sure that
  214. 10:59each score is handling an adequate
  215. 11:01amount of data yeah so the way we would
  216. 11:04do that is we would simply change the
  217. 11:07number of Shuffle partition so we would
  218. 11:09say the number of Shuffle partition is
  219. 11:12simply is going to be 300
  220. 11:19gab and I want each Shuffle partition to
  221. 11:23be of 200 megab yeah so I simply put
  222. 11:27this over here so I put the total data
  223. 11:29size by the optimal data size and this
  224. 11:33is going to give me the number of
  225. 11:34Shuffle partition so this is simply
  226. 11:37going to be
  227. 11:391,500 Shuffle partition so if we set
  228. 11:43this number to be
  229. 11:451,500 this number to be
  230. 11:481,500 what essentially we are going to
  231. 11:50get is
  232. 11:52that each of this each of this core over
  233. 11:56here is going to get only 200 MB of data
  234. 12:00and this is very much adequate that each
  235. 12:03core can handle right so what we are
  236. 12:06essentially ensuring is that the
  237. 12:09utilization for each of the core is
  238. 12:11adequate and each of the core is being
  239. 12:13given an adequate workload yeah so this
  240. 12:17is the first case where you see that the
  241. 12:19data per Shuffle partition was very
  242. 12:21large and we tuned it to a reasonable
  243. 12:24number by changing the number of Shuffle
  244. 12:27partition yeah okay so now let's have a
  245. 12:30look at the second scenario where the
  246. 12:33data per Shuffle partition is very small
  247. 12:37yeah so scenario 2 and the data per
  248. 12:41Shuffle
  249. 12:43partition is very
  250. 12:48small so let's have a look at the
  251. 12:50parameters that have already been given
  252. 12:52to us so we have three executors and
  253. 12:55four cores that means a total total of 3
  254. 13:00into 4 which is 12 cores and the data
  255. 13:04that is being shuffled is 50 megabytes
  256. 13:08so the
  257. 13:09shuffle right data is 50
  258. 13:14megab and we are not changing the
  259. 13:17shuffle partition over here the number
  260. 13:18of Shuffle partition it is been set to
  261. 13:21200 yeah so the data that is being
  262. 13:24shuffled is 50
  263. 13:27mbes
  264. 13:29the number of Shuffle partition is 200
  265. 13:32now if I were to calculate the data per
  266. 13:35Shuffle
  267. 13:37partition data per Shuffle partition it
  268. 13:40is going to be 50 megabyte by
  269. 13:44200 which is simply 250 KB and this is a
  270. 13:50very small number this is a very small
  271. 13:54part Shuffle partition size the optimal
  272. 13:57recommendation is
  273. 13:59always to have a range somewhere between
  274. 14:021 to 200
  275. 14:04mgab this is the optimal size of a
  276. 14:07shuffle partition and this is a very
  277. 14:10small number yeah so we have two options
  278. 14:12over here the first one is that we
  279. 14:15change the shuffle partition size over
  280. 14:18here yeah so the first option is we
  281. 14:20change the number of Shuffle partition
  282. 14:24we say that the data size is 50 mb and
  283. 14:27we choose any number between 1 to 200
  284. 14:30and we say that that is going to be the
  285. 14:32optimal Shuffle partition s so we say we
  286. 14:36choose let's say 10 megab if we say that
  287. 14:40we want each Shuffle partition to be of
  288. 14:4310
  289. 14:45megab so this is going to give me five
  290. 14:48Shuffle partition right so this simply
  291. 14:51gives me five Shuffle partition so what
  292. 14:53I'm going to do is that I'm just going
  293. 14:54to set this value to five and this is is
  294. 14:58going to ensure that each
  295. 15:02of each of the CES over here these guys
  296. 15:06they going to be five cores and each is
  297. 15:09going to process 10 megab of data which
  298. 15:13is a good number right so here I have
  299. 15:16ensured that the shuffle partitions are
  300. 15:19not very small and it falls within the
  301. 15:22optimal range that you see over here
  302. 15:25yeah but there is a problem the problem
  303. 15:27is that all of the other
  304. 15:30cores all of the other cores that you
  305. 15:33see over here they are sitting
  306. 15:37idle these guys are sitting idle so the
  307. 15:40other option that you have is that you
  308. 15:43could utilize all of the cores in your
  309. 15:46cluster so we've seen here that we have
  310. 15:49a total of 12 cores right we have a
  311. 15:52total of 12 cores and we know that one
  312. 15:56partition is processed by one core so
  313. 16:00what we could do instead is that we can
  314. 16:02say
  315. 16:03alternatively number of Shuffle
  316. 16:05partitions could be
  317. 16:07somewhere between
  318. 16:1050 50 megab by 12 so I'm going to say
  319. 16:14that 50 megab is my total size and I
  320. 16:17have 12 cores so I'm going to give some
  321. 16:20amount of data to each core and that
  322. 16:23amount of data is simply going to be
  323. 16:24this value which is going to be 4
  324. 16:27something right let's assume to be
  325. 16:3242 4.2 megabyte right so now what I've
  326. 16:37essentially ensured is that each of my
  327. 16:40cores is going to get 4.2 megabytes of
  328. 16:44data they're going to get 4.2 megab of
  329. 16:47data and all of them are going to be
  330. 16:52utilized and this would simply ensure
  331. 16:56that my job are completed much faster
  332. 17:00because there are more people who are
  333. 17:02working on a smaller data set right so
  334. 17:05this is again another way uh at which
  335. 17:08you can look at how your data set is
  336. 17:11structured uh what is the size of your
  337. 17:13data set what is the size of your
  338. 17:16cluster and then according to that you
  339. 17:18can tune the number of Shuffle partition
  340. 17:21you would have to set this number to 12
  341. 17:24for the
  342. 17:26second approach yeah
  343. 17:29okay while still keeping the two
  344. 17:30scenarios in mind even after adjusting
  345. 17:33the number of Shuffle partition there
  346. 17:36may be cases where your job is still
  347. 17:39running very slow right and in those
  348. 17:42cases you might need to look outside of
  349. 17:45Shuffle partitioning and think of maybe
  350. 17:47issues like data
  351. 17:50Q because in cases of data skew what
  352. 17:53happens is let's say you're doing a join
  353. 17:56operation if there is a if if the
  354. 17:59operation is queued on a particular
  355. 18:01value most of the keys are going to go
  356. 18:04to a particular partition or a set of
  357. 18:07partitions right and those partitions
  358. 18:09are going to be heavily loaded so a few
  359. 18:13cores are going to process those heavily
  360. 18:16loaded partition and it is of course
  361. 18:18going to take a lot more time because
  362. 18:20all the other cores are sitting idle
  363. 18:22right so you can think of solving these
  364. 18:24issues using aqe by enabling aqe or by
  365. 18:28using salting and you can refer to these
  366. 18:32in my other videos so yeah in these
  367. 18:35cases you would have to think of issues
  368. 18:37like this and this may not be 100%
  369. 18:41solved by tuning or adjusting the number
  370. 18:44of Shuffle partition yeah so keep all of
  371. 18:47these in mind and I hope this video give
  372. 18:50you a complete overview of What shuffle
  373. 18:52partitioning is it's great to see that
  374. 18:54you've completed the full video if you
  375. 18:56found value please don't forget to like
  376. 18:59share and subscribe thank you again for
  377. 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.