YouTube2Text

Apache Spark Executor Tuning | Executor Cores & Memory — Transcript

by Afaque Ahmad · 5,743 words · 841 segments · language en · Watch on YouTube

Full transcript

  1. 0:03Even if you've written the bestin-class
  2. 0:06Spark code, your jobs may take forever
  3. 0:09to complete if you've not done the
  4. 0:11allocation of CPU and memory resources
  5. 0:15correctly. Hey everyone, welcome back.
  6. 0:17In this video, we are going to talk
  7. 0:19about executor tuning. basically how do
  8. 0:23you decide the number of executors that
  9. 0:26you should create and the amount of
  10. 0:29memory and the number of cores that you
  11. 0:32should be allotting to those executors.
  12. 0:35So let's first have a look at how
  13. 0:38executors are created and how executors
  14. 0:41would look like inside a node.
  15. 0:46So here I've taken one node
  16. 0:49and one node is basically one machine in
  17. 0:52a cluster and you could have several of
  18. 0:55these machines in the cluster and the
  19. 0:58configuration of this machine is that
  20. 1:02it has 17 cores and it has 20 GB of RAM.
  21. 1:08Now let's say you want to do a spark
  22. 1:11submit.
  23. 1:13And when we do a spark submit, we
  24. 1:16basically specify the number of
  25. 1:19executors we want to create, the amount
  26. 1:22of codes, and the amount of memory that
  27. 1:25we want to allot to those executors.
  28. 1:28Right? So let's let's go ahead and
  29. 1:31create three executors.
  30. 1:35Let's go ahead and create three
  31. 1:36executors. And let's say that I have
  32. 1:40assigned five cores
  33. 1:43and 6 GB of RAM to each of these
  34. 1:48executors. So in total you would end up
  35. 1:52using 5 into 3 which is 15 15 cores and
  36. 1:576 into 3 which is 18 GB of RAM from the
  37. 2:03whole cluster. Yeah. And this is how
  38. 2:05your executors would look like. So this
  39. 2:08executor, you would have three
  40. 2:10executors.
  41. 2:12You would end up having three executors.
  42. 2:14And each of them would have five cores
  43. 2:18and 6 GB of RAM.
  44. 2:21Five cores and 6 GB of RAM. Yeah. So
  45. 2:26here we saw that inside a node, this is
  46. 2:29how you're going to create executors
  47. 2:32during Spark submit. Yeah, but again
  48. 2:36whatever number I assigned over here,
  49. 2:39right, this was based on my wishful
  50. 2:41thinking, right? There are of course
  51. 2:43certain logic and certain rules that you
  52. 2:46would need to follow in order to
  53. 2:48optimally size these executor
  54. 2:52in order to optimally decide what are
  55. 2:55the number of executors you should
  56. 2:57create and how much memory and how much
  57. 3:00course you should allot to those
  58. 3:02executors. Yeah. So let's imagine you
  59. 3:05had to create and run a spark job. How
  60. 3:09would you decide what are the number of
  61. 3:11executors?
  62. 3:13What are the number of executors and the
  63. 3:16cores
  64. 3:18and memory
  65. 3:21that you would assign to those
  66. 3:24executors, right? A very critical
  67. 3:26question. So let's go ahead and take a
  68. 3:29few examples in order to understand how
  69. 3:31we would do that. Okay? So let's say you
  70. 3:34have this cluster and you have the
  71. 3:36following configuration. you have the
  72. 3:39following configuration list. So let's
  73. 3:41say you have five nodes which is five
  74. 3:43different machines and each of those
  75. 3:46machines have this configuration. They
  76. 3:50have 12 cores and 48 GB of RAM. Yeah. So
  77. 3:55basically you have five machines in a
  78. 3:58cluster you have five machines and each
  79. 4:00of those machines have 12 cores and 48
  80. 4:05GB of RAM. Now the question is how do
  81. 4:08you decide the number of executors,
  82. 4:13the number of executors
  83. 4:16that you should create when running a
  84. 4:18spark job, the course per executor and
  85. 4:23the memory per executor.
  86. 4:26Yeah. So there is a way to think about
  87. 4:30this. So you would naturally have three
  88. 4:33options as we see over here. The first
  89. 4:36option would be thin executors. The
  90. 4:39other one would be fat executors. And
  91. 4:43the last one would be optimally sized
  92. 4:46executors. Of course, by looking at
  93. 4:48this, we would always want to go ahead
  94. 4:51with optimally sized executors, right?
  95. 4:54But all three of them have advantages
  96. 4:57and disadvantages of their own. So let's
  97. 5:01first have a look at fat executors. If
  98. 5:05you were to create fat executors, how
  99. 5:08would the configuration look like and
  100. 5:11what kind of benefits or dis benefits
  101. 5:14and disadvantages that it would have?
  102. 5:16Yeah. So,
  103. 5:19let me just move my camera over here
  104. 5:23and
  105. 5:25let's go ahead and perform the
  106. 5:26calculation. Right? So, first of all,
  107. 5:29what are fat executors? Right? What are
  108. 5:32fat executors? So fat executors are
  109. 5:35those executors which occupy a large
  110. 5:39portion of the resources on a node.
  111. 5:42Yeah. So simply to keep it simply fat
  112. 5:45executors are those executors which
  113. 5:48occupy a large portion of the resources
  114. 5:52on a node. Yeah. So here we have
  115. 5:5612 core
  116. 5:5812 core and 48 GB of RAM. So fat
  117. 6:03executors are going to occupy a good
  118. 6:06portion of these resources. Yeah. So now
  119. 6:10let's go ahead and do some calculation.
  120. 6:12So in order to calculate the number of
  121. 6:14executors and the codes and all of that,
  122. 6:17we are first going to leave out one core
  123. 6:21and 1 GB of RAM for operating system
  124. 6:25Hadoop and YAN and other processes.
  125. 6:27Right? So per node we are going to leave
  126. 6:30out one core and 1 GB of RAM. So that is
  127. 6:34the first thing that we should do. So if
  128. 6:37we had if we had 12 cores and 48 GB of
  129. 6:40RAM we just subtract one. We just
  130. 6:43subtract one
  131. 6:45and what we are going to get is the
  132. 6:47number that we have over here. So this
  133. 6:49is per node
  134. 6:52you're going to have 11 cores and 47 GB
  135. 6:56of RAM. Yeah. So now what we said that
  136. 7:01fat executors are going to occupy a
  137. 7:03large portion of the resources. Yeah. So
  138. 7:05what we are going to say is one
  139. 7:07executor, one fat executor will take all
  140. 7:11of 11 cores
  141. 7:14and it is going to take all of 47 GB of
  142. 7:17RAM. So this simply means that one node
  143. 7:22is simply going to have one executor. It
  144. 7:25is going to have one executor which
  145. 7:27takes up all of the 11 cores and 47 GB
  146. 7:31of RAM and
  147. 7:33this much amount of space is left for
  148. 7:36the operating system and other
  149. 7:37processes. Yeah. So this would simply
  150. 7:41mean now that one node
  151. 7:44has one executor
  152. 7:48and our cluster has five nodes as you've
  153. 7:51seen over here.
  154. 7:54So one one one cluster has five nodes.
  155. 7:57So five nodes are going to simply have
  156. 8:00five executors.
  157. 8:03We are going to have five executors. Now
  158. 8:09the number of executors
  159. 8:13the number of executors is simply going
  160. 8:15to be five. And executor codes
  161. 8:21how much is it going to be? Take a
  162. 8:23guess. It is going to be 11. We already
  163. 8:26decided that over here. Yeah. So it is
  164. 8:28going to be 11. And executor memory
  165. 8:34executor memory is going to be 47
  166. 8:38GB over here. So this configuration is
  167. 8:41the configuration that we supply when we
  168. 8:44do a spark submit when we are creating a
  169. 8:47job.
  170. 8:50Yeah. So basically if we were to use fat
  171. 8:54executor we are going to have five fat
  172. 8:58executors. Yeah one present on each of
  173. 9:01the nodes and you see that these are
  174. 9:04very strong very powerful executors
  175. 9:07because each of them has 11 cores and 47
  176. 9:12GB of RAM. Yeah. So you're going to end
  177. 9:16up with five executors each having 11
  178. 9:19cores and 47 GB of RAM. So this is how
  179. 9:23the scenario would look like for fat
  180. 9:25executors. Now let's have a look at thin
  181. 9:28executors.
  182. 9:32So first of all, what are thin
  183. 9:33executors? Thin executors are just the
  184. 9:35opposite of fat executors. Thin
  185. 9:38executors occupy minimal resources from
  186. 9:42the node. Yeah. So they would occupy
  187. 9:45minimal
  188. 9:49resources from the node. So let's
  189. 9:52quickly do a few calculations.
  190. 9:54[clears throat] So here we've seen that
  191. 9:56we've already left out one core and 1 GB
  192. 9:59of RAM. So per node again we are left
  193. 10:03with 11 cores and 47 GB of RAM. Yeah. So
  194. 10:08one node has 11 cores and 47 GB of RAM.
  195. 10:12Now what we decide is that one executor
  196. 10:16because it contains minimal it takes up
  197. 10:19minimal resources. I'm going to only
  198. 10:21give it one core. One executor is only
  199. 10:24going to take up one core and we have 11
  200. 10:27cores. So this simply means that one
  201. 10:29node is going to end up with 11
  202. 10:33executors.
  203. 10:35Yeah, simple unitary method. We
  204. 10:40have one core per exeutor and there are
  205. 10:4411 cores in total. So there are going to
  206. 10:48be 11 executors. Yeah. Now what is the
  207. 10:54memory per executor?
  208. 10:58What is the memory per executor?
  209. 11:02So we know that we have 11 executors in
  210. 11:05total and 47 GB of RAM is all that I
  211. 11:10have on node. So 47 is my total RAM and
  212. 11:15in one node I have 11 executors. So this
  213. 11:17is somewhere going to be 4 GB.
  214. 11:21Yeah. So what we've essentially found
  215. 11:24out here is one executor
  216. 11:27is going to contain one core and 4 GBTE
  217. 11:32of RAM.
  218. 11:35Yeah. Now the last thing that we need to
  219. 11:37find out is what is the total number of
  220. 11:41cores and that is very simple to find
  221. 11:43out. So one node has
  222. 11:4711 executors as we've seen over here and
  223. 11:50we have a total of five nodes. So this
  224. 11:53simply means we are going to end up with
  225. 11:5555 executors.
  226. 11:5811 into 5 is 55. So again the
  227. 12:02configuration is going to look something
  228. 12:04like this. Num executors
  229. 12:09is going to be 55.
  230. 12:11The executor course
  231. 12:15the executor core is going to be 1 and
  232. 12:20the executor
  233. 12:23memory
  234. 12:26is going to be close to 4 GB. Yeah. So
  235. 12:29this is how the configuration is going
  236. 12:31to look like for a thin and a fat
  237. 12:34executor. Now let's understand the
  238. 12:39differences, the advantages and the
  239. 12:41disadvantages for each of them and then
  240. 12:45we would go ahead and understand how an
  241. 12:48optimally sized executor would look
  242. 12:50like. Okay. So now let's have a look at
  243. 12:52the advantages of a fat executor. So the
  244. 12:56first advantage of it is increase
  245. 12:58parallelism. So we've seen that fat
  246. 13:01executors are pretty powerful. they have
  247. 13:04a lot of cores and a lot of memory. So
  248. 13:07they'll be able to crunch a lot of data,
  249. 13:10right? So if you have a lot of cores,
  250. 13:13that means that you'll be able to run a
  251. 13:17lot of task, right? And because many
  252. 13:20codes are present on one executor, many
  253. 13:24tasks can be run together, thereby
  254. 13:26increasing the parallelism. Right? Now
  255. 13:30it is of course very beneficial because
  256. 13:35it allows you to load data which
  257. 13:37requires significant amount of memory.
  258. 13:40So tasks that require significant amount
  259. 13:44of memory can be easily consumed by
  260. 13:48consumed and processed by fat executors.
  261. 13:52Another advantage of it is if managing a
  262. 13:56lot of executors is a concern in any
  263. 13:58case. Right? We've seen that one node
  264. 14:03one node only has one executor
  265. 14:08and similarly all the other nodes would
  266. 14:11only have one executor either one or
  267. 14:14minimal executors. So in cases where
  268. 14:17managing executors is a concern you
  269. 14:20could think of fat executors.
  270. 14:23Now the other advantage is enhanced data
  271. 14:26locality. Now given that the executor
  272. 14:31already has a large memory,
  273. 14:37it is going to be able to fit a lot of
  274. 14:41partitions lot of partitions in this
  275. 14:43memory. Right? So that would mean that
  276. 14:46the data it wants to process is already
  277. 14:49local to it. It's already loaded in
  278. 14:52memory. Right? So it wouldn't need to
  279. 14:55shuffle data from the other nodes on uh
  280. 14:57in the clusters. Right? So it is for
  281. 15:00that reason that the data locality is
  282. 15:03enhanced and this overall reduces the
  283. 15:07network traffic and the overall
  284. 15:09application piece. First of all, because
  285. 15:11you don't need to move data across the
  286. 15:13cluster, right? And because of this, it
  287. 15:17improves the overall application speed.
  288. 15:21Now coming over to the disadvantages.
  289. 15:24Of course we have a lot of resources
  290. 15:26within an executor. If we don't fully
  291. 15:30utilize it, we are going to end up pay
  292. 15:33for resources which are sitting idle for
  293. 15:36resources which are not which are which
  294. 15:37you're not even using. Right? The second
  295. 15:41one is fault tolerance. So let's imagine
  296. 15:44that you have
  297. 15:47two executors and these decide on node
  298. 15:50one and node two and these two executors
  299. 15:53are crunching a lot of data. So let's
  300. 15:56say they are crunching 32 GB and 32 GB
  301. 15:59of data. Right? Now let's say for
  302. 16:02whatever reason this executor failed,
  303. 16:07something happened and this executor
  304. 16:09crashed.
  305. 16:11The amount of computation, the amount of
  306. 16:14effort that needs to be done in order to
  307. 16:16recomputee this is going to be large
  308. 16:19because it was processing a huge amount
  309. 16:21of data. It was processing 32 GB of
  310. 16:24data. So there is going to be a good
  311. 16:26amount of time in terms of computation
  312. 16:29that is going to be lost in order to
  313. 16:34recomputee this size of data right and
  314. 16:36this is going to reduce the application
  315. 16:39reliability.
  316. 16:41Now the last one is HDFS throughput. So
  317. 16:45for those of you using HDFS, HDFS
  318. 16:48throughput simply mean the rate at which
  319. 16:51you can write data to HDFS or the rate
  320. 16:56at which you can read data from HDFS.
  321. 16:59Right? Write data to HDFS or read data
  322. 17:02from HDFS. And the rate at which you can
  323. 17:05do these two is called HDFS throughput.
  324. 17:08Now if you use a lot of codes
  325. 17:12basically more than three to five cores
  326. 17:16actually more than five cores
  327. 17:19then this is going to cause a lot of
  328. 17:22garbage collection and I've discussed
  329. 17:23about garbage collection in my previous
  330. 17:25video on spark memory management. If
  331. 17:28you've not watched it please go ahead
  332. 17:29and watch it. So it is going to cause a
  333. 17:32lot of garbage collection and garbage
  334. 17:34collection is basically a process within
  335. 17:36the JVM
  336. 17:38in which if in which it cleans up your
  337. 17:42memory for unwanted objects right so if
  338. 17:45your memory is full it basically cleans
  339. 17:48up your memory of the unwanted objects
  340. 17:51and during this time it pauses your
  341. 17:53program. So it pauses your program,
  342. 17:56cleans up the memory for unwanted
  343. 17:57objects and then resumes back your
  344. 18:00program. Now imagine if this is going to
  345. 18:02happen, this pause is going to happen
  346. 18:04again and again and again. It is going
  347. 18:07to take a performance toll on your
  348. 18:10program. Right? So that is one of the
  349. 18:14reasons why you're recommended to have
  350. 18:16something between 3 to five executor. So
  351. 18:20these are some of the advantages and
  352. 18:23disadvantages of fat executors. Yeah.
  353. 18:27Okay. So now let's talk about the
  354. 18:29advantages and disadvantages of thin
  355. 18:32executors. So the first advantage is
  356. 18:36increased parallelism again. And this
  357. 18:38might confuse you a little bit because
  358. 18:39you saw increased parallelism for fat
  359. 18:42executors as well. But this is a little
  360. 18:44different, right? in the sense that you
  361. 18:48would have a lot of executors, right?
  362. 18:52You would have a lot of executors.
  363. 18:57So this is basically executor level
  364. 19:00parallelism. The last one was task level
  365. 19:03parallelism. In one executor, you had a
  366. 19:06lot of cores, right?
  367. 19:10So each core is capable of processing
  368. 19:13one task and you can ex you can process
  369. 19:16many of such task in parallel. So it was
  370. 19:19task level parallelism.
  371. 19:22Task level parallelism but this is
  372. 19:25executor level parallelism and that's
  373. 19:27how it's different. So the advantage of
  374. 19:29it is that you would still be able to
  375. 19:33process a lot of things parallelly but
  376. 19:36these jobs need to be lightweight. The
  377. 19:40amount of work that each executor is
  378. 19:43doing needs to be lightweight because
  379. 19:44you cannot stuff in a lot of memory
  380. 19:47because these guys have very small
  381. 19:49memory. So you can still do lightweight
  382. 19:53jobs. Yeah, lightweight task.
  383. 19:56Now the second advantage of it is fault
  384. 19:59tolerance. So we've seen earlier that
  385. 20:02one executor was processing a huge
  386. 20:04amount of data and then that worker
  387. 20:07crashed for some reason and because of
  388. 20:10that you lost a huge amount of data on
  389. 20:14which already computation had happened.
  390. 20:16But in this case your executors are very
  391. 20:20small right?
  392. 20:22So even if you lose an executor, it is
  393. 20:26easy to be able to recomputee
  394. 20:30whatever has been done over here
  395. 20:32already. Yeah. So the fall tolerance is
  396. 20:35pretty good when comparing it to fat
  397. 20:39executor. Now the disadvantage is high
  398. 20:43network traffic. Now because this
  399. 20:46executor has a very small memory, there
  400. 20:49are chances that the data it needs might
  401. 20:52not be fully present on this executor.
  402. 20:55So it needs to move data across the
  403. 20:58cluster in order to bring the relevant
  404. 21:01data into this executor. Yeah. And this
  405. 21:06is going to happen to all the executors
  406. 21:08within the cluster. So that is one
  407. 21:10reason why it is going to increase the
  408. 21:13network traffic. Now the last
  409. 21:14disadvantage is reduce data locality.
  410. 21:18Now we know that each of these executors
  411. 21:21have a small amount of memory. Right? So
  412. 21:24the amount of partitions it is going to
  413. 21:27be able to load in this memory is also
  414. 21:30going to be small. So the number of
  415. 21:32partitions which are local to this
  416. 21:34executor will therefore be small. So
  417. 21:36because this this memory is small, the
  418. 21:40amount of partitions that could be
  419. 21:42loaded into this memory is also going to
  420. 21:44be small. Let's say it need P10.
  421. 21:47It would need to load this in this
  422. 21:50memory. And P10 is not local to this
  423. 21:53executor, right? And the reason why it's
  424. 21:56not local because it couldn't be loaded
  425. 21:58into this memory because the memory
  426. 22:00itself was small. So this leads to a
  427. 22:04reduced data locality. Yeah. So these
  428. 22:07are the overall advantages and
  429. 22:10disadvantage of both thin and fat
  430. 22:13executors. Yeah. Okay. Now that we've
  431. 22:17seen the advantages and disadvantage of
  432. 22:20a fat and thin executor, let's try to
  433. 22:23understand how do we size or how do we
  434. 22:27create an optimal executor. Yeah. And
  435. 22:31there are a few rules that we should
  436. 22:34always keep in mind when trying to size
  437. 22:37an optimal executor. And there are four
  438. 22:40of them over here. The first one is
  439. 22:43leave out one core and 1 GB of RAM for
  440. 22:48Hadoop Yan and operating system
  441. 22:51processes. Right. So you always leave
  442. 22:54out one core and 1 GB of RAM.
  443. 22:59Yeah.
  444. 23:01This is the first one. Now the second
  445. 23:04one is the yan application master. Now
  446. 23:07yan application master is basically the
  447. 23:11one which is responsible for negotiating
  448. 23:14resources to the resource manager. So
  449. 23:16the application master basically
  450. 23:19negotiates for resources
  451. 23:22from the resource manager. So basically
  452. 23:25when you say that I want to create an
  453. 23:28executor with 11 cores and 47 GB of RAM,
  454. 23:33it is this guy who is going to go and
  455. 23:35ask the resource manager that I need
  456. 23:38these resources in order to be able to
  457. 23:40provision
  458. 23:41an executor. Yeah. So we also need to
  459. 23:45leave out something for this guy to
  460. 23:48function properly. So you can either
  461. 23:51leave out one executor
  462. 23:54or you can leave out one core and one GB
  463. 23:59of RAM. The reason why we have two
  464. 24:02options over here is because application
  465. 24:06master generally works quite well with
  466. 24:08one core and 1 GB of RAM. Yeah, but just
  467. 24:12for simplicity, you may just want to
  468. 24:14remove out one executor.
  469. 24:16But this may not be very suitable for
  470. 24:19cases where you have a fat executor.
  471. 24:21Yeah. So in fat executor you have
  472. 24:24configurations like 11 cores and 47 GB
  473. 24:28of RAM. You wouldn't want to give away
  474. 24:30such a big executor for an application
  475. 24:33master which just needs one core or and
  476. 24:361 GB of RAM. So if your executor is
  477. 24:39small just subtract one executor when
  478. 24:42you define the num executors right when
  479. 24:45you define the num executor just
  480. 24:47subtract one executor. We are going to
  481. 24:49look at these example. So either you can
  482. 24:52just subtract one executor or you can
  483. 24:56spare out one core and 1 GB of RAM. So
  484. 25:00that's the second important thing to
  485. 25:03remember. The third one is 3 to five
  486. 25:08task per executor. And we mentioned
  487. 25:11we've discussed earlier that the HDFS
  488. 25:13throughput deteriorates if we have more
  489. 25:16than five ex uh cores per executor. It
  490. 25:20leads to a lot of garbage collection.
  491. 25:23So a general rule of thumb and this is a
  492. 25:26general practice rather than just saying
  493. 25:28HDFS throughput. This is a general good
  494. 25:31practice to have three to five cores per
  495. 25:34executor.
  496. 25:38Yeah. And the last one is when you
  497. 25:42define your executor memory, right?
  498. 25:47When you define your executor memory,
  499. 25:49this executor memory should exclude
  500. 25:53the memory overhead. So the overhead
  501. 25:55memory as we discussed in one of my last
  502. 25:58videos on park memory management this is
  503. 26:01basically used for internal system
  504. 26:04system processes right. So we need to
  505. 26:07spare out some memory. So the actual
  506. 26:10executor memory should always exclude
  507. 26:16the overhead memory.
  508. 26:21And we are going to look at an example
  509. 26:23taking into consideration all of these
  510. 26:26rules. So don't worry if this this this
  511. 26:29sounds this sounds intimidating right
  512. 26:31now. Yeah. Okay. So now let's go ahead
  513. 26:34and try to understand how would we size
  514. 26:37an optimal executor. Right. So let me
  515. 26:41quickly iterate. We have a five node
  516. 26:44cluster and each of those nodes have 12
  517. 26:48cores and 48 GB of RAM. So let's go
  518. 26:52through the rules that we just walk
  519. 26:54through and follow each of them. Yeah.
  520. 26:58So the first one that we saw was we
  521. 27:01leave out one core and 1 GB of RAM for
  522. 27:05Hadoop, Pan and other operating system
  523. 27:08processes, right? So we would do the
  524. 27:12same over here. We would leave out one
  525. 27:14core and 1 GB of RAM for Hadoop and
  526. 27:18another operating system demon. Right?
  527. 27:20[snorts] And this is basically done per
  528. 27:22node level. So it's quite important to
  529. 27:25understand that this is per node. Per
  530. 27:28node we have to leave out one core and 1
  531. 27:31GB of RAM. So per node
  532. 27:36after subtracting this we are left with
  533. 27:3911 cores and 47 GB of RAM. Yeah. So this
  534. 27:45is the amount of resources that we would
  535. 27:48be left with. Now let's go ahead and
  536. 27:52follow the f the second rule. The second
  537. 27:55rule was that we leave out either one
  538. 27:59executor
  539. 28:01or we leave out one core and 1 GB of RAM
  540. 28:09for the application master and this is
  541. 28:13at the cluster level.
  542. 28:17So it's important to know that this is
  543. 28:19at the cluster level. This was at a node
  544. 28:23level. So let's go ahead and do those
  545. 28:26calculation. So before doing that
  546. 28:27calculation, let's first calculate
  547. 28:31the total
  548. 28:33memory that we have.
  549. 28:35The total memory that we have is per
  550. 28:38node we have 47 GB and we have five
  551. 28:42nodes in total. So this is going to be
  552. 28:4547 into 5 which is going to be 235.
  553. 28:50The total core
  554. 28:53the total cores is going to be 11 core
  555. 28:56in one node and we have five nodes. So
  556. 29:00it is going to be 55
  557. 29:02cores. Yeah. Now this is the resources
  558. 29:07that we have at a cluster level.
  559. 29:12So now we would simply go ahead and
  560. 29:15follow this rule which is the
  561. 29:18subtracting out either one core 1 GB of
  562. 29:20RAM or one exeutor. So we'll go ahead
  563. 29:22with one core and 1 GB of RAM. Yeah. So
  564. 29:26we subtract out 1 GB of RAM which gives
  565. 29:29me 234 GB
  566. 29:32and
  567. 29:35I subtract out one core which gives me
  568. 29:3854 cores.
  569. 29:41So this is the configuration of my
  570. 29:44cluster. Yeah. So now what I need to do
  571. 29:48is I need to find out how many executors
  572. 29:52I need to create and for each of those
  573. 29:55executors what are the cores the number
  574. 29:58of cores and the memory
  575. 30:02that I should be allotting per executor.
  576. 30:06Yeah. So the first thing that is very
  577. 30:08simplified is we can find out the number
  578. 30:11of cores by simply following one of the
  579. 30:13rules.
  580. 30:15So we already have these rules which
  581. 30:18would say that we should assign three to
  582. 30:21five ex 3 to five coursees per executor.
  583. 30:25So this should be cores. Yeah, it's
  584. 30:29interchangeable course or task because
  585. 30:31one core would process one task. So
  586. 30:33quotes per exeutor.
  587. 30:37So I would simply say that for each
  588. 30:41executor I am going to assign five cores
  589. 30:46and this is going to help me find out
  590. 30:49the total number of executor.
  591. 30:51So the total executors
  592. 30:56is simply going to be
  593. 31:02the total course which is 54 cores by
  594. 31:05the course per exeutor
  595. 31:08and this is simply going to be giving me
  596. 31:11some number which is close to 10. So I'm
  597. 31:15going to have 10 executors
  598. 31:18and each executor is going to have five
  599. 31:21cores. Yeah. So the last calculation
  600. 31:24that I now need to do is to find out the
  601. 31:27memory per executor.
  602. 31:30So in order to find out the memory per
  603. 31:33executor,
  604. 31:37I simply need to take the total memory
  605. 31:39which is 234 GB and divided by the
  606. 31:43number of executors which gives me
  607. 31:45somewhere close to 23.4
  608. 31:49GB.
  609. 31:51So each executor is going to have 23.4
  610. 31:54GB. But remember this rule. Remember the
  611. 31:58last rule that we are yet to follow was
  612. 32:00that we need to subtract the overhead
  613. 32:03memory from the executor memory and that
  614. 32:07will be the actual executor memory. And
  615. 32:09this is the calculation for it. So it's
  616. 32:12maximum of 384 MB
  617. 32:16or 10% of whatever executor memory
  618. 32:21we have. Yeah.
  619. 32:24So let's go ahead and do the
  620. 32:25calculation. So the actual
  621. 32:29memory per executor
  622. 32:34per exeutor is simply going to be
  623. 32:3923. So let me just round this off to 23
  624. 32:43GB. 23 minus maximum of 384 MB or 10% of
  625. 32:5123GB is 2.3GB.
  626. 32:54Yeah. So 23 minus 2.3GB
  627. 33:00and let's let's also round this again to
  628. 33:0420GB.
  629. 33:07Yeah. So we've rounded it down to 20 GB.
  630. 33:10So this is the actual memory per
  631. 33:13executor. Yeah. So we finally found out
  632. 33:17the number of executors.
  633. 33:22The number of executors the number of
  634. 33:24executors are 10. And let me just
  635. 33:26highlight this quickly. So the number of
  636. 33:30executors oops
  637. 33:33the number of executors are 10.
  638. 33:37The cores per executor is five cores and
  639. 33:41the memory per executor is 20 GB.
  640. 33:47So the number of executor is 10.
  641. 33:50The executor
  642. 33:54cores
  643. 33:56is five and the executor memory
  644. 34:02is 20 GB.
  645. 34:06Yeah. Now a possible question that you
  646. 34:09may have is why are we not talking about
  647. 34:12the size of data?
  648. 34:16The size of data.
  649. 34:19Why are we not talking about let's say
  650. 34:21we we are processing 10 GB of data over
  651. 34:24here or 100 GB of data over here.
  652. 34:27Wouldn't the sizes of data affect all of
  653. 34:30these calculation that we are doing over
  654. 34:32here? Yeah. So the reason the reason why
  655. 34:35I'm not talking about the sizes of data
  656. 34:38particularly is because
  657. 34:41I want you to focus on one very
  658. 34:43important metric which is
  659. 34:48the memory per core
  660. 34:52which is the memory per core and let's
  661. 34:54calculate what is the memory per core.
  662. 34:56So we here we have cores. The executor
  663. 35:00core is five and five cores are getting
  664. 35:0420 GB of memory. That mean one core is
  665. 35:08going to get 20 by 5 which is close to 4
  666. 35:12GB of memory. Yeah. So that means one
  667. 35:16core is simply going to get 4 GB of
  668. 35:19memory and one core can process one
  669. 35:23partition. Right.
  670. 35:26That means as long as my partition is
  671. 35:29less than equal to 4 GB
  672. 35:32the processing should happen seamlessly
  673. 35:34without any issues. Right? So whether is
  674. 35:3610 GB or 100 GB
  675. 35:40we need to talk about this at a very
  676. 35:43granular level. What is the size of my
  677. 35:46partition? Right? If my partition if my
  678. 35:49data is partitioned into one into sizes
  679. 35:51of 128 MB into partitioned of 128 MB
  680. 35:55this configuration that we've created
  681. 35:57over here is very good because each core
  682. 36:01can live up to 4 GB of data. Yeah. So I
  683. 36:08believe that is how you should be
  684. 36:09looking at this problem. Try to
  685. 36:12understand what is your partition size
  686. 36:14and then try to figure out that the
  687. 36:17configuration that you've created what
  688. 36:20is the partition that it can process
  689. 36:22what is the partition size that it can
  690. 36:24process right what is the memory that
  691. 36:27has been allocated to each core yeah
  692. 36:31because that memory would be the amount
  693. 36:35of data that that it would be able to
  694. 36:37process right for one partition so That
  695. 36:41was one of the most important reason but
  696. 36:43I hope all of this makes sense and it
  697. 36:44made um things clear for you right okay
  698. 36:47so to quickly summarize the benefits of
  699. 36:49an optimal executor we see that we've
  700. 36:52tried to maintain a balance between the
  701. 36:54thin executor and the fat executor right
  702. 36:57so here we see that we have assigned we
  703. 37:01have assigned five cores and 20 GB of
  704. 37:05RAM in this example right which is not
  705. 37:08too low and not too high as we've seen
  706. 37:12in a thin and a fat executor. So this is
  707. 37:15a good configuration for good
  708. 37:17parallelism right. We also avoid falling
  709. 37:22into issues with the HDFS throughput.
  710. 37:26HDFS throughput right because one of the
  711. 37:29consequence of this is that it is going
  712. 37:31to lead to large GC cycles which is
  713. 37:34going to pause your program. So we have
  714. 37:36assigned five codes. So this should be
  715. 37:39good with HDFS throughput. The last one
  716. 37:42is of course data locality and we see
  717. 37:45that the amount of memory is 20 GB of
  718. 37:49RAM. So the number of partitions this
  719. 37:52executor can hold it should be a good
  720. 37:55number, right? Should be a good number
  721. 37:56of partitions. So these partitions are
  722. 38:00going to be local to that executor. So
  723. 38:03data locality is still preserved is
  724. 38:06still enhanced. Right? So these are some
  725. 38:09of the advantages uh of an optimal size
  726. 38:13executor, right? Okay. So let's solve
  727. 38:16one more example to size an optimal
  728. 38:20executor, right? So this is the
  729. 38:24configuration of your cluster. Here you
  730. 38:27have three nodes. Each of those nodes is
  731. 38:3016 core and 48 GB of RAM. Yeah. So now
  732. 38:36let's again follow those those four
  733. 38:39rules right we are first going to leave
  734. 38:42out
  735. 38:44one core and 1 GB of RAM per node.
  736. 38:48So once we leave out one core and 1 GB
  737. 38:51of RAM per node we are going to be left
  738. 38:53with 15 cores and 47 GB of RAM and this
  739. 38:59is per node. Yeah. So now let's
  740. 39:03calculate the total memory and the total
  741. 39:06code that we are going to have in the
  742. 39:07cluster that we can use. So total
  743. 39:12cores is going to be 15 into 3 which is
  744. 39:1645 cores.
  745. 39:19And then total memory is going to be
  746. 39:2447 into 3
  747. 39:28which is going to be 73.21
  748. 39:33141. Yeah. So this is going to be 141 GB
  749. 39:37of memory. Now let's go ahead and follow
  750. 39:42the second rule which in which we are
  751. 39:44going to leave out 1 GB of
  752. 39:48RAM and one core for the application
  753. 39:52master. Yeah. So let's go ahead and do
  754. 39:55that. So we are going to leave out one
  755. 39:57core which gives me 44 cores and 1 GB of
  756. 40:02RAM which gives me 140 GP of RAM. Yeah.
  757. 40:06So this is now going to be
  758. 40:10the resources in my cluster. Yeah. So
  759. 40:14this is the cluster resources
  760. 40:17that I can use.
  761. 40:20So now let's go ahead and find out how
  762. 40:23many executors
  763. 40:25how many executors we want to create and
  764. 40:30what is the number of course that we
  765. 40:34should give to these executor and what
  766. 40:36is the memory.
  767. 40:39Yeah. So let's first find out the
  768. 40:42course. Right. So again we are going to
  769. 40:43follow this.
  770. 40:45We are going to follow this and we are
  771. 40:47going to give five executors and we also
  772. 40:50see that if we give five executor it is
  773. 40:52going to okay it's not going to fully
  774. 40:55use 44
  775. 40:5744 uh cores right so instead let's say
  776. 41:01we give four execute uh four cores so if
  777. 41:05we give four coursees per executor we
  778. 41:08are going to end up with 11 executors
  779. 41:11so the number of executor this is simply
  780. 41:13going to be total course by the core per
  781. 41:17executor.
  782. 41:20So this is going to give you 11
  783. 41:22executors.
  784. 41:24So I'm going to have 11 executors in
  785. 41:27total. Each of them is going to have
  786. 41:30four cores. And now let's see what is
  787. 41:32the amount of memory. So I have a total
  788. 41:36of 140 GB and I have 11 executor. So
  789. 41:41this is simply going to give me
  790. 41:43something very close to 1
  791. 41:4812 something GB right so let let's
  792. 41:51assume it to be 12 GB so it's going to
  793. 41:55be give me giving me some number which
  794. 41:58is very close to 12 GB so this would
  795. 42:00simply mean that the total memory that I
  796. 42:05would have is 12 GB now if you remember
  797. 42:08I also need to subtract out the memory
  798. 42:13overhead. Yeah. So the memory overhead
  799. 42:16is again going to be max of 384 MB or
  800. 42:1910% of executor memory. So the 10% is
  801. 42:23going to be 1.2 GB, right? 10% of 12 GB
  802. 42:28is going to be 1.2 GB. So overhead
  803. 42:33is going to be max of 384 MB or 10% of
  804. 42:4012 GB and the answer to this is going to
  805. 42:43be 1.2 GB. Yeah. So let let me just
  806. 42:48assume again this is I'm doing this for
  807. 42:51to make the calculations very simple. So
  808. 42:53I would just assume it to be 1 GB. Yeah.
  809. 42:57So now the actual memory is going to be
  810. 43:0312 minus overhead memory which is 1 GB
  811. 43:06it is going to be 11 GB. Yeah. So if we
  812. 43:12were to narrow down the final
  813. 43:14calculations it is going to be something
  814. 43:16like this. Num executors
  815. 43:21num executors is going to be 11 as we
  816. 43:24calculated over here.
  817. 43:27The executor core
  818. 43:31is going to be four as we calculated
  819. 43:34over here. And
  820. 43:37the executor
  821. 43:40memory
  822. 43:42is going to be
  823. 43:4611 GB.
  824. 43:48So this is how your calculation is going
  825. 43:50to look like. Yeah, it's great to see
  826. 43:53that you've reached the end of the
  827. 43:55video. So to summarize, we've learned
  828. 43:58how to size an executor and what are the
  829. 44:02approaches that you could follow in
  830. 44:04order to size these executors optimally.
  831. 44:07Right? So we've seen various options
  832. 44:09like what are fat executors, what are
  833. 44:11thin executors, what are the advantages
  834. 44:14and disadvantage of both of them and
  835. 44:16then finally how do you size an optimal
  836. 44:20executor? what are the rules that are
  837. 44:23going to govern the sizing of that
  838. 44:25optimal executor. So, I hope all of that
  839. 44:28made sense and if it did, please don't
  840. 44:30forget to like this video, share this
  841. 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.