YouTube2Text

Mastering Joins In Apache Spark: Complete Deep Dive — Transcript

by Afaque Ahmad · 12,154 words · 1,852 segments · language en · Watch on YouTube

Full transcript

  1. 0:00Joins are one of the most commonly and
  2. 0:03widely used operations in Apache Spark.
  3. 0:06In most of the code that you write on a
  4. 0:09day-to-day basis in Apache Spark, you
  5. 0:12must have used joins, right? And joins
  6. 0:15is actually where most of the problems
  7. 0:18in Apache Spark lives, right? So, in
  8. 0:20this video, we are going to cover four
  9. 0:24very important Spark physical joins in
  10. 0:27depth, right? So, we are going to cover
  11. 0:31We are going to cover sort merge join.
  12. 0:35We are going to cover broadcast hash
  13. 0:36join, shuffle hash join, and broadcast
  14. 0:41nested loop join. Of course, there are
  15. 0:42other joins, but understanding these
  16. 0:44four core joins is going to help you
  17. 0:47develop an intuition for the rest of the
  18. 0:50joins, right? So, for each of these
  19. 0:53joins, right? So, for each of these four
  20. 0:57joins, we are going to do three things.
  21. 1:00Right? The first one is we are going to
  22. 1:03understand what are the exact conditions
  23. 1:06that make Spark choose that join
  24. 1:09strategy, right? So, Spark has sort
  25. 1:12merge join available, Spark has
  26. 1:14broadcast hash join available, and then
  27. 1:16it also has these other two types of
  28. 1:18join available, right?
  29. 1:20What exactly are the conditions that
  30. 1:23make Spark choose, let's say, a sort
  31. 1:26merge join strategy over a broadcast
  32. 1:29hash join? So, that is the first thing
  33. 1:31that we are going to understand, the
  34. 1:33exact conditions that make Spark choose
  35. 1:36a particular join strategy, right?
  36. 1:39Number two, we are going to walk through
  37. 1:42step-by-step how each of these join
  38. 1:45strategy work under the hood, right? So,
  39. 1:47we are going to walk through visuals and
  40. 1:49understand each of the steps on how the
  41. 1:53join strategy work under the hood. Yeah?
  42. 1:56And finally, we are going to prove it
  43. 1:57with real code by looking at the Spark
  44. 2:00physical plan, the Spark physical UI,
  45. 2:03right? On how each of these joins are
  46. 2:08going to look like, right? So, how does
  47. 2:10the Spark physical plan look like for a
  48. 2:13sort merge join? How does it look like
  49. 2:16on the Spark UI, right? You're going to
  50. 2:18have a lot more clarity after going
  51. 2:22through all of these three steps. So,
  52. 2:24now let's get started with the sort
  53. 2:26merge join. Now, as data and AI
  54. 2:28engineers, we obsess a lot over how data
  55. 2:32moves, which is through joins,
  56. 2:34partitions, and shuffles. But, getting
  57. 2:37clean and structured data from the web
  58. 2:40in the first place, that is still one of
  59. 2:42the messiest parts of the job. SerpApi
  60. 2:45takes care of exactly that. This video
  61. 2:49is sponsored by SerpApi. SerpApi gives
  62. 2:52you clean, structured JSON results from
  63. 2:55Google, YouTube, Google Scholar, and
  64. 2:58many more, all through a single API
  65. 3:02call. So, they take care of all of the
  66. 3:04infrastructure headaches, CAPTCHA
  67. 3:06solving, proxy rotation, layout changes,
  68. 3:09none of that is your problem. So, you
  69. 3:12get reliable, structured data ready to
  70. 3:14plug straight into your pipeline. For
  71. 3:17people building in this space, two
  72. 3:19things stand out for me. Firstly, if you
  73. 3:21are training ML or AI models, SerpApi
  74. 3:25lets you build custom data sets using
  75. 3:28real-time search data. So, things like
  76. 3:30image titles, thumbnails, URLs,
  77. 3:33peer-reviewed articles from Google
  78. 3:35Scholar, all of them structured and
  79. 3:38ready to go. Secondly, if you are
  80. 3:41working on portfolio projects, this is
  81. 3:43genuinely a great way to power your apps
  82. 3:46with live, production-like data from
  83. 3:49some of the world's top search engines
  84. 3:52in JSON format that plays out really
  85. 3:54well with Spark and Databricks. They run
  86. 3:57at 99.9%
  87. 3:59uptime with a 1.2 second average
  88. 4:03response time. So, it's really
  89. 4:05production grade. Get started with 250
  90. 4:07free credits by clicking the link in the
  91. 4:09description or scanning the QR code
  92. 4:12here.
  93. 4:13All right. Now, let's get back to sort
  94. 4:15merge join. Joins are one of the most
  95. 4:19commonly and widely used operations in
  96. 4:21Apache Spark. In most of the code that
  97. 4:24you write on a day-to-day basis in
  98. 4:26Apache Spark, you must have used joins,
  99. 4:30right? And joins is actually where most
  100. 4:33of the problems in Apache Spark lives,
  101. 4:35right? So, in this video, we are going
  102. 4:38to cover four very important Spark
  103. 4:42physical joins in depth, right? So, we
  104. 4:45are going to cover
  105. 4:47We are going to cover sort merge join.
  106. 4:51We are going to cover broadcast hash
  107. 4:52join, shuffle hash join, and broadcast
  108. 4:56nested loop join. Of course, there are
  109. 4:58other joins, but understanding these
  110. 5:00four core joins is going to help you
  111. 5:03develop an intuition for the rest of the
  112. 5:06joins, right? So, for each of these
  113. 5:09join, right? So, for each of these four
  114. 5:12joins, we are going to do three things,
  115. 5:16right? The first one is we are going to
  116. 5:18understand what are the exact conditions
  117. 5:22that make Spark choose that join
  118. 5:25strategy, right? So, Spark has sort
  119. 5:28merge join available. Spark has
  120. 5:30broadcast hash join available. And then,
  121. 5:32it also has these other two types of
  122. 5:34join available, right?
  123. 5:36What exactly are the conditions that
  124. 5:39make Spark choose, let's say, a sort
  125. 5:42merge join strategy over a broadcast
  126. 5:45hash join? So, that is the first thing
  127. 5:47that we are going to understand. The
  128. 5:49exact conditions that make Spark choose
  129. 5:52a particular join strategy, right?
  130. 5:55Number two, we are going to walk through
  131. 5:58step by step how each of the join
  132. 6:01strategy work under the hood, right? So,
  133. 6:03we are going to walk through visuals and
  134. 6:05understand each of the steps on how the
  135. 6:08join strategy work under the hood, yeah?
  136. 6:11And finally, we are going to prove it
  137. 6:13with real code by looking at the Spark
  138. 6:16physical plan, the Spark physical UI,
  139. 6:19right? On how each of the joins are
  140. 6:23going to look like, right? So, how does
  141. 6:26the Spark physical plan look like for a
  142. 6:28sort merge join? How does it look like
  143. 6:31on the Spark UI, right? You're going to
  144. 6:34have a lot more clarity after going
  145. 6:37through all of these three steps. So,
  146. 6:39now let's get started with the sort
  147. 6:41merge join. Okay, so the first type of
  148. 6:43join that we are going to study is a
  149. 6:45sort merge join, right? So, first of
  150. 6:49all, let's understand when is it exactly
  151. 6:52performed? When is it When is it exactly
  152. 6:54chosen by Spark over any other join
  153. 6:58strategy. So, most importantly, there
  154. 7:00are
  155. 7:01two conditions, right? The first one is
  156. 7:04when both the data frames are large,
  157. 7:07right? And what this means is that when
  158. 7:10both the sides exceed a property, right?
  159. 7:14So, there is a property which is
  160. 7:20spark.sql.autobroadcastjointhreshold.
  161. 7:24And this
  162. 7:25the default
  163. 7:27value is set to 10 MB. It is sent set to
  164. 7:3110 MB, right? So, when both of the sides
  165. 7:34exceed this value, neither of them can
  166. 7:37be broadcasted, right? So, Spark is
  167. 7:39going to fall back to sort merge join.
  168. 7:45Right? And it is the only join that
  169. 7:47handles unlimited size on both sides
  170. 7:51without loading either of the data sets
  171. 7:53fully into memory, right? So, the first
  172. 7:55condition is both of the data frame
  173. 7:58should be large. And what this means is
  174. 8:01when both of the sides is greater than
  175. 8:04this auto broadcast join threshold
  176. 8:06property, which is
  177. 8:0710 megabytes, and as a result of which
  178. 8:10because
  179. 8:11they are higher than this threshold,
  180. 8:14none of them can be broadcasted, and
  181. 8:15therefore Spark falls back to this sort
  182. 8:19merge join strategy. The second one is
  183. 8:22equi-join on sortable keys, right? So,
  184. 8:26the key the keyword over here is
  185. 8:28equi-join.
  186. 8:30What equi-join mean over here is when
  187. 8:32you have your join condition, which is
  188. 8:34let's say
  189. 8:37DF1.x
  190. 8:39= DF2.x,
  191. 8:41right? So, this operator over here this
  192. 8:44operator over here needs to be an equal
  193. 8:47to.
  194. 8:48Right? It cannot be greater than,
  195. 8:50greater than equal to, less than, less
  196. 8:52than equal to, not equal to.
  197. 8:54Right? For all of such conditions, sort
  198. 8:58merge join will not work, right? And the
  199. 9:00keys needs to be sortable so that they
  200. 9:02can be put in order. And as you will see
  201. 9:04that sorting is one of the steps of sort
  202. 9:08merge join, right? So,
  203. 9:10the keys need to be sortable, yeah? So,
  204. 9:13these are the two most important
  205. 9:14conditions that need to be satisfied in
  206. 9:18order for a sort merge join to be chosen
  207. 9:21by Apache Spark, right? So, the three
  208. 9:24steps the three steps of sort merge
  209. 9:27join, and we're going to see it in a
  210. 9:28visualization,
  211. 9:30uh but to keep it concrete, the the
  212. 9:32three steps is shuffle,
  213. 9:35sort, and merge. By shuffling, what we
  214. 9:38mean is that
  215. 9:40all of the rows with the same joint key
  216. 9:42end up together in the same machine and
  217. 9:45in the same partition, right? So, the
  218. 9:48same keys after a particular hash after
  219. 9:51doing a particular hash hash of the key
  220. 9:55and this is the
  221. 9:57joint key.
  222. 10:00Right? Whatever value they get, right?
  223. 10:03Let's say this is X.
  224. 10:05And some other keys get Y. All of the
  225. 10:07keys which return X, they are going to
  226. 10:10end up in the same partition. Let's say
  227. 10:12it's called partition zero. Y is going
  228. 10:14to end up in partition one and so on,
  229. 10:16right? And then there is the sort step
  230. 10:19and finally a merge step. So, we are
  231. 10:20going to see all of this in action now.
  232. 10:23Okay, so let's understand it with this
  233. 10:25visualization, yeah?
  234. 10:28So, here we have two executors first of
  235. 10:31all, right?
  236. 10:33And in each of these two executors, we
  237. 10:35have two data sets. As you see, we have
  238. 10:38a customer data set and we have an order
  239. 10:42data set. In the first executor, we have
  240. 10:44two partitions for the customer data
  241. 10:46set, which is partition zero and
  242. 10:49partition number one.
  243. 10:51And we have one partition
  244. 10:55for the customer data set in the second
  245. 10:57executor and two partitions for the
  246. 11:01orders data set
  247. 11:03in the second executor, right? So, this
  248. 11:05is how the data set is distributed. You
  249. 11:07have a customer data set and an order
  250. 11:10data set. And now we want to do a join
  251. 11:14between the two. So, the join key is
  252. 11:17customer ID, right?
  253. 11:19The joint key is
  254. 11:25The joint key is customer ID, right?
  255. 11:29So, now once we figured out that the
  256. 11:31joint key is customer ID, the first step
  257. 11:34is going to be the shuffle step, right?
  258. 11:37The first step is going to be the
  259. 11:39shuffle step.
  260. 11:41And as we've understood that in a
  261. 11:43shuffle, the same keys are going to land
  262. 11:47in the same partition, right? Now, a
  263. 11:50very important parameter to define is
  264. 11:52the number of shuffle partition.
  265. 11:56Right? Is the number of shuffle
  266. 11:58partition. Shuffle partition is
  267. 11:59basically the number of partitions that
  268. 12:01will be produced for one single data set
  269. 12:04in the shuffle step, right? So, let's
  270. 12:07say we are assuming this to be two,
  271. 12:10right? Now, in order for the shuffle
  272. 12:11step to proceed and the same keys to
  273. 12:13land in the same partition, there is a
  274. 12:16logic that is followed. And that logic
  275. 12:19is basically
  276. 12:21that logic is basically actually let me
  277. 12:23erase
  278. 12:24some of the stuff here.
  279. 12:27So, that logic is basically hash
  280. 12:30of
  281. 12:31the join key
  282. 12:33mod shuffle partition, right? So,
  283. 12:36shuffle partition is two. So, this is
  284. 12:39always going to produce a number between
  285. 12:41zero and one, right? So, whatever the
  286. 12:44join key is, it is always going to
  287. 12:46produce a number between zero and one
  288. 12:49because shuffle partition is two, right?
  289. 12:52So,
  290. 12:53this means that the key is going to go
  291. 12:55to either partition number zero or it is
  292. 12:58going to go to partition number one.
  293. 13:00Okay, so, let's take an example.
  294. 13:02So, we are going to calculate hash of
  295. 13:04the join key, which is your customer
  296. 13:06one,
  297. 13:07mod two, right? So, let's assume for
  298. 13:11simplicity that hash of CX,
  299. 13:14where X is a number, right? These are
  300. 13:16all numbers, C1, C3, C6, and C4.
  301. 13:21So, we are going to assume hash of CX
  302. 13:23mod shuffle partitions, this is simply
  303. 13:27going to be X mod shuffle partition,
  304. 13:30right? So, in this case, what this means
  305. 13:32is we are going to take this X one mod
  306. 13:35two, and this is one. And again, this is
  307. 13:37just for simplicity. In reality, as I
  308. 13:40mentioned, this hash is going to be
  309. 13:43MurmurHash or some other kind of hash
  310. 13:46that is widely used. Yeah? So, now we
  311. 13:49are going to do a one mod two, which is
  312. 13:52one. So, this is going to go to
  313. 13:54partition number one, right? Similarly,
  314. 13:58C3 is going to be three mod two, and
  315. 14:01this is also going to go to partition
  316. 14:02one. C6 is going to be six mod two,
  317. 14:06which is zero. So, this is going to
  318. 14:08partition number zero, partition number
  319. 14:10zero,
  320. 14:11partition number zero, partition number
  321. 14:13one, partition number zero. And
  322. 14:15similarly, we are going to repeat the
  323. 14:17same for this one. So, this will be
  324. 14:18partition one, partition zero,
  325. 14:21partition one, partition zero, partition
  326. 14:24zero,
  327. 14:25partition one, and partition one. Right?
  328. 14:29So, this logic is going to decide in
  329. 14:32which partition are each of these rows
  330. 14:35going to go to, right? So, the ones that
  331. 14:37end with modulo one, they are going to
  332. 14:39go to partition number one. And the ones
  333. 14:43that give you a result of zero after the
  334. 14:45modulo, they are going to partition
  335. 14:47number zero. Right?
  336. 14:50So, in the next step,
  337. 14:52as we said earlier, the shuffle
  338. 14:55partition the shuffle partitions is two
  339. 14:57in number, and basically the shuffle
  340. 15:00partition means the output partition per
  341. 15:02data set after or during the shuffle
  342. 15:06phase, right? So, here you see for
  343. 15:09executor number one, for customer, this
  344. 15:12is the first partition, or rather, I
  345. 15:14should just highlight this partition
  346. 15:16number zero, and partition number one.
  347. 15:19This should be partition number one.
  348. 15:21Right? And all the ones all the customer
  349. 15:24ID
  350. 15:25with the modulo of 06 mod
  351. 15:282 is equal to 0, they have all ended up
  352. 15:32in partition number 0. Right?
  353. 15:35This is all
  354. 15:37have ended up in partition number 0 for
  355. 15:39the orders partition. Similarly, this
  356. 15:42has ended up in partition number one
  357. 15:44because 3 mod 2 is equal to 1. And
  358. 15:48similarly, all of this has also up in
  359. 15:51partition number one. Right? So, this
  360. 15:53complete the shuffle phase in which the
  361. 15:56same key the same key as in this the
  362. 15:59keys which have the same modulo value
  363. 16:02after doing the hash, right? They end up
  364. 16:05in the same partition. So, that
  365. 16:07completes the shuffle phase. Right? So,
  366. 16:11now once the shuffle phase is completed,
  367. 16:13we then go ahead to the sort phase.
  368. 16:16Right? And in the sort phase, it is
  369. 16:18simply going to take up this customer ID
  370. 16:20right over here and it is going to sort
  371. 16:24the row by that customer ID. So, let's
  372. 16:26go ahead to the next step over here and
  373. 16:29we see that this is all sorted which is
  374. 16:31C2, C4, C6, right? And this is sorted.
  375. 16:35C2, C4, and C6. And wherever there is a
  376. 16:39tie, it has taken up order one and order
  377. 16:42number seven. And similarly, C1, C3, C5
  378. 16:46and C1, C3, C5. Right? So, this then
  379. 16:50complete the sort phase over here.
  380. 16:54And finally, we are then going to go
  381. 16:57ahead with the merge
  382. 17:01step.
  383. 17:02Right?
  384. 17:03And the merge step is very similar to
  385. 17:06the merge step in merge sort.
  386. 17:10Right? In merge sort. So, actually, let
  387. 17:12me remove
  388. 17:14all of this quickly
  389. 17:16and show you how the merge is going to
  390. 17:18look like. So,
  391. 17:19both of them for both of them a pointer
  392. 17:21is going to be placed, right? And a
  393. 17:24comparison is going to happen whether
  394. 17:26this C2 is the same as this C2. And in
  395. 17:29this case, it is. Therefore, a join is
  396. 17:31going to happen and the resulting rows
  397. 17:33you're going to see over here. Right?
  398. 17:35These are going to be your joined rows.
  399. 17:39Yeah?
  400. 17:41Then it goes to the next row. It
  401. 17:42compares, is this value and this value
  402. 17:45the same? Which it is.
  403. 17:47So, a join is going to happen over here
  404. 17:49between this row and this row. Now it
  405. 17:52comes over here and it compares the if
  406. 17:54this value C4 same as this one, which is
  407. 17:57not the case, then it is going to move
  408. 17:59the pointer ahead. It is going to come
  409. 18:02to C4 and then again a comparison happen
  410. 18:04between this value C4 and this value.
  411. 18:07And now the values are the same. So,
  412. 18:10therefore this there is going to be a
  413. 18:12join between this row
  414. 18:13and this row. And then the value the
  415. 18:16pointer comes to the next row and it
  416. 18:18compares whether this C6 and this value
  417. 18:21over here is the same. It is not. So,
  418. 18:23then the pointer move over here and
  419. 18:25again the comparison happen. This time
  420. 18:27both of them are the same. So, then you
  421. 18:29see the join result over here, which is
  422. 18:31Adam and Mumbai, right? Similar thing is
  423. 18:35going to be repeated over here for this
  424. 18:37one and this one and the matching is
  425. 18:39going to happen, right? So, a join
  426. 18:41between these two data set is going to
  427. 18:44produce a new joined
  428. 18:48partition.
  429. 18:51Right? So, this is going to be the first
  430. 18:53partition that is going to be produced
  431. 18:55and both of them is going to be joined
  432. 18:57and this is also going to produce a new
  433. 19:00joined partition.
  434. 19:04And as a result, this is what is the
  435. 19:07output that you see.
  436. 19:09Right? So, on executor number one, this
  437. 19:11is your first output, which is basically
  438. 19:15this one right over here
  439. 19:17on executor number one. And for executor
  440. 19:20number two, this is going to produce the
  441. 19:22second partition, which is right over
  442. 19:25here, partition number one and partition
  443. 19:27number zero. Right? And this is the
  444. 19:29joined result for both of the data set.
  445. 19:33Right? So, this is how a sort merge join
  446. 19:37works. Okay, so now let's go ahead with
  447. 19:39the lab for the sort merge join. Right?
  448. 19:44And I'm going to use two data sets. So,
  449. 19:47I have already put
  450. 19:49uh
  451. 19:50in this default schema. So, for me,
  452. 19:54um when you create a workspace,
  453. 19:57it creates a catalog by default. And
  454. 20:00this is the default catalog for me.
  455. 20:02Right? And similarly, when you create a
  456. 20:04workspace, you are also going to have a
  457. 20:07default catalog. And inside of this,
  458. 20:09there is going to be a default schema.
  459. 20:13Right? So, I have created a volume where
  460. 20:18I can store all of the data. Right? A
  461. 20:22CSV file. Right? So, I created this
  462. 20:24volume called ecom, and you can also
  463. 20:26create it by going over here and
  464. 20:28clicking on volume. And this can be a
  465. 20:31managed volume. Right? So, I when I
  466. 20:34click on ecom, I see that there are two
  467. 20:37files under the folder ecom. Right? So,
  468. 20:41I created a directory which is ecom, and
  469. 20:44then I put these two CSVs, which we are
  470. 20:48going to use in all of the exercises.
  471. 20:50So, don't worry, all of this will be
  472. 20:52made available in GitHub, from where you
  473. 20:55can download and follow the exercises.
  474. 20:58Right? So, now let's get started. So,
  475. 21:01here I'm going to show you the physical
  476. 21:03plan and how a sort merge join will look
  477. 21:08like on the Spark UI, right? So, let's
  478. 21:11set up a few things first. Uh we are
  479. 21:14going to import PySpark uh dot SQL uh
  480. 21:17this functions module as F.
  481. 21:20And what we are doing over here
  482. 21:22essentially is I am setting the auto
  483. 21:25broadcast join threshold to minus one,
  484. 21:28right? So, because the two data sets
  485. 21:30that I have uploaded, they are small in
  486. 21:33size, so naturally Spark is going to
  487. 21:35force a broadcast hash join, right? But
  488. 21:39we want to see an example of a sort
  489. 21:42merge join, right? So, for that reason,
  490. 21:44I am setting this to minus one
  491. 21:47and the number of shuffle partitions to
  492. 21:50four, yeah. So, let's go ahead and run
  493. 21:52this. Another interesting property is
  494. 21:55the spark.sql.adaptive.enabled,
  495. 21:58which is basically turning on AQE, which
  496. 22:02is adaptive query execution, right? So,
  497. 22:05if we set this on, you are going to see
  498. 22:08a lot of things in the Spark UI like AQE
  499. 22:10shuffle read and other optimizations,
  500. 22:12right? So, we are turning this off just
  501. 22:16to have a clean and aesthetic plan when
  502. 22:19we look at the Spark UI and the physical
  503. 22:22plan, right? So, let's go ahead and run
  504. 22:25this.
  505. 22:26And
  506. 22:27uh we are setting the paths over here,
  507. 22:29which is basically uh the path that I
  508. 22:32have got from here. If you click on copy
  509. 22:34path, this is going to give you the same
  510. 22:36path, right? For customer.
  511. 22:39So, let's go ahead and run this. And
  512. 22:41now, we are going to read both of the
  513. 22:43data frame with the header set to true
  514. 22:46and the infer schema set to true, so
  515. 22:49that it is automatically able to
  516. 22:51understand the data types,
  517. 22:53right? Okay, so the two data frames have
  518. 22:55been read and here you see the columns
  519. 23:00and the data type. The customer ID is an
  520. 23:02integer. All the others, which is first
  521. 23:04name, last name, email, phone, city,
  522. 23:07these are all strings. The sign-up date
  523. 23:09is date, and whether the customer is
  524. 23:12active or not is a boolean.
  525. 23:14And similarly for orders, we have all
  526. 23:16the relevant
  527. 23:19columns and the data type, which is
  528. 23:21order ID and customer ID is integer. The
  529. 23:24order date is a date. Category, product,
  530. 23:27and all of this is string. Unit
  531. 23:30uh quantity is integer and unit price is
  532. 23:32double. The amount is double, and so on,
  533. 23:36right? So, let's also quickly have a
  534. 23:38look at both the data sets. And this is
  535. 23:40how it looks like. Customer ID, first
  536. 23:41name, email,
  537. 23:43phone number, and all of this, right?
  538. 23:45And similarly, the orders is where the
  539. 23:47order ID is over here. The customer who
  540. 23:50placed the order, when did they place an
  541. 23:52order, what exactly did they place an
  542. 23:54order, and the total prices and the
  543. 23:56amount, right?
  544. 23:58Now, let's go ahead. Let me skip this
  545. 24:01for a while. Well, actually, let's go
  546. 24:03ahead and run this. Let's see how many
  547. 24:06rows are there in each of the data sets.
  548. 24:08So, there are 50,000 rows in each of the
  549. 24:11data sets. Now, let's go ahead and do a
  550. 24:14join. So, I'm joining orders with the
  551. 24:17customer data set. We are joining it on
  552. 24:19the customer ID because this is what is
  553. 24:22common between the two data sets, right?
  554. 24:25And now, we are going to do an inner
  555. 24:28join. Yeah? So, let's go ahead and run
  556. 24:30this.
  557. 24:32Yeah, let's go ahead and do a show.
  558. 24:40So, this is how my data set is going to
  559. 24:43look like. So, you have all the
  560. 24:46order level details over here, right?
  561. 24:48The order date, and
  562. 24:53the customer level details over here,
  563. 24:54which is the first name, last name,
  564. 24:56email, and all of this, right? Now,
  565. 24:58before we go to the Spark UI, let's
  566. 25:00first have a look at the plan, right?
  567. 25:03So, this is the physical plan. So, let
  568. 25:05me copy the physical plan, and let's
  569. 25:08paste it over here.
  570. 25:10Yeah?
  571. 25:11So, let's have a look at this line by
  572. 25:14line, right? So, the way we read the
  573. 25:16physical plan is from the bottom to the
  574. 25:19top, right?
  575. 25:21So, the first line is basically a file
  576. 25:23scan CSV.
  577. 25:25Right? And this is your customer ID,
  578. 25:27first name, last name, all of this. So,
  579. 25:28this simply means that we are or Spark
  580. 25:32is reading the customer's data set,
  581. 25:35right? And let's have a look at a few
  582. 25:37other things, right? An interesting
  583. 25:40thing to note over here is this
  584. 25:41in-memory file index, and this is
  585. 25:44basically the path that we have
  586. 25:46specified, right? So, in-memory file
  587. 25:50index is basically Spark file listing
  588. 25:53component, which finds out the file that
  589. 25:56is present at a given location, right?
  590. 25:59So, it basically finds out the metadata,
  591. 26:01sizes, partition values, and all of it,
  592. 26:04right? So, it is going to list all of
  593. 26:06the files specified at that location,
  594. 26:09and it stores it in the driver's memory.
  595. 26:13So, here we have specified one single
  596. 26:16file. If we specified a directory, the
  597. 26:19number is going to be a lot larger,
  598. 26:21right? So, that is one. And this is the
  599. 26:25schema of the data set that has been
  600. 26:27read in.
  601. 26:29Now, if we go to the second line,
  602. 26:31you see that filter not null,
  603. 26:34and this is basically filtering on the
  604. 26:36customer ID. So, it is basically making
  605. 26:39sure that the customer IDs that have
  606. 26:41been read in
  607. 26:42from the data set is not null. And a
  608. 26:45very important point to remember is that
  609. 26:47the customer ID is your join key. Now,
  610. 26:51we did not add this condition anywhere,
  611. 26:53right? Did we add this condition
  612. 26:55anywhere? No.
  613. 26:56Right? We haven't added this condition
  614. 26:58anywhere. So, this is a part of Spark
  615. 27:01internal optimization strategy. Yeah?
  616. 27:05To make sure that the data is clean.
  617. 27:07The third step is exchange hash
  618. 27:10partitioning on the customer ID. And the
  619. 27:13number four that you see over here,
  620. 27:15yeah? The number four that you see over
  621. 27:17here is basically the shuffle partition.
  622. 27:20So, we certify the shuffle partition as
  623. 27:23four over here, right? So, basically
  624. 27:25before going ahead and doing the
  625. 27:27shuffle, exchange over here means
  626. 27:29shuffle.
  627. 27:31It basically
  628. 27:33makes sure that it has specified that
  629. 27:36the number of shuffle partitions are
  630. 27:38four. And what does hash partitioning
  631. 27:40mean? So, if you remember, we we learned
  632. 27:43that during a shuffle, the same keys go
  633. 27:46to the same partition. And how does that
  634. 27:48happen? It happens by taking up the join
  635. 27:51key, doing a hash, and then a modulo of
  636. 27:56the part of the shuffle partition,
  637. 27:58right? So, basically the hash of the
  638. 28:01join key,
  639. 28:03and then modulo the shuffle partition.
  640. 28:07Right? And in this case, the shuffle
  641. 28:09partition is four. So, this is simply
  642. 28:12going to be
  643. 28:13this one
  644. 28:14mod four, right?
  645. 28:17And this is what is hash partitioning
  646. 28:19over here. Right? So, this step
  647. 28:21basically distributes all of the data,
  648. 28:24making sure that the same keys fall in
  649. 28:27the same partition, right? So, this is
  650. 28:30the shuffle step of the sort merge join.
  651. 28:33And then, as we've seen,
  652. 28:35each of the rows is sorted within the
  653. 28:38partition, right? And it is sorted by
  654. 28:41the customer ID. And the logic that is
  655. 28:43used is ascending and nulls first, yeah?
  656. 28:46So, the sort is complete, the shuffle is
  657. 28:50complete, and then we repeat the same
  658. 28:53for the orders data set. Right? So, here
  659. 28:57you see a file scan, and then the order
  660. 29:00customer ID, order date, and all of
  661. 29:02that. And here you would see a similar
  662. 29:05in-memory file index, which means that
  663. 29:08it is reading Basically, it is not
  664. 29:11reading actually, it is listing the
  665. 29:13files that have been specified in the
  666. 29:16path that we have given. Right? And here
  667. 29:19it has found only one file. Yeah? And it
  668. 29:22is doing something very similar, which
  669. 29:24is filter not null, making sure that the
  670. 29:26customer ID is clean. And it is then
  671. 29:29doing a shuffle.
  672. 29:32Yeah? So, exchange specifies a shuffle,
  673. 29:35and then it does a hash partitioning.
  674. 29:37So, hash partitioning is the logic that
  675. 29:39it uses in order to do the shuffle. And
  676. 29:42you see four shuffle partitions over
  677. 29:44here. Yeah? And then the sort step
  678. 29:47again. And finally, when both of the
  679. 29:49data frames, right? When both of the
  680. 29:51data frames have been shuffled, they
  681. 29:53have been sorted, the next step is the
  682. 29:55merge step. Right? And this is where
  683. 29:58your final sort merge join happens, and
  684. 30:01it has specified that this is an inner
  685. 30:05join. Yeah? So, once the join has
  686. 30:07happened, project basically picks up the
  687. 30:09relevant columns that the user needs to
  688. 30:13see, or the ones that need to be
  689. 30:15displayed. Right? So, this is how the
  690. 30:17physical plan for a sort merge join is
  691. 30:20going to look like. Now, let's have a
  692. 30:22look at how the Spark UI is going to
  693. 30:24look like. So, here is where
  694. 30:27the join happened. Right?
  695. 30:31So, this is where the action was called,
  696. 30:33and that is why it triggered a Spark
  697. 30:35job. Now, let's have a look at
  698. 30:39the Spark UI, and the job number is
  699. 30:42eight. So, let me directly go to the SQL
  700. 30:47{slash} data frame and here you see that
  701. 30:49this is job number eight. Let's go ahead
  702. 30:51over here and you are going to see a
  703. 30:54DAG, right? So, this DAG basically scans
  704. 30:59the scans the orders data set and you
  705. 31:02would see that this basically reads up
  706. 31:0550,000 rows, yeah?
  707. 31:07It reads up 50,000 rows and similarly
  708. 31:10this is the customer data set, yeah? So,
  709. 31:13you see customer ID, first name,
  710. 31:15uh last name, email and all of that. And
  711. 31:19the rows output is 10,000,
  712. 31:23yeah?
  713. 31:24Now,
  714. 31:25what we are going to do is
  715. 31:28it is going to
  716. 31:30do a filter of not null on both the data
  717. 31:34sets, right? Something similar to what
  718. 31:36we just saw in the Spark physical plan.
  719. 31:39And then it is going to do an exchange
  720. 31:41hash partitioning
  721. 31:43of
  722. 31:44both the data sets, right? On the left
  723. 31:47and on the right,
  724. 31:49yeah?
  725. 31:50On the left and on the right and then it
  726. 31:52does a sort of both the data frames and
  727. 31:55finally it does a merge, right? Which is
  728. 31:58specified as a sort merge join over here
  729. 32:01and the key that is used, the join key
  730. 32:04is basically the customer ID and it
  731. 32:06specifies an inner join and finally
  732. 32:08there are 10 columns which is selected
  733. 32:10as a result of the final data set being
  734. 32:13produced. So, this is how overall
  735. 32:16all sort merge joins are going to look
  736. 32:18like, right? All sort merge join, the
  737. 32:20DAG for all the sort merge join is going
  738. 32:23to look something very similar like
  739. 32:26this. Okay, so the next type of join we
  740. 32:28are going to study is broadcast hash
  741. 32:31join, right? And let's first understand
  742. 32:34when exactly is this chosen by Spark,
  743. 32:38right? So, broadcast hash join is chosen
  744. 32:41or it is performed when, let's say, you
  745. 32:43have two tables and one of the tables in
  746. 32:47the join is small enough to fit entirely
  747. 32:50in memory in each of the executor,
  748. 32:52right? So, the key point to note is one
  749. 32:54of the tables or one of the data frame
  750. 32:56is small enough to fit entirely in the
  751. 33:01executor's memory on each of the
  752. 33:03executor's memory, right? So, the Spark
  753. 33:06based Spark cost based optimizer, right?
  754. 33:09So, the Spark
  755. 33:12CBO, which is cost based optimizer,
  756. 33:15automatically chooses broadcast hash
  757. 33:17join when the estimated size of one of
  758. 33:21the table is below the threshold defined
  759. 33:24by this parameter.
  760. 33:26And we've seen this parameter earlier,
  761. 33:27which is
  762. 33:27spark.sql.autoBroadcastJoinThreshold.
  763. 33:36The default value of which is 10 MB,
  764. 33:39right? So, the Spark cost based
  765. 33:41optimizer automatically is going to
  766. 33:43automatically choose this when the
  767. 33:46estimated size of one of the table is
  768. 33:49below this threshold, right? So, that is
  769. 33:53the case when
  770. 33:55broadcast hash join is going to be
  771. 33:57chosen, right? So, to quickly summarize,
  772. 33:59you have two tables, one table is large,
  773. 34:01the other table is small, and the other
  774. 34:04table, the size of the other table is
  775. 34:06below the threshold, which is defined by
  776. 34:10this property over here. In this case,
  777. 34:13Spark is going to choose broadcast hash
  778. 34:16join, yeah? So, there are three steps to
  779. 34:19performing a broadcast hash join. The
  780. 34:21first one is, of course, broadcast,
  781. 34:24right? And we're going to understand
  782. 34:26what exactly each of these steps are,
  783. 34:28but the three steps quickly are
  784. 34:30broadcast, build, and probe. Build is
  785. 34:33building the hash map, and and probe is
  786. 34:35actually finding for finding the right
  787. 34:37key, finding the match in a hash map,
  788. 34:41right? So, let's understand this
  789. 34:42visually with an example all of these
  790. 34:45three steps. Okay, so this is our setup.
  791. 34:48We have a customer table, and we have
  792. 34:51the orders table, right?
  793. 34:53There are two executors. Each of the
  794. 34:55executor have one order partition each,
  795. 34:59right? And right now
  796. 35:02the customer table resides on the
  797. 35:05driver. It is very small in size. As you
  798. 35:07see, we have listed that this is 320
  799. 35:10bytes in size, and the default value of
  800. 35:14auto broadcast join threshold is 10
  801. 35:17megabytes, right? So, what is going to
  802. 35:19happen is that the first step, which is
  803. 35:22the broadcast
  804. 35:25step,
  805. 35:27is going to take place, right? So, what
  806. 35:29is the broadcast step? The broadcast
  807. 35:31step is the step in which the driver is
  808. 35:35going to serialize the small table and
  809. 35:38push a full copy to every executor,
  810. 35:41right? So, here you see we have two
  811. 35:42executors. It is going to send two
  812. 35:45copies of this small customer table to
  813. 35:50these two executors over here, right?
  814. 35:53So, what the word serialize is mean
  815. 35:55which I I said that the driver is going
  816. 35:57to serialize the small table and then
  817. 35:59push a full copy. So, what that means is
  818. 36:02that this in-memory table, this table is
  819. 36:05the residing in memory, right?
  820. 36:06Everything resides in memory in Spark
  821. 36:08during computation. So, in-memory table
  822. 36:12is going to be converted into a flat
  823. 36:15sequence of bytes that can travel over
  824. 36:19the network, right? So, this is going to
  825. 36:22be serialized, and this table is going
  826. 36:24to be broadcasted
  827. 36:29on to executor number one
  828. 36:33on to executor number one and this is
  829. 36:35also going to be broadcasted on to
  830. 36:38executor number two. So, you will have
  831. 36:40one full copy
  832. 36:42You will have one full copy of the
  833. 36:44customer's
  834. 36:47partition on executor number two
  835. 36:50and on executor number one.
  836. 36:55Right? So, that is the first step which
  837. 36:57is the broadcast step. Now, coming over
  838. 37:00to the second step which is the build
  839. 37:03step. Right? The build step.
  840. 37:08Now, in the build step each executor is
  841. 37:11going to iterate over the received copy.
  842. 37:14So, it just now received the copy of the
  843. 37:16customer's partition. Yeah? And it is
  844. 37:19going to iterate over each of the rows
  845. 37:21and it is going to build a hash map and
  846. 37:24that is what you see over here.
  847. 37:27Right? So, this basically has four rows
  848. 37:30which is C1, C2, C3, C4 of Fark, Naved,
  849. 37:33Rohan, and Imdad.
  850. 37:36So, basically it is going to convert
  851. 37:39this into a hash map where
  852. 37:44it is going to be in the format of
  853. 37:46customer
  854. 37:49ID and the row value. Right? So, that is
  855. 37:53what you see over here. You have the
  856. 37:55customer IDs and the row value. Yeah?
  857. 38:01And this is going to take place at both
  858. 38:04of the executors.
  859. 38:06Right?
  860. 38:07So, that is the build phase and finally
  861. 38:10we come over to the probe phase.
  862. 38:15Right? We come over to the probe phase
  863. 38:17and in this step what happens is each
  864. 38:19executor is going to scan its partition
  865. 38:22of the large table row by row. So, the
  866. 38:25large table that you see over here, it
  867. 38:27is going to be scanned row by row and we
  868. 38:30are going to look if there is a match
  869. 38:33for C1 in this hash map.
  870. 38:38Right? So, it is basically going to
  871. 38:40check whether C1 exist in this hash map
  872. 38:43or not. Right? And this is simply an
  873. 38:47order of one operation, a very quick
  874. 38:49order of one operation. So, if there is
  875. 38:51a hit, it produces a joined row. So, C1
  876. 38:55for C1, it is a hit because you see this
  877. 38:59row over here. For C3,
  878. 39:01it is a hit and for C2, it is a hit.
  879. 39:04Right? And similarly for C4, you find C4
  880. 39:07over here, C1 over here and C3 over
  881. 39:11here. Right? So, that means all of these
  882. 39:14rows have found a hit and they are going
  883. 39:19to be joined.
  884. 39:21And the value is going to be produced
  885. 39:22over here at the joint data set. Right?
  886. 39:25So, for C1, you're going to have the
  887. 39:27name
  888. 39:29which is a park.
  889. 39:32And the city equal Bengaluru.
  890. 39:35Right? Similarly, this is going to
  891. 39:37happen for all. This is going to happen
  892. 39:38for all and you see the resulting data
  893. 39:42set over here.
  894. 39:44This is the resulting data set, the two
  895. 39:46partitions.
  896. 39:48And the joined value on the right-hand
  897. 39:50side. So, this is how a broadcast join
  898. 39:54broadcast hash join is going to work.
  899. 39:57Now, let's go ahead with the lab for
  900. 40:00the broadcast hash join. Right? And we
  901. 40:02are going to follow something very
  902. 40:04similar what we've done for the sort
  903. 40:06merge join. So, let's go ahead and set
  904. 40:08all of this up. Let's quickly check what
  905. 40:12is the auto broadcast join threshold,
  906. 40:14which is set to 10 megabytes. right?
  907. 40:18This is the default value. And as we
  908. 40:20discussed earlier, we are going to set
  909. 40:23AQE to false.
  910. 40:25And
  911. 40:26basically setting up the paths for both
  912. 40:28of these data sets.
  913. 40:30Defining both of these data sets, right?
  914. 40:34Setting the header to true and info
  915. 40:35schema to false.
  916. 40:37Sorry, true, not false. Yeah.
  917. 40:39And let's quickly have a look at both of
  918. 40:42these data sets. So, these data sets are
  919. 40:44just the same as what we've used for the
  920. 40:47sort merge join, right?
  921. 40:51So, now let's go ahead and
  922. 40:53perform
  923. 40:55the join, right? So, you see that the
  924. 40:56code is literally the same. DF orders
  925. 41:00uh joins with DF customers on customer
  926. 41:03ID and the join type is inner, yeah.
  927. 41:06So, one key difference is last time we
  928. 41:09specified we were forcing a sort merge
  929. 41:13join by specifying this to be minus one,
  930. 41:15but this time we have allowed it to have
  931. 41:19its default value, right?
  932. 41:22So, in this case, this should ideally
  933. 41:25trigger a broadcast hash join.
  934. 41:28So, let's go ahead and run this. Let's
  935. 41:30do a show. The results are going to be
  936. 41:33the same as the sort merge join, but
  937. 41:36what is important is let's look at
  938. 41:39the physical plan, yeah. So, let's copy
  939. 41:42this and let me put this over here.
  940. 41:45Which is broadcast hash join.
  941. 41:49Let's put the physical plan over here.
  942. 41:53And let's read this from bottom to top.
  943. 41:55So, it followed the normal drill, which
  944. 41:57is file scan CSV. It reads the customer
  945. 42:00data set, and then it basically figures
  946. 42:02out
  947. 42:03that it ensure that none of the customer
  948. 42:06ID should be null. So, this is a part of
  949. 42:09Spark's own optimization strategy. Now,
  950. 42:11the interesting part is this broadcast
  951. 42:14exchange and the hash relation broadcast
  952. 42:18mode, right? So, this broadcast exchange
  953. 42:20basically says
  954. 42:22that the driver on which the customer
  955. 42:25data set is residing, right? It is
  956. 42:28basically going to take that, serialize
  957. 42:30it, and send it over the network to the
  958. 42:34executors, right?
  959. 42:36So, that the executors can have their
  960. 42:38own copy of the small data set. Now,
  961. 42:42what does hash relation broadcast mode
  962. 42:44mean, right? So, the broadcast mode,
  963. 42:48yeah, this broadcast mode
  964. 42:50that we see over here, this broadcast
  965. 42:52mode
  966. 42:53it basically tells how the broadcast
  967. 42:56side is packaged before it is shipped to
  968. 42:58the executor, right? So, there are two
  969. 43:00types of broadcast mode. One is the hash
  970. 43:02relation broadcast mode, the other one
  971. 43:05is the identity broadcast mode. And we
  972. 43:07are going to look into identity
  973. 43:09broadcast mode a little later.
  974. 43:11The broadcast mode basically decides or
  975. 43:13it basically tells you how the broadcast
  976. 43:16side is packaged and sent to the
  977. 43:20executor. So, in case of a hash relation
  978. 43:23broadcast mode, which is used by a
  979. 43:25broadcast hash join,
  980. 43:26the broadcast rows are basically
  981. 43:28transformed into a hash map, right? It
  982. 43:31is transformed into hash map with the
  983. 43:33key and a value. In our case, the key is
  984. 43:37going to be the customer ID and the
  985. 43:39value is going to be the the row value,
  986. 43:42right? So, how is the key decided and
  987. 43:45that is what you see over here, right?
  988. 43:47So, input of zero is basically pointing
  989. 43:50to index zero, which is the customer ID.
  990. 43:52It says that customer ID is of data type
  991. 43:55integer and false means that it is
  992. 43:58non-nullable. It converts it into a long
  993. 44:02or big int for optimization purposes,
  994. 44:05right? And all of this is converted into
  995. 44:07a list. Right? So, it is a list where it
  996. 44:10has the join key. Now, the reason why
  997. 44:13it's a list because sometimes join key
  998. 44:14can be
  999. 44:16more than one. Right? It can be a
  1000. 44:17composite join key. So, that is why it's
  1001. 44:19put into a list. So, basically it
  1002. 44:21specifies that a broadcast exchange is
  1003. 44:24going to happen. The data on the driver
  1004. 44:27is going to be broadcasted over to the
  1005. 44:29executor. And the mode that is going to
  1006. 44:32be followed is a hash relation. Where
  1007. 44:35the key is going to look like this.
  1008. 44:38Right? It is a list with the customer
  1009. 44:41ID. Yeah?
  1010. 44:43So, that is what means from for these
  1011. 44:47three lines. And then
  1012. 44:49it reads the order data set.
  1013. 44:52Filters out all of the customer ID
  1014. 44:54wherever it is null. And finally, it
  1015. 44:57does a broadcast hash join on the
  1016. 45:01customer ID. And it also tells you that
  1017. 45:02this is an inner join. And build right
  1018. 45:05basically means which side of the join
  1019. 45:07is going to be used as a hash table. And
  1020. 45:11in this case, this is the right side.
  1021. 45:13And this is the left side. The top one
  1022. 45:16is the left side. Yeah? So, the right
  1023. 45:19side is the customer data set which is
  1024. 45:21going to be converted into a hash map.
  1025. 45:24So, by this step over here the join has
  1026. 45:27completed. And finally
  1027. 45:29the columns that need to be selected as
  1028. 45:32output are projected.
  1029. 45:34So, also have a look at the Spark UI.
  1030. 45:38Which is
  1031. 45:41An action was called over here. So, job
  1032. 45:4315 and 16 are the relevant ones. So, if
  1033. 45:46I go to jobs over here over here, these
  1034. 45:49are the two jobs. Yeah? So, let me go to
  1035. 45:53SQL {slash} data frame.
  1036. 45:55And this is the one that I should look
  1037. 45:57at. Yeah?
  1038. 45:58So, here
  1039. 46:01and let me open this up in another tab.
  1040. 46:06So here, what we see is
  1041. 46:09this data set, which is the customer
  1042. 46:11data set, has been scanned,
  1043. 46:13and then a filter has been applied over
  1044. 46:16the customer ID, which is what we saw in
  1045. 46:19the physical plan as well. And then you
  1046. 46:21see a broadcast exchange, right? And
  1047. 46:25this goes to the same place on the
  1048. 46:27left-hand side, where you have the
  1049. 46:29orders data set. On the left-hand side,
  1050. 46:32you have the order data set that was
  1051. 46:33scanned. It also went through a filter,
  1052. 46:36and finally both of them are broadcast
  1053. 46:40hash join, right? So you also see
  1054. 46:44there is time to read broadcast, right?
  1055. 46:46Which was basically the broadcast
  1056. 46:48exchange data set that comes over here,
  1057. 46:51yeah? And the reason why you see rows
  1058. 46:53output over here is six, because I
  1059. 46:55specified a show of five, yeah?
  1060. 46:59So finally, a broadcast hash join
  1061. 47:00happens in this step, and you see the
  1062. 47:04inner build right, just what we have
  1063. 47:06understood right now.
  1064. 47:08And finally, we select seven columns as
  1065. 47:10a result of the final data set. So this
  1066. 47:14is how a DAG is going to look like for a
  1067. 47:18broadcast hash join. Almost all
  1068. 47:22data sets which undergo a broadcast hash
  1069. 47:24join, they are going to look somewhat
  1070. 47:26similar. Okay, so the next join we are
  1071. 47:28going to study is shuffle hash join,
  1072. 47:32right? And shuffle hash join lies in
  1073. 47:35between a broadcast hash join and a sort
  1074. 47:38merge join.
  1075. 47:40So this is an interesting one. It lies
  1076. 47:42between a broadcast hash join
  1077. 47:46and a sort merge join.
  1078. 47:51Right? And why do I say that?
  1079. 47:53I say that because it gets picked up
  1080. 47:55when the situation is too big
  1081. 47:58for a broadcast hash join, but it also
  1082. 48:01at the same time doesn't need the full
  1083. 48:02machinery of the sort merge join, right?
  1084. 48:06So, when the situation is too big for a
  1085. 48:09broadcast hash join, broadcast hash join
  1086. 48:12will not work in that case, but you also
  1087. 48:14don't need the full machinery of a sort
  1088. 48:18merge join. So, in that case, Spark is
  1089. 48:21going to choose a shuffle hash join.
  1090. 48:23Yeah. But, what are the internal
  1091. 48:25mechanics? What are What is the logic
  1092. 48:27using which it is going to decide that,
  1093. 48:29right? So, the smaller table is too
  1094. 48:31large to broadcast to every executor.
  1095. 48:34Yeah. So, the smaller table basically
  1096. 48:37exceeds that property that we've been
  1097. 48:40discussing, which is Spark
  1098. 48:48Yeah. So, the smaller table is too large
  1099. 48:50to be able to broadcast to each
  1100. 48:52executor. So, the smaller table exceeds
  1101. 48:55this default value, which is 10 MB.
  1102. 49:00Yeah. And
  1103. 49:02one partition
  1104. 49:04one partition of that smaller table
  1105. 49:06after shuffling the the redistribution
  1106. 49:08of data can still fit into a single
  1107. 49:12executor memory in order to build the
  1108. 49:15hash table, right? So, the key point to
  1109. 49:17note is the smaller table is too large
  1110. 49:20to be able to fit into every executor,
  1111. 49:22but however one partition of the smaller
  1112. 49:26table after So, shuffling is going to
  1113. 49:28redistribute all of the data to the
  1114. 49:30executor.
  1115. 49:32Now, after redistribution, the
  1116. 49:34partitions of the smaller table that
  1117. 49:36land in each of the executors
  1118. 49:39they can still fit in the executor
  1119. 49:41memory in order to build the hash table,
  1120. 49:43right? So, that is the case when a
  1121. 49:47shuffle hash join is going to be chosen.
  1122. 49:49Now, a very important thing to note is
  1123. 49:52there is also a preference gate. Yeah.
  1124. 49:54Uh by preference gate, what I mean is uh
  1125. 49:57by default, Spark is going to prefer a
  1126. 50:00sort merge join.
  1127. 50:02It's going to prefer a sort merge join
  1128. 50:04over a shuffle hash join.
  1129. 50:08Yeah. Even when shuffle hash join is
  1130. 50:10value uh viable. Yeah. Even when shuffle
  1131. 50:13hash join is viable, it is going to
  1132. 50:15prefer a sort merge join. And the reason
  1133. 50:17for this is sort merge join is more
  1134. 50:21memory stable.
  1135. 50:25Right? It is more memory stable. And the
  1136. 50:27reason why I say it's more memory stable
  1137. 50:29is because shuffle hash join has to keep
  1138. 50:32the hash map in memory. And if there is
  1139. 50:35a key which is skewed, then there is a
  1140. 50:38good chance of an out of memory error.
  1141. 50:41Right? So, we just said that in a
  1142. 50:43shuffle hash join,
  1143. 50:45>> [snorts]
  1144. 50:46>> after shuffling, when shuffling does the
  1145. 50:48redistribution of data,
  1146. 50:50one of the partition, the partition that
  1147. 50:52gets redistributed to each of the
  1148. 50:54executor, that is small enough to build
  1149. 50:57the hash table.
  1150. 50:58But, what if when you are building the
  1151. 51:00hash table, the key there is a key which
  1152. 51:03which is skewed, right? So, in that
  1153. 51:05case, there are good chances of an out
  1154. 51:09of memory error. Yeah. So, if you want
  1155. 51:11to override Spark default setting, you
  1156. 51:13have to set this property, which is
  1157. 51:16spark
  1158. 51:17.sql.
  1159. 51:20join
  1160. 51:21.prefer
  1161. 51:24prefer sort merge join.
  1162. 51:27What are we going to do?
  1163. 51:28Right? So, once you do that, Spark is
  1164. 51:31going to pick up shuffle hash join when
  1165. 51:34the smaller side per partition data fits
  1166. 51:37in memory.
  1167. 51:38So, the shuffle hash join has three
  1168. 51:40steps. The first one is the shuffle step
  1169. 51:44where both the tables are reshuffled
  1170. 51:47across executor simply by following the
  1171. 51:49logic which is hash of the join key
  1172. 51:52modulo shuffle partitions.
  1173. 51:54Yeah.
  1174. 51:55So, this is basically to make sure that
  1175. 51:57all the rows with the same join key land
  1176. 52:00in the same partition in the same
  1177. 52:01executor, right? The second one is the
  1178. 52:04build step where each of the executors
  1179. 52:07takes its partition of the smaller side,
  1180. 52:09right? We discussed the smaller side and
  1181. 52:12then it builds a hash map in memory
  1182. 52:15keyed on the join column, right? So,
  1183. 52:18this is going to be join key and then
  1184. 52:20the row value.
  1185. 52:22And then finally the probe phase where
  1186. 52:25we're going to take the larger partition
  1187. 52:28row by row and then we're going to look
  1188. 52:30up in the hash table and see if there is
  1189. 52:32a hit, right? If there is a match and a
  1190. 52:34match is going to produce a joined row.
  1191. 52:38So, these are the three steps. Now,
  1192. 52:39let's look into it with visuals and
  1193. 52:42example. Okay, so these are our data
  1194. 52:46sets, right? We again have a customer
  1195. 52:48data set
  1196. 52:50and an order data set. We have two
  1197. 52:52executors,
  1198. 52:53right? One partition each
  1199. 52:56for customer and orders in each of the
  1200. 52:59executors, right? Now, the first step
  1201. 53:02that we mentioned was the shuffle step.
  1202. 53:07Right?
  1203. 53:09The shuffle step and
  1204. 53:11we're going to assume
  1205. 53:13three shuffle partition.
  1206. 53:15Right? This time we're going to assume
  1207. 53:17three shuffle partition. This means that
  1208. 53:19in the shuffle stage it is going to
  1209. 53:22produce three partitions per data set,
  1210. 53:24right? So, in order to understand where
  1211. 53:27this row is going to go to, we simply
  1212. 53:29have to follow the same logic which is
  1213. 53:31hash of the join key
  1214. 53:35modulo shuffle partition. And our join
  1215. 53:37key is basically the customer ID, right?
  1216. 53:40So, let's use the same logic, which is
  1217. 53:42hash of customer one modulo three. For
  1218. 53:46simplicity, we said we are just going to
  1219. 53:48do one mod three.
  1220. 53:51Right?
  1221. 53:52C1 is going to become one, and one mod
  1222. 53:54three, this is just one. Right? So, this
  1223. 53:57is going to be one. This is C2 mod
  1224. 54:02three is going to be two mod three,
  1225. 54:06which is equal to two. Right? So, this
  1226. 54:08is going to go to partition number two,
  1227. 54:11partition number zero, partition number
  1228. 54:13one, partition number one, partition
  1229. 54:15number zero, partition number zero,
  1230. 54:17partition number one, partition number
  1231. 54:19two, partition number zero. Right?
  1232. 54:22Similarly, this is also going to go to
  1233. 54:24partition number two,
  1234. 54:27and partition number zero. This is going
  1235. 54:29to go to partition one, partition two,
  1236. 54:31partition zero, partition one,
  1237. 54:34partition two, partition two,
  1238. 54:36partition zero, and partition number
  1239. 54:39two. Right? So, now, as a result, in the
  1240. 54:43next step, what you're going to see is
  1241. 54:45three partitions.
  1242. 54:48Three partitions, because we defined
  1243. 54:51three shuffle partition.
  1244. 54:54So, this is partition number zero,
  1245. 54:56partition number one, they reside on
  1246. 54:58executor number one, and partition
  1247. 55:01number two, this is residing on executor
  1248. 55:05number two.
  1249. 55:06Now, why two partitions on one executor
  1250. 55:08and one partition on the other?
  1251. 55:10And this is something this the Spark
  1252. 55:12scheduler decides. Right? So,
  1253. 55:15this one where all of this is modulo
  1254. 55:17zero, goes to partition number zero.
  1255. 55:23Wherever it was modulo one, it goes to
  1256. 55:25partition number one, and similarly,
  1257. 55:27wherever it was modulo two, it goes to
  1258. 55:30partition number two. And similarly for
  1259. 55:32the orders as well. This is all modulo
  1260. 55:34zero,
  1261. 55:36this is all modulo one and this is all
  1262. 55:39modulo two.
  1263. 55:40Right? So, this is post
  1264. 55:43shuffle
  1265. 55:46state.
  1266. 55:47Right?
  1267. 55:49This is how the partitions are going to
  1268. 55:51look like. Now, the next step the next
  1269. 55:54step was the build step. Right? So, in
  1270. 55:58the build step, what we are going to do
  1271. 56:00is we take the smaller side
  1272. 56:03we take the smaller side and we build a
  1273. 56:06hash map out of it. Right? So, we take
  1274. 56:09this side
  1275. 56:10we take this partition and we build a
  1276. 56:13hash map. And the hash map is simply
  1277. 56:15going to be of the format the join key
  1278. 56:19the join key to the row.
  1279. 56:22Right? So, here the join key is C3 and
  1280. 56:25C6
  1281. 56:26which is the customer ID. Customer ID is
  1282. 56:28the join key and this is the row value.
  1283. 56:31Right?
  1284. 56:32And similarly, it happens for this one
  1285. 56:33as well. The large partition stays as it
  1286. 56:36is during this step. So, in this step,
  1287. 56:38we have now created a hash map from the
  1288. 56:43smaller partition. Right?
  1289. 56:45Now, we are going to do the probing.
  1290. 56:48And in probing, what is going to happen
  1291. 56:50is each executor is going to stream its
  1292. 56:53partition of the larger side
  1293. 56:56which is this one, these three
  1294. 56:58partitions, right? Row by row and it is
  1295. 57:01going to look up each of the rows keys
  1296. 57:04in the local hash map. So, C3 is going
  1297. 57:07to check whether it exists over here or
  1298. 57:09not. And it did find a
  1299. 57:13hit. Therefore, it is going to match and
  1300. 57:15then you're going to get name
  1301. 57:18and then the city over here which is
  1302. 57:21Rohan
  1303. 57:22and Bengaluru. Right? Similarly, C3 is
  1304. 57:25going to match match 60 C6 is also going
  1305. 57:28to match over here.
  1306. 57:29Right? Similarly, this is also going to
  1307. 57:31happen over here. It is going to probe
  1308. 57:32this hash map
  1309. 57:35and find a hit. And similarly for this
  1310. 57:37one as well.
  1311. 57:40Okay, I think this has been
  1312. 57:43there's a mistake over here. This should
  1313. 57:46be C2 and this should be C5 and this
  1314. 57:49should be Okay, this These two names are
  1315. 57:51fine. And C2 is going to match over here
  1316. 57:54and C5 is going to match with Imdad over
  1317. 57:57here, right? So, this is how the matches
  1318. 58:00the join is going to produce hits rows,
  1319. 58:03right? And finally, as a result, you see
  1320. 58:06that there are three partition, right?
  1321. 58:08As a result of joining these two,
  1322. 58:12this one and this one.
  1323. 58:15Joining these two, this one and this
  1324. 58:17one. So, there are going to be two
  1325. 58:18output partitions
  1326. 58:20on
  1327. 58:21the first executor.
  1328. 58:27And one output partition on the second
  1329. 58:31executor. And that is what you see over
  1330. 58:33here.
  1331. 58:34Number one, number two, and
  1332. 58:37number three.
  1333. 58:39Yeah. So, this is how a shuffle hash
  1334. 58:43join works.
  1335. 58:44Okay, so let's go ahead with the lab for
  1336. 58:47the shuffle hash join. And you'll find a
  1337. 58:49lot of things to be similar. We do the
  1338. 58:51import. We set
  1339. 58:53spark.sql.adaptive.enabled
  1340. 58:56basically the AQE to false and shuffle
  1341. 58:59partitions to four. And we've read about
  1342. 59:03this very important property, which is
  1343. 59:05spark.sql.join.preferSortMergeJoin
  1344. 59:09to be false. So, this is the property
  1345. 59:12which enables Spark to choose shuffle
  1346. 59:15hash join whenever it is viable, right?
  1347. 59:19And the second important property is we
  1348. 59:21have lowered the threshold for broadcast
  1349. 59:24join from 10 MB to 1 MB, right? So, the
  1350. 59:27default is 10 MB. We have lowered it to
  1351. 59:311 MB. If we don't do this, our data set
  1352. 59:34is less than 10 MB, but it is greater
  1353. 59:36than 1 MB, right? So, it is
  1354. 59:39automatically going to trigger a
  1355. 59:40broadcast hash join, but we don't want
  1356. 59:43that, right? So, if you remember the
  1357. 59:45original definition,
  1358. 59:46the situation or the condition should be
  1359. 59:48such that
  1360. 59:49it should be bigger, it should be uh
  1361. 59:53higher for a broadcast hash join, but it
  1362. 59:56also doesn't need the full machinery of
  1363. 59:58a sort merge join. And that is why a
  1364. 1:00:01shuffle hash join falls in between.
  1365. 1:00:03Yeah? So, that is why we have lowered
  1366. 1:00:05the threshold so Spark doesn't
  1367. 1:00:08automatically do a broadcast hash join.
  1368. 1:00:11Yeah? But also, when the shuffle
  1369. 1:00:13happens, the data set is broken into
  1370. 1:00:15small small partition, and those small
  1371. 1:00:17partition, it is able to build a hash
  1372. 1:00:20table on those small partition, right?
  1373. 1:00:23So, I hope this makes sense. Now, we
  1374. 1:00:25this is the normal drill that we've been
  1375. 1:00:28repeating for quite some time now. Uh
  1376. 1:00:31we basically run this, and we set the
  1377. 1:00:34paths.
  1378. 1:00:35We read both the data frame. We have a
  1379. 1:00:38look at both of them, and then finally,
  1380. 1:00:40we do a join between the two data sets.
  1381. 1:00:44Yeah?
  1382. 1:00:46Now, before running show, let me just go
  1383. 1:00:48ahead and print the plan.
  1384. 1:00:51Yeah? Let me copy this.
  1385. 1:00:54Let's put
  1386. 1:00:57a shuffle hash join title over here.
  1387. 1:01:01Yeah. So, let's go ahead and read this
  1388. 1:01:04from the bottom to the top.
  1389. 1:01:06So, we first read the file, which is the
  1390. 1:01:08customer data set.
  1391. 1:01:10We simply filter out all of the null for
  1392. 1:01:13customer ID, and then a normal shuffle
  1393. 1:01:17happen, right? A normal shuffle happen.
  1394. 1:01:20Similarly, for the second data set, the
  1395. 1:01:22order data set, it is read, it is
  1396. 1:01:24scanned, and then a normal shuffle
  1397. 1:01:28happen, right? So, the three steps for
  1398. 1:01:31both of them, they look identical. And
  1399. 1:01:34this is normal, right? The hash is built
  1400. 1:01:37after the shuffling, and that is what
  1401. 1:01:39you see over here.
  1402. 1:01:41Once both of them have been shuffled, a
  1403. 1:01:43shuffle hash join takes place, and you
  1404. 1:01:46see this is an inner join. The keyword
  1405. 1:01:49to note over here is build right. So,
  1406. 1:01:52build right basically specifies the
  1407. 1:01:55table on which the hash map is to be
  1408. 1:01:58built, right? And this is the right
  1409. 1:02:00table, which is the customer data set.
  1410. 1:02:02This is the left table. Yeah? So, uh
  1411. 1:02:06hash table is built on the customer data
  1412. 1:02:09set, and then a shuffle hash join is
  1413. 1:02:12performed. And finally, the data uh a
  1414. 1:02:15project operation is called, which
  1415. 1:02:16selects the relevant columns. Yeah?
  1416. 1:02:19Okay, let's go ahead and also have a
  1417. 1:02:21look at the DAG. And for this, we are
  1418. 1:02:24going to run this action over here, and
  1419. 1:02:27this produces one Spark job, which is
  1420. 1:02:29job number 40. Yeah? So, let's quickly
  1421. 1:02:32open this up. If I look at job, this is
  1422. 1:02:35basically job number 40, and it has
  1423. 1:02:39three stages, right? It has three stages
  1424. 1:02:42that you would find over here.
  1425. 1:02:44But, let me go to the SQL data frame.
  1426. 1:02:48Yeah, let's click over here.
  1427. 1:02:51And actually, the structure looks very
  1428. 1:02:53similar to the sort merge join. Yeah?
  1429. 1:02:56So, so on the left-hand side, you see
  1430. 1:02:59scanning both of the data sets, right?
  1431. 1:03:02The this one on the right is the
  1432. 1:03:04customer ID, and then on the left is the
  1433. 1:03:06order ID, order order data set, right?
  1434. 1:03:10Um then, you are going to filter on the
  1435. 1:03:13customer ID not equal to null, and
  1436. 1:03:15similarly for this one, something that
  1437. 1:03:16we've seen in the plan already.
  1438. 1:03:19And now there is going to be an exchange
  1439. 1:03:21hash partitioning that we see through
  1440. 1:03:23this step, the left and the right. And
  1441. 1:03:26finally, there is going to be a shuffle
  1442. 1:03:28hash join. But where do we actually see
  1443. 1:03:31that the hash map is being built? So
  1444. 1:03:33that is what you see over here, time to
  1445. 1:03:36build hash map, which is 6 ms, right? So
  1446. 1:03:40this is how
  1447. 1:03:42it shows that a hash map was built and
  1448. 1:03:44then it was probed to find the matches,
  1449. 1:03:48right? And that is what is also
  1450. 1:03:50specified by build right. And on the
  1451. 1:03:52right side, you have the customer data
  1452. 1:03:55set. And finally, the shuffle hash join
  1453. 1:03:58is performed and it projects eight
  1454. 1:04:00columns, which is the final output of
  1455. 1:04:03the data set. This is how a DAG will
  1456. 1:04:06look like for a shuffle hash join. Okay,
  1457. 1:04:09so the final type of join that we're
  1458. 1:04:12going to study is the broadcast nested
  1459. 1:04:15loop join, right? So let's first
  1460. 1:04:18understand when is it performed, when is
  1461. 1:04:20it chosen by Spark
  1462. 1:04:23over any other join, right?
  1463. 1:04:25So the first point is a non-equi join
  1464. 1:04:29condition, right? So when the join
  1465. 1:04:31predicate contains no equality clause.
  1466. 1:04:34So when you're doing a join, let's say
  1467. 1:04:35you are doing DF1.x
  1468. 1:04:38= DF
  1469. 1:04:412.x.
  1470. 1:04:43Instead of equal to, you have a non-equi
  1471. 1:04:46join predicate, right? And what that
  1472. 1:04:48means is you have something like equal
  1473. 1:04:51to, greater than equal to, less than
  1474. 1:04:53equal to, less than,
  1475. 1:04:55not equal to, or something like a range,
  1476. 1:04:59right? So for example, if you use
  1477. 1:05:01between,
  1478. 1:05:03right? So DF1.x
  1479. 1:05:05between
  1480. 1:05:07some value and DF2.x, right? Something
  1481. 1:05:10like that. So basically, it has to be a
  1482. 1:05:12non-equi join condition. It is either an
  1483. 1:05:16inequality, a range, or a complex
  1484. 1:05:19condition, right? So, in that case, a
  1485. 1:05:22broadcast nested loop join will be
  1486. 1:05:25chosen. Another case where a broadcast
  1487. 1:05:28nested loop join is chosen is during a
  1488. 1:05:32cross join, when there is no join
  1489. 1:05:34condition at all, right? So, every row
  1490. 1:05:36from one table must be paired with every
  1491. 1:05:39row from the other table, right? So,
  1492. 1:05:42these are the two conditions when Spark
  1493. 1:05:45is going to choose a broadcast nested
  1494. 1:05:48loop join.
  1495. 1:05:49So, now, what exactly are the steps? The
  1496. 1:05:52steps are twofold. The first one is the
  1497. 1:05:54broadcast step. And this broadcast is
  1498. 1:05:58identical to the broadcast hash join, in
  1499. 1:06:00which the smaller table is serialized by
  1500. 1:06:04the driver, and it is pushed in full to
  1501. 1:06:07every executor, right? No shuffle of the
  1502. 1:06:10large table.
  1503. 1:06:11And this also brings me to one more
  1504. 1:06:13condition, right? If this broadcast
  1505. 1:06:16needs to happen, then
  1506. 1:06:19if we have two tables,
  1507. 1:06:21if we have two tables, then one of the
  1508. 1:06:24table needs to be small to be able to be
  1509. 1:06:26broadcasted, and the other table can be
  1510. 1:06:30large, right? So, along with these two
  1511. 1:06:32conditions, this is also a very
  1512. 1:06:36important condition. Otherwise, the
  1513. 1:06:38broadcast step cannot happen, right? So,
  1514. 1:06:41that is why the first step is the
  1515. 1:06:43broadcast step, where you take the
  1516. 1:06:44smaller table, it's serialized by the
  1517. 1:06:46driver, and it is pushed in full to
  1518. 1:06:49every executor. Yeah? And there is no
  1519. 1:06:52shuffle of the large table.
  1520. 1:06:54And step number two,
  1521. 1:06:57step number two is the nested loop,
  1522. 1:06:59right? So, every executor runs a double
  1523. 1:07:02loop over its partition of the large
  1524. 1:07:04table, and the full broadcast table, and
  1525. 1:07:07then it evaluates the join condition.
  1526. 1:07:10So, it looks something like this. So,
  1527. 1:07:11for
  1528. 1:07:12row
  1529. 1:07:14in large
  1530. 1:07:16partition right?
  1531. 1:07:20For row in large partition, and then
  1532. 1:07:21again for row
  1533. 1:07:23in broadcast table
  1534. 1:07:28For row in broadcast table, yeah.
  1535. 1:07:30And if the condition
  1536. 1:07:34if the condition
  1537. 1:07:36of this row
  1538. 1:07:38which is the la L row and B row
  1539. 1:07:42Let me name it like that.
  1540. 1:07:44L row and B row
  1541. 1:07:47they match, then we are going to emit
  1542. 1:07:51the joined row.
  1543. 1:07:54Right? Then we are going to call that a
  1544. 1:07:57join has happened. Yeah. So, essentially
  1545. 1:08:00there are two steps, which is the
  1546. 1:08:01broadcast step, and then a nested loop
  1547. 1:08:04step. The nested loop step is basically
  1548. 1:08:06a
  1549. 1:08:07two loop, a dual loop between the large
  1550. 1:08:10and large partition and the broadcast
  1551. 1:08:13table. And if the condition matches,
  1552. 1:08:14then we emit the joined row. Now, let's
  1553. 1:08:17understand this through visuals and an
  1554. 1:08:20example. Okay, now let's take this
  1555. 1:08:22example where I have a customer data set
  1556. 1:08:26and I also have an order data set,
  1557. 1:08:28right?
  1558. 1:08:29The customer has a credit limit, and the
  1559. 1:08:33orders has which customer
  1560. 1:08:36spent what amount on a particular order,
  1561. 1:08:39right? And the problem statement is that
  1562. 1:08:42I want to figure out for each of the
  1563. 1:08:44customers
  1564. 1:08:45what are the orders in which the amount
  1565. 1:08:48is greater than my credit limit, right?
  1566. 1:08:51So, for each of the customers, I want to
  1567. 1:08:53figure out all orders. It can be my
  1568. 1:08:55order, it can be somebody else's order.
  1569. 1:08:57The only thing I want to figure out
  1570. 1:08:59which are the orders in which the amount
  1571. 1:09:02of the order is greater than my credit
  1572. 1:09:05limit, right? So, the join condition
  1573. 1:09:08the join condition is this order.amount
  1574. 1:09:13orders.amount
  1575. 1:09:17should be greater than
  1576. 1:09:19order.amount should be greater than
  1577. 1:09:21customer's
  1578. 1:09:25.credit limit.
  1579. 1:09:29Yeah?
  1580. 1:09:30And we see that in the join condition
  1581. 1:09:32there is an inequality over here. And
  1582. 1:09:36for that reason, we are going to choose,
  1583. 1:09:39or rather Spark is going to choose a
  1584. 1:09:42broadcast nested loop join.
  1585. 1:09:46Right? So, this table, let's assume that
  1586. 1:09:48the size of this table is 320 bytes. And
  1587. 1:09:51this is much lesser than the 10
  1588. 1:09:53megabytes default auto broadcast join
  1589. 1:09:56threshold size. That means this table
  1590. 1:09:59can be broadcasted
  1591. 1:10:01this table can be broadcasted to both of
  1592. 1:10:05the executors.
  1593. 1:10:06And both of them are going to receive
  1594. 1:10:08their individual copy of the customer's
  1595. 1:10:12table.
  1596. 1:10:13Right? So, both of them are going to
  1597. 1:10:16receive an individual copy of the
  1598. 1:10:19customer's table.
  1599. 1:10:21Right? So, this is the first step, which
  1600. 1:10:23is the broadcast step.
  1601. 1:10:28The broadcast step, right? Now, let's
  1602. 1:10:30move over to the next step, which is the
  1603. 1:10:33nested loop join step.
  1604. 1:10:36Nested loop step.
  1605. 1:10:42So, in order for the nested loop join to
  1606. 1:10:44happen,
  1607. 1:10:46let's revisit this condition, right?
  1608. 1:10:48Where we are saying order amount is
  1609. 1:10:50greater than customer.credit limit. So,
  1610. 1:10:52if we were to look at this row,
  1611. 1:10:55this row, the order amount has to be
  1612. 1:10:57compared with this row, this row, and
  1613. 1:11:00this row.
  1614. 1:11:01So, we have to figure out whether this
  1615. 1:11:03order amount is it greater than this
  1616. 1:11:05credit limit, which is yes, 250 is
  1617. 1:11:08greater than 200. Then, it is going to
  1618. 1:11:10join this and this row. Similarly, this
  1619. 1:11:13row is going to be
  1620. 1:11:15checked against this row,
  1621. 1:11:17where 250 is it greater than 350? No,
  1622. 1:11:21that means this row is not going to be
  1623. 1:11:22joined, right? And for this one, 250 is
  1624. 1:11:25it greater than 100? Yes, that means
  1625. 1:11:28this row is going to be joined. So, in
  1626. 1:11:29order to visualize that, in order to
  1627. 1:11:32visualize that, what I've done is
  1628. 1:11:34comparing this row with all of the three
  1629. 1:11:37customers' credit limit, right? So, 250
  1630. 1:11:40and
  1631. 1:11:42where the credit limit was 200 and 100,
  1632. 1:11:45these are the output, which is 200 and
  1633. 1:11:48100 over here, C1
  1634. 1:11:52and C3, right? C2 will not be the
  1635. 1:11:55output, right?
  1636. 1:11:57So, this row will join with this one,
  1637. 1:11:58which is C1 and C3. Similarly, this row,
  1638. 1:12:02which is
  1639. 1:12:03this order number 102, will join or will
  1640. 1:12:06not join with anything, actually,
  1641. 1:12:08because this is the amount is 80, the
  1642. 1:12:10credit limit is 200, 350, and 100.
  1643. 1:12:14And similarly, this row will join with
  1644. 1:12:16all of the rows, because 400 is greater
  1645. 1:12:18than 200, 350, and 100, right? So, as a
  1646. 1:12:23result, you're going to see
  1647. 1:12:25two rows, which is one and
  1648. 1:12:28one and two
  1649. 1:12:32will be the output produced for this
  1650. 1:12:33row.
  1651. 1:12:34Zero for this one.
  1652. 1:12:38And three for this one, right? One
  1653. 1:12:42one, two, and three. So, as a result,
  1654. 1:12:44five rows are going to be produced,
  1655. 1:12:46right? And that is what you see over
  1656. 1:12:48here, one
  1657. 1:12:50one, two, three, three, four, and 5,
  1658. 1:12:53right? And this is basically C1 will
  1659. 1:12:56have two rows.
  1660. 1:12:58C3 will have three rows.
  1661. 1:13:00C1 is having two rows and C3 will have
  1662. 1:13:03three rows, right? And the similar case
  1663. 1:13:05is going to happen for this one as well
  1664. 1:13:08where order 104, it is going to be
  1665. 1:13:10matched 150.
  1666. 1:13:13150, is it greater than 200? No, 150 is
  1667. 1:13:16not greater than 350. 150 is greater
  1668. 1:13:18than 100, so it will produce one row.
  1669. 1:13:21And similarly, this is going to produce
  1670. 1:13:23two rows.
  1671. 1:13:25Right? So, you're going to see 2 + 1
  1672. 1:13:28which is three rows as output over here.
  1673. 1:13:34Three rows as output over here.
  1674. 1:13:37Yeah.
  1675. 1:13:382 + 1
  1676. 1:13:40which is three rows over here. Yeah. So,
  1677. 1:13:42this is how a broadcast nested loop join
  1678. 1:13:45works. Okay, now let's get started with
  1679. 1:13:48the final lab.
  1680. 1:13:51Which is the lab for broadcast nested
  1681. 1:13:54loop join.
  1682. 1:13:55And again, let's go ahead and run the
  1683. 1:13:57import disable AQE.
  1684. 1:14:01Right? We read the orders data frame,
  1685. 1:14:05right? And here we are going to do
  1686. 1:14:06something a little different, right? So,
  1687. 1:14:09let's read the order data frame. Let's
  1688. 1:14:11display the orders data frame.
  1689. 1:14:14And here I'm creating a tiers data
  1690. 1:14:18frame, right? And this basically
  1691. 1:14:20contains three columns which is base
  1692. 1:14:22which is the tier economy, standard,
  1693. 1:14:26premium, or enterprise.
  1694. 1:14:28And it basically specify
  1695. 1:14:32which order amount the range of order
  1696. 1:14:34amounts that qualify for that tier,
  1697. 1:14:38right?
  1698. 1:14:39So, let's go ahead and run this as well.
  1699. 1:14:44So, the meaning of this is basically if
  1700. 1:14:46let's say I've placed an order
  1701. 1:14:48and the value of that falls somewhere
  1702. 1:14:51between 1 and 50,
  1703. 1:14:53then it is considered to be an economy
  1704. 1:14:55order.
  1705. 1:14:56Right? And similarly between 51 to 150
  1706. 1:14:59it's the standard order, premium order
  1707. 1:15:01for 151 to 400, right?
  1708. 1:15:05And finally uh between this range to be
  1709. 1:15:08enterprises. Yeah.
  1710. 1:15:10So, there are companies who would want
  1711. 1:15:12to place certain benefits for each
  1712. 1:15:14category or tier of orders. Yeah.
  1713. 1:15:17So, this is how your data frame is going
  1714. 1:15:19to look like. Now, for each of these
  1715. 1:15:21orders, I want to classify in which of
  1716. 1:15:25the tier is it going to fall, right? So,
  1717. 1:15:28the join condition is going to look
  1718. 1:15:30something like this.
  1719. 1:15:31Right?
  1720. 1:15:33So, here I am forcing a broadcast nested
  1721. 1:15:36loop join.
  1722. 1:15:37Uh actually not a broadcast nested loop
  1723. 1:15:39join, but a broadcast step over here.
  1724. 1:15:41Nested loop is something that will be
  1725. 1:15:44automatically chosen because of this
  1726. 1:15:47condition over here. This condition
  1727. 1:15:49contains inequalities, right? So, that
  1728. 1:15:51is why a nested loop join will be
  1729. 1:15:54chosen. And here we are enforcing a
  1730. 1:15:57broadcast of the smaller data set.
  1731. 1:16:00Right? So, the join condition is simply
  1732. 1:16:03uh simply this one which is amount
  1733. 1:16:05should be greater than the minimum
  1734. 1:16:07amount and it should be less than the
  1735. 1:16:10maximum amount. Yeah.
  1736. 1:16:12So, let's go ahead and this is an inner
  1737. 1:16:14join. So, let's go ahead and run this.
  1738. 1:16:17Let's do a show.
  1739. 1:16:20This is going to trigger an action and
  1740. 1:16:22hence is going to trigger Spark jobs.
  1741. 1:16:25And therefore you see over here
  1742. 1:16:28this amount which is 160
  1743. 1:16:31and this fall in the premium tier
  1744. 1:16:33because the minimum and the maximum
  1745. 1:16:36range for any order to fall in the
  1746. 1:16:39premium tier is 151 to 400, right? And
  1747. 1:16:43similarly, we do for the other rows as
  1748. 1:16:45well, right? Now, let's have a look at
  1749. 1:16:49how the plan is going to look like. So,
  1750. 1:16:52let's go ahead and run this.
  1751. 1:16:54And now, let's copy this over broadcast
  1752. 1:16:57nested loop join, and let's paste this
  1753. 1:16:59over here.
  1754. 1:17:01So, let's go through this one by one.
  1755. 1:17:03This is a local table scan basically for
  1756. 1:17:07the data frame that we created locally
  1757. 1:17:11in the Spark session, yeah. And it does
  1758. 1:17:15a filter of the not null condition
  1759. 1:17:17because we're using two of the columns
  1760. 1:17:20in the join condition, right? We're
  1761. 1:17:22using minimum amount and the maximum
  1762. 1:17:24amount. So, it is making sure that both
  1763. 1:17:26of them is not null. And finally, it
  1764. 1:17:30does a broadcast exchange with the
  1765. 1:17:34identity broadcast mode, right? So, we
  1766. 1:17:37understood that broadcast exchange
  1767. 1:17:38simply means that it is going to take
  1768. 1:17:40this data set and serialize it and send
  1769. 1:17:43it over the network to the executor,
  1770. 1:17:46right? Now, what does the identity
  1771. 1:17:48broadcast mode mean? So, first of all,
  1772. 1:17:51broadcast mode, as we discussed earlier,
  1773. 1:17:54it tells you how the broadcast side is
  1774. 1:17:57packaged before it is shipped to the
  1775. 1:18:00executor, right? How the data set data
  1776. 1:18:03set is packaged before it is shipped
  1777. 1:18:05shipped to the executor. Last time we
  1778. 1:18:08studied about the hash relation
  1779. 1:18:10broadcast mode, in which it is packaged
  1780. 1:18:13in the form of a hash map, and then
  1781. 1:18:16broadcasted over to the executor. In the
  1782. 1:18:20identity broadcast mode, this transform
  1783. 1:18:23is actually an identity function, right?
  1784. 1:18:26So, what it simply means is that it is
  1785. 1:18:27going to take the data set as is as a
  1786. 1:18:30plain array, and then send it over to
  1787. 1:18:33the executor. So, there is no hashing,
  1788. 1:18:36no keying, none of that is going to
  1789. 1:18:38happen. It is simply going to take it
  1790. 1:18:40and then pass it to the executors as an
  1791. 1:18:43array. Yeah?
  1792. 1:18:45Now, the second data set, which is the
  1793. 1:18:47orders data set, which is read in, we
  1794. 1:18:49simply do a not null over the amount
  1795. 1:18:52because the amount is the one which is
  1796. 1:18:55used for the orders data frame.
  1797. 1:18:58And finally, we do a nested loop join
  1798. 1:19:02over here. Yeah? So, build right build
  1799. 1:19:05right over here. The right data set is
  1800. 1:19:08this one over here. Build right simply
  1801. 1:19:10means that this is the data set which is
  1802. 1:19:14broadcasted, right? Which Spark holds in
  1803. 1:19:16memory and is broadcasted. And
  1804. 1:19:19therefore, an inner join happens with
  1805. 1:19:21this condition. So, that is how the
  1806. 1:19:23physical plan is going to look like for
  1807. 1:19:26a broadcast nested loop join.
  1808. 1:19:29Now, let's have a look at how the Spark
  1809. 1:19:31UI looks for
  1810. 1:19:33the broadcast nested loop join. And this
  1811. 1:19:36is job number three and four. Let's open
  1812. 1:19:39this up. Let's go to the jobs.
  1813. 1:19:43And here you see job three and four,
  1814. 1:19:44which is one and one two stages, right?
  1815. 1:19:48Uh let's go to SQL data frame, this is
  1816. 1:19:51for job three and four.
  1817. 1:19:56Yeah?
  1818. 1:19:56And you see this one, which is the first
  1819. 1:20:00one is your scan CSV, right?
  1820. 1:20:03Which is on the left-hand side, the
  1821. 1:20:05orders data set, and on the right-hand
  1822. 1:20:07side is the data set that we have
  1823. 1:20:08created. Both of them go through this
  1824. 1:20:10filter optimization that we've seen, and
  1825. 1:20:13this is again a part of Spark own
  1826. 1:20:15optimization strategy. And then we see a
  1827. 1:20:18broadcast exchange using the identity
  1828. 1:20:21broadcast mode. This is broadcasted over
  1829. 1:20:25to the same location where you have the
  1830. 1:20:29orders data set, right? And then there
  1831. 1:20:31is a broadcast nested loop join using
  1832. 1:20:34this condition that we have specified
  1833. 1:20:36that the amount should be between the
  1834. 1:20:39minimum and the maximum amount, right?
  1835. 1:20:43And finally, you see that seven columns
  1836. 1:20:45are projected as an output for the final
  1837. 1:20:49data set. So, this is how the DAG is
  1838. 1:20:51going to look like for a broadcast
  1839. 1:20:54nested loop join. I hope this gave you
  1840. 1:20:56an overview on the kind of joins we have
  1841. 1:20:58in Apache Spark and how do all of them
  1842. 1:21:01behave under the hood and behind the
  1843. 1:21:04scenes. If you like this video, you are
  1844. 1:21:06also going to love the Apache Spark
  1845. 1:21:09playlist where we go in-depth into a lot
  1846. 1:21:12of different topics in Apache Spark. And
  1847. 1:21:14I also produce content on various other
  1848. 1:21:17topics on data engineering and interview
  1849. 1:21:19preparation. Please don't forget to like
  1850. 1:21:21and share this video and subscribe to
  1851. 1:21:23the channel. I will see you in the next
  1852. 1:21:25one. Thank you for watching.

About this transcript

This page contains the full transcript of Mastering Joins In Apache Spark: Complete Deep Dive by Afaque Ahmad, generated from the public captions YouTube serves with the video. The transcript has 12,154 words across 1,852 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.