YouTube2Text

Apache Spark Was Hard Until I Learned These 30 Concepts! — Transcript

by Afaque Ahmad · 6,045 words · 961 segments · language en · Watch on YouTube

Full transcript

  1. 0:00Apache Spark is a critical topic that
  2. 0:03helped me clear interviews and land
  3. 0:05high-paying opportunities at multiple
  4. 0:08big tech companies like Apple, Uber,
  5. 0:11Atlassian, and Databricks. In this
  6. 0:13video, I am going to simplify the 30
  7. 0:17most important Spark concepts
  8. 0:19that transform the way I write,
  9. 0:23tune, and troubleshoot my Spark jobs, so
  10. 0:25that you can confidently do the same.
  11. 0:28So, let's get started. Before the world
  12. 0:30started using Apache Spark, MapReduce
  13. 0:33already existed, right? But, there were
  14. 0:36multiple performance problems with
  15. 0:38MapReduce. The most pressing problem was
  16. 0:41that MapReduce writes
  17. 0:43every stage's intermediate data to HDFS,
  18. 0:46and that is what makes it super slow.
  19. 0:49So, let's understand this with an
  20. 0:51example where we want to count the words
  21. 0:54in a file. So, let's say we have a file,
  22. 0:56and the file simply has this content.
  23. 1:00And this is going to go through the map
  24. 1:02step. The map step simply transforms
  25. 1:05each line into key-value pairs. So,
  26. 1:07basically, what it does is that it
  27. 1:09converts each of the line into key-value
  28. 1:12pairs, which looks like word, {comma}
  29. 1:15one. So, the next step is converting
  30. 1:17each of the lines into words. And this
  31. 1:20is what happens. Each of it gets split
  32. 1:22into word, hello world, hello Spark, and
  33. 1:24so on. And then,
  34. 1:26it creates key-value pairs. So, hello is
  35. 1:29going to be looking like this, hello
  36. 1:31{comma} one, similarly for world and all
  37. 1:34of this, right? So, the map step creates
  38. 1:37key-value pairs, as you see over here,
  39. 1:39and the output of the map step is going
  40. 1:42to be written back to the disk, right?
  41. 1:45So, the result is going to be written
  42. 1:48back to the disk, and this is what I was
  43. 1:49talking about earlier.
  44. 1:51Now, the next step is the shuffle and
  45. 1:54sort step.
  46. 1:56And the shuffle and sort step is going
  47. 1:59to read the data from the disk that was
  48. 2:02written in the previous step, right?
  49. 2:04In the shuffle and the sort step, the
  50. 2:06same keys are going to land at the same
  51. 2:09place or in the same partition, right?
  52. 2:11So, we see hello landing at the same
  53. 2:13partition or the same place along with
  54. 2:16all of it counts, right? And similarly
  55. 2:19for all of the other words. And the
  56. 2:21result of this is again written back to
  57. 2:23the disk. The reduce step reads the
  58. 2:27result from the previous step from the
  59. 2:30disk again, which is over here. And the
  60. 2:32job of the reduce step is to combine the
  61. 2:34output. So, it is simply going to sum
  62. 2:37these outputs or the the the individual
  63. 2:39counts, right? And then we see that
  64. 2:42hello has a final count of three, while
  65. 2:44all the others have a final count of
  66. 2:47one.
  67. 2:48And the result is finally
  68. 2:51again written back to the disk. So, this
  69. 2:53repeated read and write from the disk is
  70. 2:57what makes MapReduce super slow. The
  71. 2:59second problem, as we saw in this
  72. 3:02example, is that MapReduce has a rigid
  73. 3:06programming model in the sense that it
  74. 3:08only has two steps, which is map and
  75. 3:11reduce. So, if you have complex
  76. 3:14pipelines involving steps like filter,
  77. 3:16joins, and aggregation, you would need
  78. 3:19to break them into multiple MapReduce
  79. 3:22steps. So, how did Spark solve this
  80. 3:25problem? Spark introduced number one,
  81. 3:28in-memory processing, which means that
  82. 3:30data can be cached or be present in
  83. 3:34memory across the cluster nodes. So, in
  84. 3:38this example, what we saw was that after
  85. 3:40every intermediate step, data was
  86. 3:43flushed back to the disk, right? Data
  87. 3:46was flushed back to the disk over here,
  88. 3:48here, and here,
  89. 3:50right?
  90. 3:51Which is very costly, which is very
  91. 3:53costly. So, instead of this, across all
  92. 3:57of these steps,
  93. 3:58data can be present in the memory
  94. 4:01itself, in the RAM itself, without the
  95. 4:05need to flush back intermediate results
  96. 4:08to the disk. And this was the reason why
  97. 4:10Spark became 10 to 100 times faster.
  98. 4:13Number two, and instead of the rigid
  99. 4:16MapReduce workflow,
  100. 4:18Spark built something called the
  101. 4:19directed acyclic graph. Right? And it is
  102. 4:23going to be composed of transformations,
  103. 4:25actions, and it is going to properly
  104. 4:28divide it into something called jobs,
  105. 4:30stages, and tasks. So, it is going to
  106. 4:33look something like this.
  107. 4:35This is what exactly it is going to look
  108. 4:37like. So, don't worry, we will cover
  109. 4:39this in a lot of detail, but the idea
  110. 4:42was to tell you that it is going to be
  111. 4:44converted into something called jobs,
  112. 4:48and then stages, and then all
  113. 4:52all of these tasks. Right?
  114. 4:56So, this programming model simplifies a
  115. 4:59lot of things from the MapReduce
  116. 5:02programming model that we had earlier. I
  117. 5:04want to take a moment to talk about
  118. 5:06Educative. Educative is a fully
  119. 5:08interactive and text-based platform with
  120. 5:11excellent visuals and diagrammatic
  121. 5:14explanations, where you learn by doing.
  122. 5:16Their legendary Rocking the System
  123. 5:18Design Interview is one of the most
  124. 5:21widely accessed resources used by
  125. 5:23thousands of learners preparing for
  126. 5:25interviews at Google, Meta, Amazon, and
  127. 5:28other big tech companies. So, the course
  128. 5:30walks you through concepts for building
  129. 5:32your foundation, helps you understand
  130. 5:34how to do back-of-the-envelope
  131. 5:35calculations, and then prepares you
  132. 5:37through several examples. Right? They
  133. 5:40also have excellent courses on interview
  134. 5:43preparation, G&A, cloud platforms, labs,
  135. 5:46and mock interviews. Some of my personal
  136. 5:48favorites are Mastering MCP for advanced
  137. 5:52agentic applications and this one on
  138. 5:54cursor AI, something that [music] most
  139. 5:56of us use in our daily life. Right now,
  140. 5:58they are offering a massive 50% one-time
  141. 6:01discount and on top of that, you can use
  142. 6:03the link in the description below to get
  143. 6:06an additional 10% off. Now, let's get
  144. 6:09back to the video. Now, before we take a
  145. 6:10deep look into the Spark job itself,
  146. 6:14let's understand how and where does the
  147. 6:17Spark job get executed, right? So, it
  148. 6:20gets executed on a cluster and a cluster
  149. 6:24is simply a group of machines connected
  150. 6:28together through a network, right? So,
  151. 6:32think of it as multiple computers
  152. 6:34connected together through a network and
  153. 6:38we want to use the combined computing
  154. 6:41power of all the computers within the
  155. 6:44network, right? So, that is simply a
  156. 6:47cluster. So, let's assume that this is
  157. 6:50our cluster and our cluster has seven
  158. 6:53machines in total. So, you see six
  159. 6:55workers and one master node.
  160. 6:58And the configuration of our cluster is
  161. 7:0232 cores,
  162. 7:0432 cores for one machine and 120 GB
  163. 7:09of RAM for one machine, right? So, in
  164. 7:12total, the capacity of our cluster is 7
  165. 7:15into 32, which is 224
  166. 7:19cores
  167. 7:21and 7 into 120, which is 840
  168. 7:24GB
  169. 7:25of RAM.
  170. 7:27Right? So, this is the configuration of
  171. 7:30our cluster, right? So, there is going
  172. 7:33to be one master node, as you see over
  173. 7:36here, and the rest are worker nodes,
  174. 7:39right? So, the master node is the place
  175. 7:42where your cluster manager the cluster
  176. 7:45manager as we see over here the cluster
  177. 7:48manager is going to reside right
  178. 7:52so the cluster manager or the resource
  179. 7:54manager is the one which is responsible
  180. 7:58for managing the entire pool of the CPU
  181. 8:02and the RAM that we've got right yarn
  182. 8:05and kubernetes are the most widely used
  183. 8:08cluster manager so in our example we are
  184. 8:11simply going to stick to yarn
  185. 8:13right so in order to start the spark
  186. 8:16application the user is going to submit
  187. 8:20or the user is going to run a spark
  188. 8:23submit command along with the code the
  189. 8:26job configuration and all of that right
  190. 8:29so this request this request is going to
  191. 8:32go to yarn
  192. 8:34so this request from the user goes to
  193. 8:38yarn the resource manager is a very
  194. 8:41smart guy is a born leader and it knows
  195. 8:45how to delegate it task right so yarn
  196. 8:48the first thing that it is going to do
  197. 8:49is it is going to create the first
  198. 8:52container which is called application
  199. 8:55master so it is going to select any of
  200. 8:57the nodes the worker nodes at random and
  201. 9:00create a
  202. 9:02application master container right so
  203. 9:04let's say it creates it on worker five
  204. 9:08right
  205. 9:10so now just for your context a container
  206. 9:13is basically an isolated environment
  207. 9:15within a node with fixed amount of CPUs
  208. 9:19and RAM it is not strictly not allowed
  209. 9:23to use any resources on the node beyond
  210. 9:26what has been allocated to it right so
  211. 9:28now if we just zoom in into the
  212. 9:31application master container this is
  213. 9:34what it is going to have right this is
  214. 9:38what it is going to have.
  215. 9:41The application master container is
  216. 9:43going to start an application driver.
  217. 9:47That is what you see over here. And an
  218. 9:49application driver is nothing but a JVM
  219. 9:53process, right?
  220. 9:54Its job is to simply execute the main
  221. 9:58method of our Spark application code.
  222. 10:00So, you remember we submitted the Spark
  223. 10:02application code along with the job
  224. 10:05parameters and all of that, the
  225. 10:06configuration, right? So, the the
  226. 10:08purpose of the application driver is to
  227. 10:12execute the main method of our Spark
  228. 10:16application,
  229. 10:17right?
  230. 10:18Now, all of this is good if we've
  231. 10:21written our job in
  232. 10:23Java or Scala because they are JVM
  233. 10:26languages, right? They're JVM languages,
  234. 10:28but what if we've written our code in
  235. 10:30PySpark? JVM does
  236. 10:33doesn't just understand Python, right?
  237. 10:37So, in that case, a PySpark driver
  238. 10:41a PySpark driver is going to spin up and
  239. 10:44it is going to have its own method, own
  240. 10:47main method, the only job of which is to
  241. 10:50call the main method of the application
  242. 10:54driver, right?
  243. 10:56So, there is a PySpark driver, there is
  244. 10:57the application driver, the PySpark
  245. 10:59driver's main method calls the
  246. 11:02application driver's main method, right?
  247. 11:05And the application driver's main
  248. 11:06method, because it is written in Java or
  249. 11:09Scala, it simply comfortably runs on the
  250. 11:12JVM. A very important point to notice
  251. 11:15that Spark core is written in Scala,
  252. 11:18right? And to make it available to
  253. 11:21Python developers, it has a Java wrapper
  254. 11:24on top of it, which further has a Python
  255. 11:27wrapper on top of it called PySpark. So,
  256. 11:31the way it works is that the Python
  257. 11:33wrapper makes a call to the Java wrapper
  258. 11:36and then the Java wrapper makes a call
  259. 11:39to Spark core.
  260. 11:40Right? And this finally comfortably
  261. 11:43executes on the JVM. Right?
  262. 11:46So, going back to our
  263. 11:48application driver, now the application
  264. 11:51driver is going to execute the main
  265. 11:53method of our application.
  266. 11:55Now, it realizes that, "Okay, I need
  267. 11:57resources. Right? I need resources in
  268. 12:00order to be able to execute my Spark
  269. 12:03job." So, it looks at what were the
  270. 12:06resources that were requested and it
  271. 12:09says it understand that 30 GB of memory,
  272. 12:13four cores, and three executors
  273. 12:17were requested. Right? Three executors
  274. 12:20were requested. So, what it simply does
  275. 12:22is that it goes to YARN,
  276. 12:26which is your
  277. 12:28resource manager, and it simply tells it
  278. 12:31that, "Hey, I
  279. 12:33I need resources.
  280. 12:36Right?
  281. 12:37I need resources and the resources that
  282. 12:39I need is
  283. 12:41three executors
  284. 12:43and each of these three executors should
  285. 12:45have 30 GB of RAM
  286. 12:48and four cores. Right? And it should be
  287. 12:51three of such executors. Now, what YARN
  288. 12:55does is that it tells the application
  289. 12:58driver that, "Okay, I'm going to assign
  290. 13:00you resources for completing your job."
  291. 13:04And it is simply going to pick up
  292. 13:07randomly any three workers and create
  293. 13:10executor containers. Right? So, let's
  294. 13:13say
  295. 13:14it
  296. 13:15picks up worker number one,
  297. 13:18worker number two,
  298. 13:20and worker number four.
  299. 13:24Right? It simply creates
  300. 13:27executor containers over here. This is
  301. 13:29number one. This is number two. And this
  302. 13:33is number three, right? It simply
  303. 13:35creates them and then hands over the
  304. 13:38details to the application driver over
  305. 13:41here, right? So, let me just quickly
  306. 13:44change the color.
  307. 13:45So, it simply tells the application
  308. 13:47driver that, "Hey, these are your three
  309. 13:49executors, and here are the details.
  310. 13:52They are present on worker one,
  311. 13:54worker two, and worker four."
  312. 13:58Now, what the application driver does is
  313. 13:59that it then schedules the task across
  314. 14:03these workers. So, it then schedules the
  315. 14:05task across this one, this one, and this
  316. 14:08one. The workers
  317. 14:10um the executors then execute the task
  318. 14:12and then finally report the result back
  319. 14:16to the driver, right? Now, again, a very
  320. 14:18important point to note is that if we
  321. 14:20UDF if we've written UDFs or libraries
  322. 14:24that are not natively available in
  323. 14:26Spark, each executor container may also
  324. 14:29spin up a Python process. So, right now
  325. 14:32we see that it only has a JVM process,
  326. 14:35right? It only has a JVM process. What
  327. 14:38it might do is that it might also spin
  328. 14:40up a Python process,
  329. 14:44similar to what we see over here, right?
  330. 14:48So, along with the JVM process, it is
  331. 14:51going to spin up a Python process to
  332. 14:54help you execute your UDF code or code
  333. 14:58written in libraries that are not
  334. 15:00natively available in PySpark, right?
  335. 15:03Now, in this architecture, we see that
  336. 15:05the user or the client submits a request
  337. 15:09to the YARN resource manager, right? So,
  338. 15:12the user submits a request to the YARN
  339. 15:16resource manager, right?
  340. 15:18The resource manager then starts your
  341. 15:21application master container and the
  342. 15:24application master container then starts
  343. 15:27your application driver process, right?
  344. 15:31The application master then starts your
  345. 15:33application driver process. Now, this
  346. 15:35execution model where the application
  347. 15:39driver runs on one of the nodes in the
  348. 15:42cluster. And in our case
  349. 15:45the node of the cluster is worker number
  350. 15:49five, right?
  351. 15:51So, this execution model where the
  352. 15:53application driver runs on one of the
  353. 15:55nodes in the cluster is called the
  354. 15:59cluster mode,
  355. 16:03right?
  356. 16:04Now, there is another deployment method
  357. 16:07where the Spark submit is not routed to
  358. 16:10the YARN resource manager. Instead, when
  359. 16:12we run Spark submit, the command starts
  360. 16:15the application driver on the user or
  361. 16:18the client machine itself, right? Now,
  362. 16:20instead
  363. 16:21in the client mode, the application
  364. 16:24driver is started on the user or the
  365. 16:28client machine itself,
  366. 16:32right?
  367. 16:32And it is the application driver which
  368. 16:35is then going to coordinate with the
  369. 16:37YARN resource manager for the allocation
  370. 16:40of resources, creation of executor
  371. 16:42containers, and all of that, right? So,
  372. 16:44this execution model where the
  373. 16:47application driver runs on the client or
  374. 16:51the user machine instead of running on
  375. 16:55the worker nodes on the cluster is
  376. 16:57called the client mode.
  377. 17:02So, let's assume that our deployment has
  378. 17:04now happened using the cluster mode and
  379. 17:07our Spark application contains a mix of
  380. 17:11transformations and actions, right? So,
  381. 17:15what exactly are these transformations
  382. 17:18and action? And we'll refer back to the
  383. 17:21same code the word count example that we
  384. 17:24were referring to earlier.
  385. 17:26Transformations are simply operations
  386. 17:29that produce a new data set from an
  387. 17:31existing one, right? So, they are lazy
  388. 17:34operations simply meaning that no matter
  389. 17:37how many operations you perform, they
  390. 17:39are not executed immediately. So, all of
  391. 17:43these steps, right? The reading of the
  392. 17:45file
  393. 17:47the reading of the file
  394. 17:49splitting of each of the line into
  395. 17:51words, exploding each of the line into
  396. 17:53words, and then actually doing a group
  397. 17:56by and the count
  398. 17:58all of this, right? They are going to be
  399. 18:01appended to the list of operations
  400. 18:03needed to be executed. It's not actually
  401. 18:06going to be executed. That is what it
  402. 18:09means when I say it's lazy in nature,
  403. 18:12right? So, getting back to
  404. 18:13transformation example Examples of
  405. 18:15transformation include, let's say, a
  406. 18:17select operation or a
  407. 18:20filter operation. Let's say I had a
  408. 18:23filter operation over here saying that,
  409. 18:24"Okay, filter out all of the words which
  410. 18:27start with the letter W." Right? Or
  411. 18:31let's say you want to add a column using
  412. 18:33with column in Spark. All of these are
  413. 18:38transformations, right? So, they keep on
  414. 18:40get add getting added to the list of
  415. 18:42operations until an action is called.
  416. 18:47Right? Now, what is an action? An action
  417. 18:50is an operation that triggers actual
  418. 18:53computation, right? So, some examples of
  419. 18:55it is collect or a show or a count or a
  420. 18:59save or let's say you write to a parquet
  421. 19:01file. All of this, they are actions. So,
  422. 19:04let's assume that I would call a dot
  423. 19:07show after
  424. 19:10computing the count, this is going to be
  425. 19:13an action.
  426. 19:15Yeah?
  427. 19:16So, it's important to understand that
  428. 19:19again, transformations can be classified
  429. 19:21into narrow and wide. In narrow
  430. 19:24transformations, each output partition
  431. 19:26depends on a single input partition and
  432. 19:30there is no shuffle involved. The
  433. 19:32keyword here is there is no shuffle
  434. 19:34involved. So, we are going to understand
  435. 19:36what exactly shuffle is, but
  436. 19:39in wide transformation, a shuffle is
  437. 19:42involved. So, transformation they are of
  438. 19:43two types, narrow and wide. Narrow
  439. 19:46doesn't require shuffle. One input
  440. 19:48partition gives one output partition,
  441. 19:50but in a wide transformation, it
  442. 19:53requires a shuffle. And shuffle is
  443. 19:56another interesting activity where your
  444. 19:58data is redistributed across the
  445. 20:01clusters. And to put it simply, the same
  446. 20:04keys land on the same partition in a
  447. 20:07shuffle. So, let's understand that with
  448. 20:09an example over here.
  449. 20:11The same word count example. And don't
  450. 20:15worry about the stages,
  451. 20:17the exchange and all of this over here
  452. 20:20for now. We are going to go into a lot
  453. 20:22of details on the job stages, tasks,
  454. 20:24shuffle and all of this, right? So, we
  455. 20:26read in this file, it got split up into
  456. 20:28words, and then we got the key value
  457. 20:32pairs, right? Similar to what we were
  458. 20:34doing in MapReduce. So, let's exactly
  459. 20:36replicate the same step. Now, I want to
  460. 20:39do a group by, right? So, this means
  461. 20:42that I need to count the hellos, I need
  462. 20:45to count all of the other letters,
  463. 20:47right?
  464. 20:48So, in a shuffle, the same keys land on
  465. 20:52the same partition. So, you see all the
  466. 20:54hellos are going to land in the same
  467. 20:56partition, right? And then you are going
  468. 20:59to have the word
  469. 21:01Spark and MapReduce, right?
  470. 21:03Then finally, the count over here. So,
  471. 21:06shuffle is going to redistribute your
  472. 21:09data to make your keys. In this case,
  473. 21:12the keys is this one over here because
  474. 21:14we are doing a group by word, word
  475. 21:16becomes a key. It is going to make sure
  476. 21:18that the same keys land in the same
  477. 21:21partition, right? Now, this is becoming
  478. 21:24a little over simplified. What would
  479. 21:26generally happen is there would be a
  480. 21:28rule which would say, let's say key
  481. 21:30modulo some number, right? And let's say
  482. 21:33there are three partitions.
  483. 21:361 2 and
  484. 21:383, right? So, I would say key modulo 3
  485. 21:42and it is always going to give me a
  486. 21:43number between 0 1 and 2 and whatever
  487. 21:46number I get
  488. 21:48it is going to decide which
  489. 21:50partition the key is going to go into.
  490. 21:52So, let's say for hello,
  491. 21:54I'm going to say simply hello mod 2.
  492. 21:57Now, because this is not a number, let's
  493. 21:59say I'm going to do a hash of hello, it
  494. 22:02is going to give me an integer and this
  495. 22:05integer mod 3 is going to give me a
  496. 22:07number between 0 1 and 2
  497. 22:10and this is going to decide which
  498. 22:12partition it is
  499. 22:13going to land into. So, this makes it
  500. 22:15obvious that all the hellos are going to
  501. 22:18land in the same partition because
  502. 22:21this function is going to give the same
  503. 22:24values for all the hellos, right? So, I
  504. 22:27hope this makes sense. Now, every time a
  505. 22:30user calls an action, the driver is a
  506. 22:34very smart and organized guy. So, before
  507. 22:36it sends the task to the executors, it
  508. 22:39creates a plan.
  509. 22:42It is going to take up the code that
  510. 22:44we've written in the form of SQL or data
  511. 22:47frames and it checks whether the syntax
  512. 22:50of the code is correct or not, right?
  513. 22:53Once validated, it creates what is known
  514. 22:56as the unresolved logical plan.
  515. 22:59Yeah?
  516. 23:00Internally, it maintains something
  517. 23:03called a catalog. So, it has a structure
  518. 23:06called a catalog.
  519. 23:08And it has details about the table,
  520. 23:10databases, data types, and all of it,
  521. 23:13right? So, let's say you're writing a
  522. 23:14query which looks like something like
  523. 23:16select a {comma} b from
  524. 23:21table a.
  525. 23:24Right? So, let's say you wrote this
  526. 23:25query. What it is going to do is it
  527. 23:28verifies whether table a is present
  528. 23:32within the Spark ecosystem or not,
  529. 23:34right? And whether column a and b are a
  530. 23:38part of table a, right? So, this
  531. 23:40basically makes sure that your query is
  532. 23:43semantically correct. Now, once this is
  533. 23:46verified, it creates the logical plan.
  534. 23:50So, the unresolved logical plan goes
  535. 23:52through the catalog, and then finally it
  536. 23:54creates the logical plan. The logical
  537. 23:58plan then goes through the Catalyst
  538. 24:00Optimizer,
  539. 24:02and finally produces the optimized
  540. 24:04logical plan. Now, you can ask me that
  541. 24:07what exactly are these optimizations,
  542. 24:11right?
  543. 24:12What exactly are these optimizations?
  544. 24:15So, some examples could be filter push
  545. 24:17down,
  546. 24:18projection push down, right? Some
  547. 24:20examples can be filter push down,
  548. 24:25and projection push down.
  549. 24:28Filter push down simply means that the
  550. 24:30filtering logic is directly executed at
  551. 24:34the data source, right? So, let's say
  552. 24:36you wrote a filter
  553. 24:39city equals Boston.
  554. 24:43So, this filtering logic is directly
  555. 24:46executed at the data source without
  556. 24:48having to bring back the data to the
  557. 24:51executor, and then removing the unwanted
  558. 24:54record, right? Because then you've lost
  559. 24:56a lot of energy and bandwidth in
  560. 24:59bringing all of the chunk of records to
  561. 25:02the executors and then the executors
  562. 25:04doing the heavy lifting in removing the
  563. 25:06unwanted rows, right? So, that is where
  564. 25:08filter push down can be very, very
  565. 25:11helpful, yeah? And projection push down
  566. 25:14reduces the number of columns that you
  567. 25:17need to read that needs to be read from
  568. 25:19the data source by fetching only the
  569. 25:22columns that are written within the
  570. 25:23query, right? So, let's say if you have
  571. 25:25a query which looks like select name,
  572. 25:28{comma} age and then inside of this
  573. 25:33inside of this you have a query which
  574. 25:35says select star
  575. 25:37from whatever
  576. 25:40instead of bringing all of the data
  577. 25:43because we mentioned a select star, it
  578. 25:45is specifically going to bring the name
  579. 25:47and the age because it knows that the
  580. 25:50output only requires name and age,
  581. 25:53right? So, this is where also projection
  582. 25:56push down reduces the amount of data
  583. 25:59that we are pulling and executing on our
  584. 26:02executors. So, now after the
  585. 26:04optimization that I said and they have
  586. 26:06been applied, it creates an optimized
  587. 26:09logical plan and this optimized logical
  588. 26:12plan is converted into several physical
  589. 26:16plans, right? But before choosing which
  590. 26:19physical plan is going to run on the
  591. 26:20cluster all of these generated physical
  592. 26:24plans all of these generated physical
  593. 26:26plans go through something which is the
  594. 26:29cost model, right? It goes through the
  595. 26:32cost model
  596. 26:34which is going to decide which is the
  597. 26:36most optimal plan that can be run on the
  598. 26:39cluster, right? So, once it goes through
  599. 26:41the cost model, the final physical plan
  600. 26:44comes out and now the DAG scheduler
  601. 26:46kicks in.
  602. 26:48Right? So, the DAG scheduler kicks in
  603. 26:52and it creates a job. It takes up this
  604. 26:55final physical plan and it converts it
  605. 26:58into something called a DAG of stages
  606. 27:02and task and finally executes it on the
  607. 27:05cluster. Now, what exactly are the jobs,
  608. 27:08stages, and tasks? So, let's understand
  609. 27:11this with this example. Let's assume
  610. 27:13that we submitted this piece of code.
  611. 27:16So,
  612. 27:17what exactly is happening here is we are
  613. 27:19reading a parquet file and then we are
  614. 27:21applying a filter uh where the amount is
  615. 27:24greater than 1,500. And then we do a
  616. 27:28repartition of three, meaning that if we
  617. 27:30read in one partition, it is now going
  618. 27:33to be split into three partition, right?
  619. 27:35Then we select the region and amount and
  620. 27:39then we add a new column which is tax
  621. 27:42and finally we do a group by by the
  622. 27:45region and sum the sale amount and the
  623. 27:49tax amount, right?
  624. 27:50The sale amount over here and the tax
  625. 27:52amount. And lastly, we finally add what
  626. 27:55exactly was the timestamp where this
  627. 27:57computation happened. And we trigger an
  628. 27:59action which is collect. So, let's first
  629. 28:02break this down into different steps,
  630. 28:04right?
  631. 28:05So, the first step over here is the
  632. 28:08reading of the file which is spark.read.
  633. 28:11Number one. And then the second step
  634. 28:13over here is a filter. The third step
  635. 28:17over here is a repartition,
  636. 28:19right? The fourth step over here is a
  637. 28:22select. The fifth step over here is a
  638. 28:26withColumn where we add in the tax
  639. 28:30column.
  640. 28:31The next step over here is group by
  641. 28:33which is sixth step.
  642. 28:35Then we do an aggregation
  643. 28:37which is the seventh step and then the
  644. 28:40last step over here before
  645. 28:42the action is a withColumn and then
  646. 28:45finally we do a collect, which is an
  647. 28:49action, right? Now, the the
  648. 28:53the ones which are marked in gray, these
  649. 28:56are narrow transformations, right?
  650. 28:59These are narrow transformations.
  651. 29:04Because they don't require any shuffle,
  652. 29:06right? Repartition number three and
  653. 29:08group by, these are wide transformations
  654. 29:12because they are going to require a
  655. 29:14shuffle, right? So, a very important
  656. 29:16point to note is
  657. 29:19of course, we understood about narrow
  658. 29:21transformations.
  659. 29:23Narrow {slash} wide transformations.
  660. 29:30Number one, number two, what exactly are
  661. 29:33the job stages and tasks, right? So,
  662. 29:35whenever we see an action being called,
  663. 29:38for example, {dot} collect, right?
  664. 29:40Uh this is an action being invoked by
  665. 29:43the user. Spark is going to create job
  666. 29:46in order to get the result of this
  667. 29:49action, right? So, this is going to
  668. 29:51create a job where all of these eight
  669. 29:54steps are going to be executed to
  670. 29:56produce the output and give you back the
  671. 29:59result, right?
  672. 30:00Now, that is a job. Now, stages
  673. 30:03are created at the boundary of a
  674. 30:06shuffle. So, this is where a stage is
  675. 30:08going to be created, right? Stage one is
  676. 30:11going to be created. So, this is all a
  677. 30:13part of stage one.
  678. 30:15And then a stage is going to be created
  679. 30:17over here again.
  680. 30:19Stage two. This is going to be a part of
  681. 30:23stage two. And this is going to be stage
  682. 30:26three.
  683. 30:28So, if you have n shuffles
  684. 30:31or n wide transformations,
  685. 30:34n shuffles or wide transformations,
  686. 30:40you are going to have n plus one
  687. 30:45stages.
  688. 30:46So, in this case, we have two wide
  689. 30:48transformations.
  690. 30:50We are going to have three stages. So,
  691. 30:53stage one over here
  692. 30:56is going to look like this.
  693. 31:00This one over here. Stage two, what we
  694. 31:03discussed over here is going to look
  695. 31:05like this.
  696. 31:06And the final stage three
  697. 31:10is
  698. 31:11the one over here, right? So, let's go
  699. 31:14through this example
  700. 31:16through all of the items in the stages,
  701. 31:19right? The first step was a read, and
  702. 31:21then a filter, and then a repartition,
  703. 31:23right? Let's say we read in this file,
  704. 31:26which is the first step, and then we do
  705. 31:28a filter, where filter amount was
  706. 31:30greater than 1,500, and that is what it
  707. 31:33leaves you with over here. And then we
  708. 31:35did a repartition.
  709. 31:38So, you see that we read in one
  710. 31:40partition over here,
  711. 31:42but now we repartition it into three
  712. 31:44partitions. So, this is what number one,
  713. 31:47partition number two, and number three.
  714. 31:49Three partitions, right? And this get
  715. 31:52written to the shuffle write exchange.
  716. 31:55So, this is basically a buffer storage,
  717. 31:57which is where your data gets written
  718. 32:00when shuffling is happening, right? So,
  719. 32:02that's the first step, first stage. And
  720. 32:06then we go to the second stage, where
  721. 32:08all of this is read in, and we are going
  722. 32:11to do a select. The select was
  723. 32:15to select the region and the amount.
  724. 32:18And then we added a tax column, right?
  725. 32:20You remember that the tax the the
  726. 32:22formula for the tax
  727. 32:24tax column was 0.8 into the amount, and
  728. 32:27that is what we are going to apply over
  729. 32:28here.
  730. 32:31So, we simply select the relevant
  731. 32:33column, which is these two columns, and
  732. 32:35then we add a new column called tags,
  733. 32:38and this is going to apply for all of
  734. 32:40the three partitions.
  735. 32:42Right? These two steps, number one and
  736. 32:44number two. Right? And then a group by
  737. 32:47is going to happen where we are going to
  738. 32:49write to the shuffle write exchange.
  739. 32:52Now, a very important thing to note is
  740. 32:55these steps, right?
  741. 32:57These steps are executed on each of the
  742. 33:00partition. Right? This is going to be
  743. 33:02executed on this partition, on this
  744. 33:04partition, on and on this partition.
  745. 33:06Right? So, when it's executed on one
  746. 33:09partition, it is called a task. Right?
  747. 33:12And one task is operated on one
  748. 33:15partition
  749. 33:17by one core.
  750. 33:20Right? To keep it To keep things simple,
  751. 33:22by one core. Right? So, we have three
  752. 33:24partitions over here, number one, two,
  753. 33:27and three. So, there are going to be
  754. 33:29three tasks.
  755. 33:31So, remember that if you have n
  756. 33:34partitions,
  757. 33:37there are going to be n parallel
  758. 33:40tasks.
  759. 33:42Right?
  760. 33:42So, these can be executed parallelly
  761. 33:44because there's no dependency of one on
  762. 33:46the other. The select and with column on
  763. 33:49one partition can completely go on in
  764. 33:52parallel with this partition and this
  765. 33:54partition. Right?
  766. 33:56So, then we do a group by, and this data
  767. 33:59is written back to the shuffle write
  768. 34:01exchange, and this is written in as
  769. 34:04three partitions.
  770. 34:06Yeah?
  771. 34:06So, when we do a group by, the first
  772. 34:08thing that happens is a shuffling, and
  773. 34:10in shuffling, the same keys come to the
  774. 34:12same partition, and you see that Kolkata
  775. 34:16has come to the same partition. Earlier,
  776. 34:17it was over here and over here. Now, it
  777. 34:19has come to the same partition. Right?
  778. 34:21It gets written to the shuffle write
  779. 34:23exchange, and then it is read up
  780. 34:27in the shuffle read exchange, which is
  781. 34:29stage number three over here. Right?
  782. 34:32Now, once this is read up, we do a sum
  783. 34:35and it is simply going to sum up the
  784. 34:37total amount
  785. 34:38and the tax over here. And this is going
  786. 34:41to happen for all of the three
  787. 34:43partition, right? And the last
  788. 34:46step is a width column. The width column
  789. 34:48is simply going to add a timestamp and
  790. 34:51this gives us the final
  791. 34:54result.
  792. 34:56Right? So, what we've seen here is, to
  793. 34:59summarize, we called an action as a
  794. 35:01result of it and that action was this
  795. 35:04collect action over here. As a result of
  796. 35:06it, a job was created and when a job was
  797. 35:09created, we figured out the narrow and
  798. 35:12the wide dependencies,
  799. 35:14the wide transformation, and we figured
  800. 35:16out that there are going to be three
  801. 35:20stages.
  802. 35:21And this is how the stages are created.
  803. 35:23Number one, number two, and number
  804. 35:26three.
  805. 35:27Right? And we created three partitions,
  806. 35:30because of which we had three parallel
  807. 35:33task in place. So, now I believe you
  808. 35:35understand how job stages and task are
  809. 35:39created. So, finally, you must be
  810. 35:41thinking that Spark does all of this
  811. 35:44magic in memory and what exactly goes
  812. 35:47behind the scenes in memory to enable
  813. 35:51Spark to be able to do all of this
  814. 35:53magic. So, we've seen YARN spin up
  815. 35:56executor containers on the worker nodes,
  816. 35:59right? Now, let's have a deeper look
  817. 36:01into the memory section of the executor
  818. 36:05container, yeah? So, the container has
  819. 36:09three important regions, right? The
  820. 36:11first one is the on-heap memory, right?
  821. 36:15This is the most important section. This
  822. 36:17is the main area which is managed by the
  823. 36:21JVM and this is where most of the Spark
  824. 36:24operations run, right? And this is
  825. 36:26divided into four sections, which as you
  826. 36:29see over here is execution memory,
  827. 36:31storage memory, user memory, and
  828. 36:34reserved memory. So, execution memory is
  829. 36:37the place where your joins, your sorts,
  830. 36:40your shuffle, and your aggregation, some
  831. 36:43of the things that we do very frequently
  832. 36:45in our code, all of this happen in the
  833. 36:48execution memory, right? The storage
  834. 36:51memory is the place where caching, let's
  835. 36:54say we've cached our data frame or we've
  836. 36:56created broadcast variable, we are doing
  837. 36:58broadcast joins, right? So, this is the
  838. 37:01place where all of those entities live,
  839. 37:04where our cached variable,
  840. 37:06where our cached data frame, where our
  841. 37:08broadcast variables live, right? The
  842. 37:10third one is user memory, yeah? User
  843. 37:14memory is the one which is used by our
  844. 37:17own program's objects, right? So, let's
  845. 37:19say we create a lot of list, variables,
  846. 37:22user-defined functions, and all of that.
  847. 37:25All of that is going to live in the user
  848. 37:28memory. And lastly, is the reserved
  849. 37:31memory. Reserved memory is a very small
  850. 37:34slice. So, Spark keeps it for its own
  851. 37:38internal housekeeping and working. The
  852. 37:40next important section is the off-heap
  853. 37:43memory,
  854. 37:45which we see over here. This is
  855. 37:47generally disabled by default. It is
  856. 37:50used when you enable
  857. 37:52this setting using
  858. 37:53spark.memory.offHeapEnabled,
  859. 37:56and actually give it a size by using
  860. 37:58spark.memory.offHeap.size, right? So,
  861. 38:00this this section of the memory
  862. 38:05is useful when Spark is handling large
  863. 38:08data sets or doing heavy shuffle, and
  864. 38:11frequent garbage collection is
  865. 38:13happening, right? Frequent garbage
  866. 38:14collection generally happens when
  867. 38:17there's a lot of aggravation of objects
  868. 38:20within the memory, and Spark needs to
  869. 38:23free up that memory in order to be able
  870. 38:25to use it for subsequent steps. That is
  871. 38:28when garbage collection happens, so that
  872. 38:31it removes those unwanted objects,
  873. 38:34right? And garbage collection is a
  874. 38:36costly operation because it stops your
  875. 38:38program and then does the cleanup and
  876. 38:40then resumes your program, right? So, in
  877. 38:44such a case where frequent GC or garbage
  878. 38:47collection is happening, the data can be
  879. 38:49kept outside the JVM on the off-heap
  880. 38:53memory because the off-heap memory is
  881. 38:55generally immune to garbage collection.
  882. 38:58But, then it's very important to note
  883. 39:00then that you would have to do them the
  884. 39:03the cleaning of the unwanted objects in
  885. 39:05order to avoid memory leaks, right?
  886. 39:08If you want to completely avoid garbage
  887. 39:10collection, right? So, in this case
  888. 39:12specifically, you can use the off-heap
  889. 39:15memory.
  890. 39:16The third and the last part is the
  891. 39:19overhead memory.
  892. 39:20The overhead memory is the extra space
  893. 39:23that Spark requests from YARN
  894. 39:27in order to handle its internal
  895. 39:29system-level operations, right? Now, a
  896. 39:31very interesting thing to note is this
  897. 39:33execution plus uni- plus the storage
  898. 39:37memory is together called the unified
  899. 39:41memory
  900. 39:43because the slider that you see, right?
  901. 39:46The slider that you see is actually
  902. 39:48movable. So, they both can borrow memory
  903. 39:51from each other, but more preference is
  904. 39:54given to the execution memory because,
  905. 39:56of course, most of the important
  906. 39:58operation within our Spark lifecycle
  907. 40:01happen in the execution memory, right?
  908. 40:04So, let me take a quick example and
  909. 40:06explain you how each of the sections of
  910. 40:08this memory is calculated. So, let's say
  911. 40:10when you submit your Spark job, right?
  912. 40:13So, you say exec {hyphen} {hyphen}
  913. 40:15executor
  914. 40:16{hyphen} memory
  915. 40:18is let's say 10 GB.
  916. 40:21Right? Let's say 10 GB. So, this whole
  917. 40:24section for the on heap memory from here
  918. 40:26to here is defined by
  919. 40:29spark.executor.memory,
  920. 40:31right? So, this is going to be 10 GB.
  921. 40:35What we submitted for
  922. 40:37the executor memory over here. This is
  923. 40:39the on heap memory. Now, this portion
  924. 40:41this portion of the unified memory from
  925. 40:43here to here execution and storage
  926. 40:46memory, it is defined by
  927. 40:48spark.memory.fraction,
  928. 40:50and this is 0.6 of your
  929. 40:54spark.executor.memory,
  930. 40:55right? So, this is going to be 0.6 into
  931. 40:5710 GB, which is 6 GB, and this is going
  932. 41:01to be equally divided between execution
  933. 41:04and storage. And remember that this is
  934. 41:08movable, both of them can borrow memory
  935. 41:10from each other. Yeah? Reserved memory,
  936. 41:12as we've already seen, is 300 MB.
  937. 41:15Therefore, this portion of the memory is
  938. 41:18going to be 0.4 because this was 0.6,
  939. 41:21right? 0.4 into
  940. 41:25into 10 GB minus 300 MB,
  941. 41:29right?
  942. 41:30So, this is going to be 0.4
  943. 41:33into 10 into 1024
  944. 41:36minus 300 MB,
  945. 41:39which is simply going to give me 37
  946. 41:4396 MB,
  947. 41:44which is 3.
  948. 41:46close to 3.8 GB.
  949. 41:49Right? So, this is how different
  950. 41:50sections of the memory will be
  951. 41:52calculated. So, that's all about the 30
  952. 41:55most important Apache Spark concepts. If
  953. 41:57you found this video insightful, I
  954. 42:00believe you'll also love the Spark
  955. 42:02performance tuning playlist and the
  956. 42:046-hour long Delta Lake master class on
  957. 42:07my YouTube channel. If you're preparing
  958. 42:09for interviews or upskilling in data
  959. 42:12engineering, these videos might be very
  960. 42:14helpful for you. Thank you for watching
  961. 42:17and I will see you in the next video.

About this transcript

This page contains the full transcript of Apache Spark Was Hard Until I Learned These 30 Concepts! by Afaque Ahmad, generated from the public captions YouTube serves with the video. The transcript has 6,045 words across 961 segments, with the original timestamps preserved so you can click any line to jump to that moment in the embedded player.

What you can do with it

Use the transcript to take notes, quote the speaker, build a study guide, generate a summary with ChatGPT or Claude via the YouTube Summary tool, or export it as a timed subtitle file with YouTube to SRT. You can also re-open it in the transcriber to translate the transcript into 100+ languages.

Free YouTube transcript tool

YouTube2Text is a free YouTube transcript generator — no signup, no daily limit. Paste any YouTube link and get the full transcript instantly, with timestamps, click-to-jump, translation to 100+ languages, AI prompts for ChatGPT, Claude, and Gemini, and exports to TXT, SRT, VTT, or Markdown.