Mastering Joins In Apache Spark: Complete Deep Dive — Transcript
Full transcript
- 0:00Joins are one of the most commonly and
- 0:03widely used operations in Apache Spark.
- 0:06In most of the code that you write on a
- 0:09day-to-day basis in Apache Spark, you
- 0:12must have used joins, right? And joins
- 0:15is actually where most of the problems
- 0:18in Apache Spark lives, right? So, in
- 0:20this video, we are going to cover four
- 0:24very important Spark physical joins in
- 0:27depth, right? So, we are going to cover
- 0:31We are going to cover sort merge join.
- 0:35We are going to cover broadcast hash
- 0:36join, shuffle hash join, and broadcast
- 0:41nested loop join. Of course, there are
- 0:42other joins, but understanding these
- 0:44four core joins is going to help you
- 0:47develop an intuition for the rest of the
- 0:50joins, right? So, for each of these
- 0:53joins, right? So, for each of these four
- 0:57joins, we are going to do three things.
- 1:00Right? The first one is we are going to
- 1:03understand what are the exact conditions
- 1:06that make Spark choose that join
- 1:09strategy, right? So, Spark has sort
- 1:12merge join available, Spark has
- 1:14broadcast hash join available, and then
- 1:16it also has these other two types of
- 1:18join available, right?
- 1:20What exactly are the conditions that
- 1:23make Spark choose, let's say, a sort
- 1:26merge join strategy over a broadcast
- 1:29hash join? So, that is the first thing
- 1:31that we are going to understand, the
- 1:33exact conditions that make Spark choose
- 1:36a particular join strategy, right?
- 1:39Number two, we are going to walk through
- 1:42step-by-step how each of these join
- 1:45strategy work under the hood, right? So,
- 1:47we are going to walk through visuals and
- 1:49understand each of the steps on how the
- 1:53join strategy work under the hood. Yeah?
- 1:56And finally, we are going to prove it
- 1:57with real code by looking at the Spark
- 2:00physical plan, the Spark physical UI,
- 2:03right? On how each of these joins are
- 2:08going to look like, right? So, how does
- 2:10the Spark physical plan look like for a
- 2:13sort merge join? How does it look like
- 2:16on the Spark UI, right? You're going to
- 2:18have a lot more clarity after going
- 2:22through all of these three steps. So,
- 2:24now let's get started with the sort
- 2:26merge join. Now, as data and AI
- 2:28engineers, we obsess a lot over how data
- 2:32moves, which is through joins,
- 2:34partitions, and shuffles. But, getting
- 2:37clean and structured data from the web
- 2:40in the first place, that is still one of
- 2:42the messiest parts of the job. SerpApi
- 2:45takes care of exactly that. This video
- 2:49is sponsored by SerpApi. SerpApi gives
- 2:52you clean, structured JSON results from
- 2:55Google, YouTube, Google Scholar, and
- 2:58many more, all through a single API
- 3:02call. So, they take care of all of the
- 3:04infrastructure headaches, CAPTCHA
- 3:06solving, proxy rotation, layout changes,
- 3:09none of that is your problem. So, you
- 3:12get reliable, structured data ready to
- 3:14plug straight into your pipeline. For
- 3:17people building in this space, two
- 3:19things stand out for me. Firstly, if you
- 3:21are training ML or AI models, SerpApi
- 3:25lets you build custom data sets using
- 3:28real-time search data. So, things like
- 3:30image titles, thumbnails, URLs,
- 3:33peer-reviewed articles from Google
- 3:35Scholar, all of them structured and
- 3:38ready to go. Secondly, if you are
- 3:41working on portfolio projects, this is
- 3:43genuinely a great way to power your apps
- 3:46with live, production-like data from
- 3:49some of the world's top search engines
- 3:52in JSON format that plays out really
- 3:54well with Spark and Databricks. They run
- 3:57at 99.9%
- 3:59uptime with a 1.2 second average
- 4:03response time. So, it's really
- 4:05production grade. Get started with 250
- 4:07free credits by clicking the link in the
- 4:09description or scanning the QR code
- 4:12here.
- 4:13All right. Now, let's get back to sort
- 4:15merge join. Joins are one of the most
- 4:19commonly and widely used operations in
- 4:21Apache Spark. In most of the code that
- 4:24you write on a day-to-day basis in
- 4:26Apache Spark, you must have used joins,
- 4:30right? And joins is actually where most
- 4:33of the problems in Apache Spark lives,
- 4:35right? So, in this video, we are going
- 4:38to cover four very important Spark
- 4:42physical joins in depth, right? So, we
- 4:45are going to cover
- 4:47We are going to cover sort merge join.
- 4:51We are going to cover broadcast hash
- 4:52join, shuffle hash join, and broadcast
- 4:56nested loop join. Of course, there are
- 4:58other joins, but understanding these
- 5:00four core joins is going to help you
- 5:03develop an intuition for the rest of the
- 5:06joins, right? So, for each of these
- 5:09join, right? So, for each of these four
- 5:12joins, we are going to do three things,
- 5:16right? The first one is we are going to
- 5:18understand what are the exact conditions
- 5:22that make Spark choose that join
- 5:25strategy, right? So, Spark has sort
- 5:28merge join available. Spark has
- 5:30broadcast hash join available. And then,
- 5:32it also has these other two types of
- 5:34join available, right?
- 5:36What exactly are the conditions that
- 5:39make Spark choose, let's say, a sort
- 5:42merge join strategy over a broadcast
- 5:45hash join? So, that is the first thing
- 5:47that we are going to understand. The
- 5:49exact conditions that make Spark choose
- 5:52a particular join strategy, right?
- 5:55Number two, we are going to walk through
- 5:58step by step how each of the join
- 6:01strategy work under the hood, right? So,
- 6:03we are going to walk through visuals and
- 6:05understand each of the steps on how the
- 6:08join strategy work under the hood, yeah?
- 6:11And finally, we are going to prove it
- 6:13with real code by looking at the Spark
- 6:16physical plan, the Spark physical UI,
- 6:19right? On how each of the joins are
- 6:23going to look like, right? So, how does
- 6:26the Spark physical plan look like for a
- 6:28sort merge join? How does it look like
- 6:31on the Spark UI, right? You're going to
- 6:34have a lot more clarity after going
- 6:37through all of these three steps. So,
- 6:39now let's get started with the sort
- 6:41merge join. Okay, so the first type of
- 6:43join that we are going to study is a
- 6:45sort merge join, right? So, first of
- 6:49all, let's understand when is it exactly
- 6:52performed? When is it When is it exactly
- 6:54chosen by Spark over any other join
- 6:58strategy. So, most importantly, there
- 7:00are
- 7:01two conditions, right? The first one is
- 7:04when both the data frames are large,
- 7:07right? And what this means is that when
- 7:10both the sides exceed a property, right?
- 7:14So, there is a property which is
- 7:20spark.sql.autobroadcastjointhreshold.
- 7:24And this
- 7:25the default
- 7:27value is set to 10 MB. It is sent set to
- 7:3110 MB, right? So, when both of the sides
- 7:34exceed this value, neither of them can
- 7:37be broadcasted, right? So, Spark is
- 7:39going to fall back to sort merge join.
- 7:45Right? And it is the only join that
- 7:47handles unlimited size on both sides
- 7:51without loading either of the data sets
- 7:53fully into memory, right? So, the first
- 7:55condition is both of the data frame
- 7:58should be large. And what this means is
- 8:01when both of the sides is greater than
- 8:04this auto broadcast join threshold
- 8:06property, which is
- 8:0710 megabytes, and as a result of which
- 8:10because
- 8:11they are higher than this threshold,
- 8:14none of them can be broadcasted, and
- 8:15therefore Spark falls back to this sort
- 8:19merge join strategy. The second one is
- 8:22equi-join on sortable keys, right? So,
- 8:26the key the keyword over here is
- 8:28equi-join.
- 8:30What equi-join mean over here is when
- 8:32you have your join condition, which is
- 8:34let's say
- 8:37DF1.x
- 8:39= DF2.x,
- 8:41right? So, this operator over here this
- 8:44operator over here needs to be an equal
- 8:47to.
- 8:48Right? It cannot be greater than,
- 8:50greater than equal to, less than, less
- 8:52than equal to, not equal to.
- 8:54Right? For all of such conditions, sort
- 8:58merge join will not work, right? And the
- 9:00keys needs to be sortable so that they
- 9:02can be put in order. And as you will see
- 9:04that sorting is one of the steps of sort
- 9:08merge join, right? So,
- 9:10the keys need to be sortable, yeah? So,
- 9:13these are the two most important
- 9:14conditions that need to be satisfied in
- 9:18order for a sort merge join to be chosen
- 9:21by Apache Spark, right? So, the three
- 9:24steps the three steps of sort merge
- 9:27join, and we're going to see it in a
- 9:28visualization,
- 9:30uh but to keep it concrete, the the
- 9:32three steps is shuffle,
- 9:35sort, and merge. By shuffling, what we
- 9:38mean is that
- 9:40all of the rows with the same joint key
- 9:42end up together in the same machine and
- 9:45in the same partition, right? So, the
- 9:48same keys after a particular hash after
- 9:51doing a particular hash hash of the key
- 9:55and this is the
- 9:57joint key.
- 10:00Right? Whatever value they get, right?
- 10:03Let's say this is X.
- 10:05And some other keys get Y. All of the
- 10:07keys which return X, they are going to
- 10:10end up in the same partition. Let's say
- 10:12it's called partition zero. Y is going
- 10:14to end up in partition one and so on,
- 10:16right? And then there is the sort step
- 10:19and finally a merge step. So, we are
- 10:20going to see all of this in action now.
- 10:23Okay, so let's understand it with this
- 10:25visualization, yeah?
- 10:28So, here we have two executors first of
- 10:31all, right?
- 10:33And in each of these two executors, we
- 10:35have two data sets. As you see, we have
- 10:38a customer data set and we have an order
- 10:42data set. In the first executor, we have
- 10:44two partitions for the customer data
- 10:46set, which is partition zero and
- 10:49partition number one.
- 10:51And we have one partition
- 10:55for the customer data set in the second
- 10:57executor and two partitions for the
- 11:01orders data set
- 11:03in the second executor, right? So, this
- 11:05is how the data set is distributed. You
- 11:07have a customer data set and an order
- 11:10data set. And now we want to do a join
- 11:14between the two. So, the join key is
- 11:17customer ID, right?
- 11:19The joint key is
- 11:25The joint key is customer ID, right?
- 11:29So, now once we figured out that the
- 11:31joint key is customer ID, the first step
- 11:34is going to be the shuffle step, right?
- 11:37The first step is going to be the
- 11:39shuffle step.
- 11:41And as we've understood that in a
- 11:43shuffle, the same keys are going to land
- 11:47in the same partition, right? Now, a
- 11:50very important parameter to define is
- 11:52the number of shuffle partition.
- 11:56Right? Is the number of shuffle
- 11:58partition. Shuffle partition is
- 11:59basically the number of partitions that
- 12:01will be produced for one single data set
- 12:04in the shuffle step, right? So, let's
- 12:07say we are assuming this to be two,
- 12:10right? Now, in order for the shuffle
- 12:11step to proceed and the same keys to
- 12:13land in the same partition, there is a
- 12:16logic that is followed. And that logic
- 12:19is basically
- 12:21that logic is basically actually let me
- 12:23erase
- 12:24some of the stuff here.
- 12:27So, that logic is basically hash
- 12:30of
- 12:31the join key
- 12:33mod shuffle partition, right? So,
- 12:36shuffle partition is two. So, this is
- 12:39always going to produce a number between
- 12:41zero and one, right? So, whatever the
- 12:44join key is, it is always going to
- 12:46produce a number between zero and one
- 12:49because shuffle partition is two, right?
- 12:52So,
- 12:53this means that the key is going to go
- 12:55to either partition number zero or it is
- 12:58going to go to partition number one.
- 13:00Okay, so, let's take an example.
- 13:02So, we are going to calculate hash of
- 13:04the join key, which is your customer
- 13:06one,
- 13:07mod two, right? So, let's assume for
- 13:11simplicity that hash of CX,
- 13:14where X is a number, right? These are
- 13:16all numbers, C1, C3, C6, and C4.
- 13:21So, we are going to assume hash of CX
- 13:23mod shuffle partitions, this is simply
- 13:27going to be X mod shuffle partition,
- 13:30right? So, in this case, what this means
- 13:32is we are going to take this X one mod
- 13:35two, and this is one. And again, this is
- 13:37just for simplicity. In reality, as I
- 13:40mentioned, this hash is going to be
- 13:43MurmurHash or some other kind of hash
- 13:46that is widely used. Yeah? So, now we
- 13:49are going to do a one mod two, which is
- 13:52one. So, this is going to go to
- 13:54partition number one, right? Similarly,
- 13:58C3 is going to be three mod two, and
- 14:01this is also going to go to partition
- 14:02one. C6 is going to be six mod two,
- 14:06which is zero. So, this is going to
- 14:08partition number zero, partition number
- 14:10zero,
- 14:11partition number zero, partition number
- 14:13one, partition number zero. And
- 14:15similarly, we are going to repeat the
- 14:17same for this one. So, this will be
- 14:18partition one, partition zero,
- 14:21partition one, partition zero, partition
- 14:24zero,
- 14:25partition one, and partition one. Right?
- 14:29So, this logic is going to decide in
- 14:32which partition are each of these rows
- 14:35going to go to, right? So, the ones that
- 14:37end with modulo one, they are going to
- 14:39go to partition number one. And the ones
- 14:43that give you a result of zero after the
- 14:45modulo, they are going to partition
- 14:47number zero. Right?
- 14:50So, in the next step,
- 14:52as we said earlier, the shuffle
- 14:55partition the shuffle partitions is two
- 14:57in number, and basically the shuffle
- 15:00partition means the output partition per
- 15:02data set after or during the shuffle
- 15:06phase, right? So, here you see for
- 15:09executor number one, for customer, this
- 15:12is the first partition, or rather, I
- 15:14should just highlight this partition
- 15:16number zero, and partition number one.
- 15:19This should be partition number one.
- 15:21Right? And all the ones all the customer
- 15:24ID
- 15:25with the modulo of 06 mod
- 15:282 is equal to 0, they have all ended up
- 15:32in partition number 0. Right?
- 15:35This is all
- 15:37have ended up in partition number 0 for
- 15:39the orders partition. Similarly, this
- 15:42has ended up in partition number one
- 15:44because 3 mod 2 is equal to 1. And
- 15:48similarly, all of this has also up in
- 15:51partition number one. Right? So, this
- 15:53complete the shuffle phase in which the
- 15:56same key the same key as in this the
- 15:59keys which have the same modulo value
- 16:02after doing the hash, right? They end up
- 16:05in the same partition. So, that
- 16:07completes the shuffle phase. Right? So,
- 16:11now once the shuffle phase is completed,
- 16:13we then go ahead to the sort phase.
- 16:16Right? And in the sort phase, it is
- 16:18simply going to take up this customer ID
- 16:20right over here and it is going to sort
- 16:24the row by that customer ID. So, let's
- 16:26go ahead to the next step over here and
- 16:29we see that this is all sorted which is
- 16:31C2, C4, C6, right? And this is sorted.
- 16:35C2, C4, and C6. And wherever there is a
- 16:39tie, it has taken up order one and order
- 16:42number seven. And similarly, C1, C3, C5
- 16:46and C1, C3, C5. Right? So, this then
- 16:50complete the sort phase over here.
- 16:54And finally, we are then going to go
- 16:57ahead with the merge
- 17:01step.
- 17:02Right?
- 17:03And the merge step is very similar to
- 17:06the merge step in merge sort.
- 17:10Right? In merge sort. So, actually, let
- 17:12me remove
- 17:14all of this quickly
- 17:16and show you how the merge is going to
- 17:18look like. So,
- 17:19both of them for both of them a pointer
- 17:21is going to be placed, right? And a
- 17:24comparison is going to happen whether
- 17:26this C2 is the same as this C2. And in
- 17:29this case, it is. Therefore, a join is
- 17:31going to happen and the resulting rows
- 17:33you're going to see over here. Right?
- 17:35These are going to be your joined rows.
- 17:39Yeah?
- 17:41Then it goes to the next row. It
- 17:42compares, is this value and this value
- 17:45the same? Which it is.
- 17:47So, a join is going to happen over here
- 17:49between this row and this row. Now it
- 17:52comes over here and it compares the if
- 17:54this value C4 same as this one, which is
- 17:57not the case, then it is going to move
- 17:59the pointer ahead. It is going to come
- 18:02to C4 and then again a comparison happen
- 18:04between this value C4 and this value.
- 18:07And now the values are the same. So,
- 18:10therefore this there is going to be a
- 18:12join between this row
- 18:13and this row. And then the value the
- 18:16pointer comes to the next row and it
- 18:18compares whether this C6 and this value
- 18:21over here is the same. It is not. So,
- 18:23then the pointer move over here and
- 18:25again the comparison happen. This time
- 18:27both of them are the same. So, then you
- 18:29see the join result over here, which is
- 18:31Adam and Mumbai, right? Similar thing is
- 18:35going to be repeated over here for this
- 18:37one and this one and the matching is
- 18:39going to happen, right? So, a join
- 18:41between these two data set is going to
- 18:44produce a new joined
- 18:48partition.
- 18:51Right? So, this is going to be the first
- 18:53partition that is going to be produced
- 18:55and both of them is going to be joined
- 18:57and this is also going to produce a new
- 19:00joined partition.
- 19:04And as a result, this is what is the
- 19:07output that you see.
- 19:09Right? So, on executor number one, this
- 19:11is your first output, which is basically
- 19:15this one right over here
- 19:17on executor number one. And for executor
- 19:20number two, this is going to produce the
- 19:22second partition, which is right over
- 19:25here, partition number one and partition
- 19:27number zero. Right? And this is the
- 19:29joined result for both of the data set.
- 19:33Right? So, this is how a sort merge join
- 19:37works. Okay, so now let's go ahead with
- 19:39the lab for the sort merge join. Right?
- 19:44And I'm going to use two data sets. So,
- 19:47I have already put
- 19:49uh
- 19:50in this default schema. So, for me,
- 19:54um when you create a workspace,
- 19:57it creates a catalog by default. And
- 20:00this is the default catalog for me.
- 20:02Right? And similarly, when you create a
- 20:04workspace, you are also going to have a
- 20:07default catalog. And inside of this,
- 20:09there is going to be a default schema.
- 20:13Right? So, I have created a volume where
- 20:18I can store all of the data. Right? A
- 20:22CSV file. Right? So, I created this
- 20:24volume called ecom, and you can also
- 20:26create it by going over here and
- 20:28clicking on volume. And this can be a
- 20:31managed volume. Right? So, I when I
- 20:34click on ecom, I see that there are two
- 20:37files under the folder ecom. Right? So,
- 20:41I created a directory which is ecom, and
- 20:44then I put these two CSVs, which we are
- 20:48going to use in all of the exercises.
- 20:50So, don't worry, all of this will be
- 20:52made available in GitHub, from where you
- 20:55can download and follow the exercises.
- 20:58Right? So, now let's get started. So,
- 21:01here I'm going to show you the physical
- 21:03plan and how a sort merge join will look
- 21:08like on the Spark UI, right? So, let's
- 21:11set up a few things first. Uh we are
- 21:14going to import PySpark uh dot SQL uh
- 21:17this functions module as F.
- 21:20And what we are doing over here
- 21:22essentially is I am setting the auto
- 21:25broadcast join threshold to minus one,
- 21:28right? So, because the two data sets
- 21:30that I have uploaded, they are small in
- 21:33size, so naturally Spark is going to
- 21:35force a broadcast hash join, right? But
- 21:39we want to see an example of a sort
- 21:42merge join, right? So, for that reason,
- 21:44I am setting this to minus one
- 21:47and the number of shuffle partitions to
- 21:50four, yeah. So, let's go ahead and run
- 21:52this. Another interesting property is
- 21:55the spark.sql.adaptive.enabled,
- 21:58which is basically turning on AQE, which
- 22:02is adaptive query execution, right? So,
- 22:05if we set this on, you are going to see
- 22:08a lot of things in the Spark UI like AQE
- 22:10shuffle read and other optimizations,
- 22:12right? So, we are turning this off just
- 22:16to have a clean and aesthetic plan when
- 22:19we look at the Spark UI and the physical
- 22:22plan, right? So, let's go ahead and run
- 22:25this.
- 22:26And
- 22:27uh we are setting the paths over here,
- 22:29which is basically uh the path that I
- 22:32have got from here. If you click on copy
- 22:34path, this is going to give you the same
- 22:36path, right? For customer.
- 22:39So, let's go ahead and run this. And
- 22:41now, we are going to read both of the
- 22:43data frame with the header set to true
- 22:46and the infer schema set to true, so
- 22:49that it is automatically able to
- 22:51understand the data types,
- 22:53right? Okay, so the two data frames have
- 22:55been read and here you see the columns
- 23:00and the data type. The customer ID is an
- 23:02integer. All the others, which is first
- 23:04name, last name, email, phone, city,
- 23:07these are all strings. The sign-up date
- 23:09is date, and whether the customer is
- 23:12active or not is a boolean.
- 23:14And similarly for orders, we have all
- 23:16the relevant
- 23:19columns and the data type, which is
- 23:21order ID and customer ID is integer. The
- 23:24order date is a date. Category, product,
- 23:27and all of this is string. Unit
- 23:30uh quantity is integer and unit price is
- 23:32double. The amount is double, and so on,
- 23:36right? So, let's also quickly have a
- 23:38look at both the data sets. And this is
- 23:40how it looks like. Customer ID, first
- 23:41name, email,
- 23:43phone number, and all of this, right?
- 23:45And similarly, the orders is where the
- 23:47order ID is over here. The customer who
- 23:50placed the order, when did they place an
- 23:52order, what exactly did they place an
- 23:54order, and the total prices and the
- 23:56amount, right?
- 23:58Now, let's go ahead. Let me skip this
- 24:01for a while. Well, actually, let's go
- 24:03ahead and run this. Let's see how many
- 24:06rows are there in each of the data sets.
- 24:08So, there are 50,000 rows in each of the
- 24:11data sets. Now, let's go ahead and do a
- 24:14join. So, I'm joining orders with the
- 24:17customer data set. We are joining it on
- 24:19the customer ID because this is what is
- 24:22common between the two data sets, right?
- 24:25And now, we are going to do an inner
- 24:28join. Yeah? So, let's go ahead and run
- 24:30this.
- 24:32Yeah, let's go ahead and do a show.
- 24:40So, this is how my data set is going to
- 24:43look like. So, you have all the
- 24:46order level details over here, right?
- 24:48The order date, and
- 24:53the customer level details over here,
- 24:54which is the first name, last name,
- 24:56email, and all of this, right? Now,
- 24:58before we go to the Spark UI, let's
- 25:00first have a look at the plan, right?
- 25:03So, this is the physical plan. So, let
- 25:05me copy the physical plan, and let's
- 25:08paste it over here.
- 25:10Yeah?
- 25:11So, let's have a look at this line by
- 25:14line, right? So, the way we read the
- 25:16physical plan is from the bottom to the
- 25:19top, right?
- 25:21So, the first line is basically a file
- 25:23scan CSV.
- 25:25Right? And this is your customer ID,
- 25:27first name, last name, all of this. So,
- 25:28this simply means that we are or Spark
- 25:32is reading the customer's data set,
- 25:35right? And let's have a look at a few
- 25:37other things, right? An interesting
- 25:40thing to note over here is this
- 25:41in-memory file index, and this is
- 25:44basically the path that we have
- 25:46specified, right? So, in-memory file
- 25:50index is basically Spark file listing
- 25:53component, which finds out the file that
- 25:56is present at a given location, right?
- 25:59So, it basically finds out the metadata,
- 26:01sizes, partition values, and all of it,
- 26:04right? So, it is going to list all of
- 26:06the files specified at that location,
- 26:09and it stores it in the driver's memory.
- 26:13So, here we have specified one single
- 26:16file. If we specified a directory, the
- 26:19number is going to be a lot larger,
- 26:21right? So, that is one. And this is the
- 26:25schema of the data set that has been
- 26:27read in.
- 26:29Now, if we go to the second line,
- 26:31you see that filter not null,
- 26:34and this is basically filtering on the
- 26:36customer ID. So, it is basically making
- 26:39sure that the customer IDs that have
- 26:41been read in
- 26:42from the data set is not null. And a
- 26:45very important point to remember is that
- 26:47the customer ID is your join key. Now,
- 26:51we did not add this condition anywhere,
- 26:53right? Did we add this condition
- 26:55anywhere? No.
- 26:56Right? We haven't added this condition
- 26:58anywhere. So, this is a part of Spark
- 27:01internal optimization strategy. Yeah?
- 27:05To make sure that the data is clean.
- 27:07The third step is exchange hash
- 27:10partitioning on the customer ID. And the
- 27:13number four that you see over here,
- 27:15yeah? The number four that you see over
- 27:17here is basically the shuffle partition.
- 27:20So, we certify the shuffle partition as
- 27:23four over here, right? So, basically
- 27:25before going ahead and doing the
- 27:27shuffle, exchange over here means
- 27:29shuffle.
- 27:31It basically
- 27:33makes sure that it has specified that
- 27:36the number of shuffle partitions are
- 27:38four. And what does hash partitioning
- 27:40mean? So, if you remember, we we learned
- 27:43that during a shuffle, the same keys go
- 27:46to the same partition. And how does that
- 27:48happen? It happens by taking up the join
- 27:51key, doing a hash, and then a modulo of
- 27:56the part of the shuffle partition,
- 27:58right? So, basically the hash of the
- 28:01join key,
- 28:03and then modulo the shuffle partition.
- 28:07Right? And in this case, the shuffle
- 28:09partition is four. So, this is simply
- 28:12going to be
- 28:13this one
- 28:14mod four, right?
- 28:17And this is what is hash partitioning
- 28:19over here. Right? So, this step
- 28:21basically distributes all of the data,
- 28:24making sure that the same keys fall in
- 28:27the same partition, right? So, this is
- 28:30the shuffle step of the sort merge join.
- 28:33And then, as we've seen,
- 28:35each of the rows is sorted within the
- 28:38partition, right? And it is sorted by
- 28:41the customer ID. And the logic that is
- 28:43used is ascending and nulls first, yeah?
- 28:46So, the sort is complete, the shuffle is
- 28:50complete, and then we repeat the same
- 28:53for the orders data set. Right? So, here
- 28:57you see a file scan, and then the order
- 29:00customer ID, order date, and all of
- 29:02that. And here you would see a similar
- 29:05in-memory file index, which means that
- 29:08it is reading Basically, it is not
- 29:11reading actually, it is listing the
- 29:13files that have been specified in the
- 29:16path that we have given. Right? And here
- 29:19it has found only one file. Yeah? And it
- 29:22is doing something very similar, which
- 29:24is filter not null, making sure that the
- 29:26customer ID is clean. And it is then
- 29:29doing a shuffle.
- 29:32Yeah? So, exchange specifies a shuffle,
- 29:35and then it does a hash partitioning.
- 29:37So, hash partitioning is the logic that
- 29:39it uses in order to do the shuffle. And
- 29:42you see four shuffle partitions over
- 29:44here. Yeah? And then the sort step
- 29:47again. And finally, when both of the
- 29:49data frames, right? When both of the
- 29:51data frames have been shuffled, they
- 29:53have been sorted, the next step is the
- 29:55merge step. Right? And this is where
- 29:58your final sort merge join happens, and
- 30:01it has specified that this is an inner
- 30:05join. Yeah? So, once the join has
- 30:07happened, project basically picks up the
- 30:09relevant columns that the user needs to
- 30:13see, or the ones that need to be
- 30:15displayed. Right? So, this is how the
- 30:17physical plan for a sort merge join is
- 30:20going to look like. Now, let's have a
- 30:22look at how the Spark UI is going to
- 30:24look like. So, here is where
- 30:27the join happened. Right?
- 30:31So, this is where the action was called,
- 30:33and that is why it triggered a Spark
- 30:35job. Now, let's have a look at
- 30:39the Spark UI, and the job number is
- 30:42eight. So, let me directly go to the SQL
- 30:47{slash} data frame and here you see that
- 30:49this is job number eight. Let's go ahead
- 30:51over here and you are going to see a
- 30:54DAG, right? So, this DAG basically scans
- 30:59the scans the orders data set and you
- 31:02would see that this basically reads up
- 31:0550,000 rows, yeah?
- 31:07It reads up 50,000 rows and similarly
- 31:10this is the customer data set, yeah? So,
- 31:13you see customer ID, first name,
- 31:15uh last name, email and all of that. And
- 31:19the rows output is 10,000,
- 31:23yeah?
- 31:24Now,
- 31:25what we are going to do is
- 31:28it is going to
- 31:30do a filter of not null on both the data
- 31:34sets, right? Something similar to what
- 31:36we just saw in the Spark physical plan.
- 31:39And then it is going to do an exchange
- 31:41hash partitioning
- 31:43of
- 31:44both the data sets, right? On the left
- 31:47and on the right,
- 31:49yeah?
- 31:50On the left and on the right and then it
- 31:52does a sort of both the data frames and
- 31:55finally it does a merge, right? Which is
- 31:58specified as a sort merge join over here
- 32:01and the key that is used, the join key
- 32:04is basically the customer ID and it
- 32:06specifies an inner join and finally
- 32:08there are 10 columns which is selected
- 32:10as a result of the final data set being
- 32:13produced. So, this is how overall
- 32:16all sort merge joins are going to look
- 32:18like, right? All sort merge join, the
- 32:20DAG for all the sort merge join is going
- 32:23to look something very similar like
- 32:26this. Okay, so the next type of join we
- 32:28are going to study is broadcast hash
- 32:31join, right? And let's first understand
- 32:34when exactly is this chosen by Spark,
- 32:38right? So, broadcast hash join is chosen
- 32:41or it is performed when, let's say, you
- 32:43have two tables and one of the tables in
- 32:47the join is small enough to fit entirely
- 32:50in memory in each of the executor,
- 32:52right? So, the key point to note is one
- 32:54of the tables or one of the data frame
- 32:56is small enough to fit entirely in the
- 33:01executor's memory on each of the
- 33:03executor's memory, right? So, the Spark
- 33:06based Spark cost based optimizer, right?
- 33:09So, the Spark
- 33:12CBO, which is cost based optimizer,
- 33:15automatically chooses broadcast hash
- 33:17join when the estimated size of one of
- 33:21the table is below the threshold defined
- 33:24by this parameter.
- 33:26And we've seen this parameter earlier,
- 33:27which is
- 33:27spark.sql.autoBroadcastJoinThreshold.
- 33:36The default value of which is 10 MB,
- 33:39right? So, the Spark cost based
- 33:41optimizer automatically is going to
- 33:43automatically choose this when the
- 33:46estimated size of one of the table is
- 33:49below this threshold, right? So, that is
- 33:53the case when
- 33:55broadcast hash join is going to be
- 33:57chosen, right? So, to quickly summarize,
- 33:59you have two tables, one table is large,
- 34:01the other table is small, and the other
- 34:04table, the size of the other table is
- 34:06below the threshold, which is defined by
- 34:10this property over here. In this case,
- 34:13Spark is going to choose broadcast hash
- 34:16join, yeah? So, there are three steps to
- 34:19performing a broadcast hash join. The
- 34:21first one is, of course, broadcast,
- 34:24right? And we're going to understand
- 34:26what exactly each of these steps are,
- 34:28but the three steps quickly are
- 34:30broadcast, build, and probe. Build is
- 34:33building the hash map, and and probe is
- 34:35actually finding for finding the right
- 34:37key, finding the match in a hash map,
- 34:41right? So, let's understand this
- 34:42visually with an example all of these
- 34:45three steps. Okay, so this is our setup.
- 34:48We have a customer table, and we have
- 34:51the orders table, right?
- 34:53There are two executors. Each of the
- 34:55executor have one order partition each,
- 34:59right? And right now
- 35:02the customer table resides on the
- 35:05driver. It is very small in size. As you
- 35:07see, we have listed that this is 320
- 35:10bytes in size, and the default value of
- 35:14auto broadcast join threshold is 10
- 35:17megabytes, right? So, what is going to
- 35:19happen is that the first step, which is
- 35:22the broadcast
- 35:25step,
- 35:27is going to take place, right? So, what
- 35:29is the broadcast step? The broadcast
- 35:31step is the step in which the driver is
- 35:35going to serialize the small table and
- 35:38push a full copy to every executor,
- 35:41right? So, here you see we have two
- 35:42executors. It is going to send two
- 35:45copies of this small customer table to
- 35:50these two executors over here, right?
- 35:53So, what the word serialize is mean
- 35:55which I I said that the driver is going
- 35:57to serialize the small table and then
- 35:59push a full copy. So, what that means is
- 36:02that this in-memory table, this table is
- 36:05the residing in memory, right?
- 36:06Everything resides in memory in Spark
- 36:08during computation. So, in-memory table
- 36:12is going to be converted into a flat
- 36:15sequence of bytes that can travel over
- 36:19the network, right? So, this is going to
- 36:22be serialized, and this table is going
- 36:24to be broadcasted
- 36:29on to executor number one
- 36:33on to executor number one and this is
- 36:35also going to be broadcasted on to
- 36:38executor number two. So, you will have
- 36:40one full copy
- 36:42You will have one full copy of the
- 36:44customer's
- 36:47partition on executor number two
- 36:50and on executor number one.
- 36:55Right? So, that is the first step which
- 36:57is the broadcast step. Now, coming over
- 37:00to the second step which is the build
- 37:03step. Right? The build step.
- 37:08Now, in the build step each executor is
- 37:11going to iterate over the received copy.
- 37:14So, it just now received the copy of the
- 37:16customer's partition. Yeah? And it is
- 37:19going to iterate over each of the rows
- 37:21and it is going to build a hash map and
- 37:24that is what you see over here.
- 37:27Right? So, this basically has four rows
- 37:30which is C1, C2, C3, C4 of Fark, Naved,
- 37:33Rohan, and Imdad.
- 37:36So, basically it is going to convert
- 37:39this into a hash map where
- 37:44it is going to be in the format of
- 37:46customer
- 37:49ID and the row value. Right? So, that is
- 37:53what you see over here. You have the
- 37:55customer IDs and the row value. Yeah?
- 38:01And this is going to take place at both
- 38:04of the executors.
- 38:06Right?
- 38:07So, that is the build phase and finally
- 38:10we come over to the probe phase.
- 38:15Right? We come over to the probe phase
- 38:17and in this step what happens is each
- 38:19executor is going to scan its partition
- 38:22of the large table row by row. So, the
- 38:25large table that you see over here, it
- 38:27is going to be scanned row by row and we
- 38:30are going to look if there is a match
- 38:33for C1 in this hash map.
- 38:38Right? So, it is basically going to
- 38:40check whether C1 exist in this hash map
- 38:43or not. Right? And this is simply an
- 38:47order of one operation, a very quick
- 38:49order of one operation. So, if there is
- 38:51a hit, it produces a joined row. So, C1
- 38:55for C1, it is a hit because you see this
- 38:59row over here. For C3,
- 39:01it is a hit and for C2, it is a hit.
- 39:04Right? And similarly for C4, you find C4
- 39:07over here, C1 over here and C3 over
- 39:11here. Right? So, that means all of these
- 39:14rows have found a hit and they are going
- 39:19to be joined.
- 39:21And the value is going to be produced
- 39:22over here at the joint data set. Right?
- 39:25So, for C1, you're going to have the
- 39:27name
- 39:29which is a park.
- 39:32And the city equal Bengaluru.
- 39:35Right? Similarly, this is going to
- 39:37happen for all. This is going to happen
- 39:38for all and you see the resulting data
- 39:42set over here.
- 39:44This is the resulting data set, the two
- 39:46partitions.
- 39:48And the joined value on the right-hand
- 39:50side. So, this is how a broadcast join
- 39:54broadcast hash join is going to work.
- 39:57Now, let's go ahead with the lab for
- 40:00the broadcast hash join. Right? And we
- 40:02are going to follow something very
- 40:04similar what we've done for the sort
- 40:06merge join. So, let's go ahead and set
- 40:08all of this up. Let's quickly check what
- 40:12is the auto broadcast join threshold,
- 40:14which is set to 10 megabytes. right?
- 40:18This is the default value. And as we
- 40:20discussed earlier, we are going to set
- 40:23AQE to false.
- 40:25And
- 40:26basically setting up the paths for both
- 40:28of these data sets.
- 40:30Defining both of these data sets, right?
- 40:34Setting the header to true and info
- 40:35schema to false.
- 40:37Sorry, true, not false. Yeah.
- 40:39And let's quickly have a look at both of
- 40:42these data sets. So, these data sets are
- 40:44just the same as what we've used for the
- 40:47sort merge join, right?
- 40:51So, now let's go ahead and
- 40:53perform
- 40:55the join, right? So, you see that the
- 40:56code is literally the same. DF orders
- 41:00uh joins with DF customers on customer
- 41:03ID and the join type is inner, yeah.
- 41:06So, one key difference is last time we
- 41:09specified we were forcing a sort merge
- 41:13join by specifying this to be minus one,
- 41:15but this time we have allowed it to have
- 41:19its default value, right?
- 41:22So, in this case, this should ideally
- 41:25trigger a broadcast hash join.
- 41:28So, let's go ahead and run this. Let's
- 41:30do a show. The results are going to be
- 41:33the same as the sort merge join, but
- 41:36what is important is let's look at
- 41:39the physical plan, yeah. So, let's copy
- 41:42this and let me put this over here.
- 41:45Which is broadcast hash join.
- 41:49Let's put the physical plan over here.
- 41:53And let's read this from bottom to top.
- 41:55So, it followed the normal drill, which
- 41:57is file scan CSV. It reads the customer
- 42:00data set, and then it basically figures
- 42:02out
- 42:03that it ensure that none of the customer
- 42:06ID should be null. So, this is a part of
- 42:09Spark's own optimization strategy. Now,
- 42:11the interesting part is this broadcast
- 42:14exchange and the hash relation broadcast
- 42:18mode, right? So, this broadcast exchange
- 42:20basically says
- 42:22that the driver on which the customer
- 42:25data set is residing, right? It is
- 42:28basically going to take that, serialize
- 42:30it, and send it over the network to the
- 42:34executors, right?
- 42:36So, that the executors can have their
- 42:38own copy of the small data set. Now,
- 42:42what does hash relation broadcast mode
- 42:44mean, right? So, the broadcast mode,
- 42:48yeah, this broadcast mode
- 42:50that we see over here, this broadcast
- 42:52mode
- 42:53it basically tells how the broadcast
- 42:56side is packaged before it is shipped to
- 42:58the executor, right? So, there are two
- 43:00types of broadcast mode. One is the hash
- 43:02relation broadcast mode, the other one
- 43:05is the identity broadcast mode. And we
- 43:07are going to look into identity
- 43:09broadcast mode a little later.
- 43:11The broadcast mode basically decides or
- 43:13it basically tells you how the broadcast
- 43:16side is packaged and sent to the
- 43:20executor. So, in case of a hash relation
- 43:23broadcast mode, which is used by a
- 43:25broadcast hash join,
- 43:26the broadcast rows are basically
- 43:28transformed into a hash map, right? It
- 43:31is transformed into hash map with the
- 43:33key and a value. In our case, the key is
- 43:37going to be the customer ID and the
- 43:39value is going to be the the row value,
- 43:42right? So, how is the key decided and
- 43:45that is what you see over here, right?
- 43:47So, input of zero is basically pointing
- 43:50to index zero, which is the customer ID.
- 43:52It says that customer ID is of data type
- 43:55integer and false means that it is
- 43:58non-nullable. It converts it into a long
- 44:02or big int for optimization purposes,
- 44:05right? And all of this is converted into
- 44:07a list. Right? So, it is a list where it
- 44:10has the join key. Now, the reason why
- 44:13it's a list because sometimes join key
- 44:14can be
- 44:16more than one. Right? It can be a
- 44:17composite join key. So, that is why it's
- 44:19put into a list. So, basically it
- 44:21specifies that a broadcast exchange is
- 44:24going to happen. The data on the driver
- 44:27is going to be broadcasted over to the
- 44:29executor. And the mode that is going to
- 44:32be followed is a hash relation. Where
- 44:35the key is going to look like this.
- 44:38Right? It is a list with the customer
- 44:41ID. Yeah?
- 44:43So, that is what means from for these
- 44:47three lines. And then
- 44:49it reads the order data set.
- 44:52Filters out all of the customer ID
- 44:54wherever it is null. And finally, it
- 44:57does a broadcast hash join on the
- 45:01customer ID. And it also tells you that
- 45:02this is an inner join. And build right
- 45:05basically means which side of the join
- 45:07is going to be used as a hash table. And
- 45:11in this case, this is the right side.
- 45:13And this is the left side. The top one
- 45:16is the left side. Yeah? So, the right
- 45:19side is the customer data set which is
- 45:21going to be converted into a hash map.
- 45:24So, by this step over here the join has
- 45:27completed. And finally
- 45:29the columns that need to be selected as
- 45:32output are projected.
- 45:34So, also have a look at the Spark UI.
- 45:38Which is
- 45:41An action was called over here. So, job
- 45:4315 and 16 are the relevant ones. So, if
- 45:46I go to jobs over here over here, these
- 45:49are the two jobs. Yeah? So, let me go to
- 45:53SQL {slash} data frame.
- 45:55And this is the one that I should look
- 45:57at. Yeah?
- 45:58So, here
- 46:01and let me open this up in another tab.
- 46:06So here, what we see is
- 46:09this data set, which is the customer
- 46:11data set, has been scanned,
- 46:13and then a filter has been applied over
- 46:16the customer ID, which is what we saw in
- 46:19the physical plan as well. And then you
- 46:21see a broadcast exchange, right? And
- 46:25this goes to the same place on the
- 46:27left-hand side, where you have the
- 46:29orders data set. On the left-hand side,
- 46:32you have the order data set that was
- 46:33scanned. It also went through a filter,
- 46:36and finally both of them are broadcast
- 46:40hash join, right? So you also see
- 46:44there is time to read broadcast, right?
- 46:46Which was basically the broadcast
- 46:48exchange data set that comes over here,
- 46:51yeah? And the reason why you see rows
- 46:53output over here is six, because I
- 46:55specified a show of five, yeah?
- 46:59So finally, a broadcast hash join
- 47:00happens in this step, and you see the
- 47:04inner build right, just what we have
- 47:06understood right now.
- 47:08And finally, we select seven columns as
- 47:10a result of the final data set. So this
- 47:14is how a DAG is going to look like for a
- 47:18broadcast hash join. Almost all
- 47:22data sets which undergo a broadcast hash
- 47:24join, they are going to look somewhat
- 47:26similar. Okay, so the next join we are
- 47:28going to study is shuffle hash join,
- 47:32right? And shuffle hash join lies in
- 47:35between a broadcast hash join and a sort
- 47:38merge join.
- 47:40So this is an interesting one. It lies
- 47:42between a broadcast hash join
- 47:46and a sort merge join.
- 47:51Right? And why do I say that?
- 47:53I say that because it gets picked up
- 47:55when the situation is too big
- 47:58for a broadcast hash join, but it also
- 48:01at the same time doesn't need the full
- 48:02machinery of the sort merge join, right?
- 48:06So, when the situation is too big for a
- 48:09broadcast hash join, broadcast hash join
- 48:12will not work in that case, but you also
- 48:14don't need the full machinery of a sort
- 48:18merge join. So, in that case, Spark is
- 48:21going to choose a shuffle hash join.
- 48:23Yeah. But, what are the internal
- 48:25mechanics? What are What is the logic
- 48:27using which it is going to decide that,
- 48:29right? So, the smaller table is too
- 48:31large to broadcast to every executor.
- 48:34Yeah. So, the smaller table basically
- 48:37exceeds that property that we've been
- 48:40discussing, which is Spark
- 48:48Yeah. So, the smaller table is too large
- 48:50to be able to broadcast to each
- 48:52executor. So, the smaller table exceeds
- 48:55this default value, which is 10 MB.
- 49:00Yeah. And
- 49:02one partition
- 49:04one partition of that smaller table
- 49:06after shuffling the the redistribution
- 49:08of data can still fit into a single
- 49:12executor memory in order to build the
- 49:15hash table, right? So, the key point to
- 49:17note is the smaller table is too large
- 49:20to be able to fit into every executor,
- 49:22but however one partition of the smaller
- 49:26table after So, shuffling is going to
- 49:28redistribute all of the data to the
- 49:30executor.
- 49:32Now, after redistribution, the
- 49:34partitions of the smaller table that
- 49:36land in each of the executors
- 49:39they can still fit in the executor
- 49:41memory in order to build the hash table,
- 49:43right? So, that is the case when a
- 49:47shuffle hash join is going to be chosen.
- 49:49Now, a very important thing to note is
- 49:52there is also a preference gate. Yeah.
- 49:54Uh by preference gate, what I mean is uh
- 49:57by default, Spark is going to prefer a
- 50:00sort merge join.
- 50:02It's going to prefer a sort merge join
- 50:04over a shuffle hash join.
- 50:08Yeah. Even when shuffle hash join is
- 50:10value uh viable. Yeah. Even when shuffle
- 50:13hash join is viable, it is going to
- 50:15prefer a sort merge join. And the reason
- 50:17for this is sort merge join is more
- 50:21memory stable.
- 50:25Right? It is more memory stable. And the
- 50:27reason why I say it's more memory stable
- 50:29is because shuffle hash join has to keep
- 50:32the hash map in memory. And if there is
- 50:35a key which is skewed, then there is a
- 50:38good chance of an out of memory error.
- 50:41Right? So, we just said that in a
- 50:43shuffle hash join,
- 50:45>> [snorts]
- 50:46>> after shuffling, when shuffling does the
- 50:48redistribution of data,
- 50:50one of the partition, the partition that
- 50:52gets redistributed to each of the
- 50:54executor, that is small enough to build
- 50:57the hash table.
- 50:58But, what if when you are building the
- 51:00hash table, the key there is a key which
- 51:03which is skewed, right? So, in that
- 51:05case, there are good chances of an out
- 51:09of memory error. Yeah. So, if you want
- 51:11to override Spark default setting, you
- 51:13have to set this property, which is
- 51:16spark
- 51:17.sql.
- 51:20join
- 51:21.prefer
- 51:24prefer sort merge join.
- 51:27What are we going to do?
- 51:28Right? So, once you do that, Spark is
- 51:31going to pick up shuffle hash join when
- 51:34the smaller side per partition data fits
- 51:37in memory.
- 51:38So, the shuffle hash join has three
- 51:40steps. The first one is the shuffle step
- 51:44where both the tables are reshuffled
- 51:47across executor simply by following the
- 51:49logic which is hash of the join key
- 51:52modulo shuffle partitions.
- 51:54Yeah.
- 51:55So, this is basically to make sure that
- 51:57all the rows with the same join key land
- 52:00in the same partition in the same
- 52:01executor, right? The second one is the
- 52:04build step where each of the executors
- 52:07takes its partition of the smaller side,
- 52:09right? We discussed the smaller side and
- 52:12then it builds a hash map in memory
- 52:15keyed on the join column, right? So,
- 52:18this is going to be join key and then
- 52:20the row value.
- 52:22And then finally the probe phase where
- 52:25we're going to take the larger partition
- 52:28row by row and then we're going to look
- 52:30up in the hash table and see if there is
- 52:32a hit, right? If there is a match and a
- 52:34match is going to produce a joined row.
- 52:38So, these are the three steps. Now,
- 52:39let's look into it with visuals and
- 52:42example. Okay, so these are our data
- 52:46sets, right? We again have a customer
- 52:48data set
- 52:50and an order data set. We have two
- 52:52executors,
- 52:53right? One partition each
- 52:56for customer and orders in each of the
- 52:59executors, right? Now, the first step
- 53:02that we mentioned was the shuffle step.
- 53:07Right?
- 53:09The shuffle step and
- 53:11we're going to assume
- 53:13three shuffle partition.
- 53:15Right? This time we're going to assume
- 53:17three shuffle partition. This means that
- 53:19in the shuffle stage it is going to
- 53:22produce three partitions per data set,
- 53:24right? So, in order to understand where
- 53:27this row is going to go to, we simply
- 53:29have to follow the same logic which is
- 53:31hash of the join key
- 53:35modulo shuffle partition. And our join
- 53:37key is basically the customer ID, right?
- 53:40So, let's use the same logic, which is
- 53:42hash of customer one modulo three. For
- 53:46simplicity, we said we are just going to
- 53:48do one mod three.
- 53:51Right?
- 53:52C1 is going to become one, and one mod
- 53:54three, this is just one. Right? So, this
- 53:57is going to be one. This is C2 mod
- 54:02three is going to be two mod three,
- 54:06which is equal to two. Right? So, this
- 54:08is going to go to partition number two,
- 54:11partition number zero, partition number
- 54:13one, partition number one, partition
- 54:15number zero, partition number zero,
- 54:17partition number one, partition number
- 54:19two, partition number zero. Right?
- 54:22Similarly, this is also going to go to
- 54:24partition number two,
- 54:27and partition number zero. This is going
- 54:29to go to partition one, partition two,
- 54:31partition zero, partition one,
- 54:34partition two, partition two,
- 54:36partition zero, and partition number
- 54:39two. Right? So, now, as a result, in the
- 54:43next step, what you're going to see is
- 54:45three partitions.
- 54:48Three partitions, because we defined
- 54:51three shuffle partition.
- 54:54So, this is partition number zero,
- 54:56partition number one, they reside on
- 54:58executor number one, and partition
- 55:01number two, this is residing on executor
- 55:05number two.
- 55:06Now, why two partitions on one executor
- 55:08and one partition on the other?
- 55:10And this is something this the Spark
- 55:12scheduler decides. Right? So,
- 55:15this one where all of this is modulo
- 55:17zero, goes to partition number zero.
- 55:23Wherever it was modulo one, it goes to
- 55:25partition number one, and similarly,
- 55:27wherever it was modulo two, it goes to
- 55:30partition number two. And similarly for
- 55:32the orders as well. This is all modulo
- 55:34zero,
- 55:36this is all modulo one and this is all
- 55:39modulo two.
- 55:40Right? So, this is post
- 55:43shuffle
- 55:46state.
- 55:47Right?
- 55:49This is how the partitions are going to
- 55:51look like. Now, the next step the next
- 55:54step was the build step. Right? So, in
- 55:58the build step, what we are going to do
- 56:00is we take the smaller side
- 56:03we take the smaller side and we build a
- 56:06hash map out of it. Right? So, we take
- 56:09this side
- 56:10we take this partition and we build a
- 56:13hash map. And the hash map is simply
- 56:15going to be of the format the join key
- 56:19the join key to the row.
- 56:22Right? So, here the join key is C3 and
- 56:25C6
- 56:26which is the customer ID. Customer ID is
- 56:28the join key and this is the row value.
- 56:31Right?
- 56:32And similarly, it happens for this one
- 56:33as well. The large partition stays as it
- 56:36is during this step. So, in this step,
- 56:38we have now created a hash map from the
- 56:43smaller partition. Right?
- 56:45Now, we are going to do the probing.
- 56:48And in probing, what is going to happen
- 56:50is each executor is going to stream its
- 56:53partition of the larger side
- 56:56which is this one, these three
- 56:58partitions, right? Row by row and it is
- 57:01going to look up each of the rows keys
- 57:04in the local hash map. So, C3 is going
- 57:07to check whether it exists over here or
- 57:09not. And it did find a
- 57:13hit. Therefore, it is going to match and
- 57:15then you're going to get name
- 57:18and then the city over here which is
- 57:21Rohan
- 57:22and Bengaluru. Right? Similarly, C3 is
- 57:25going to match match 60 C6 is also going
- 57:28to match over here.
- 57:29Right? Similarly, this is also going to
- 57:31happen over here. It is going to probe
- 57:32this hash map
- 57:35and find a hit. And similarly for this
- 57:37one as well.
- 57:40Okay, I think this has been
- 57:43there's a mistake over here. This should
- 57:46be C2 and this should be C5 and this
- 57:49should be Okay, this These two names are
- 57:51fine. And C2 is going to match over here
- 57:54and C5 is going to match with Imdad over
- 57:57here, right? So, this is how the matches
- 58:00the join is going to produce hits rows,
- 58:03right? And finally, as a result, you see
- 58:06that there are three partition, right?
- 58:08As a result of joining these two,
- 58:12this one and this one.
- 58:15Joining these two, this one and this
- 58:17one. So, there are going to be two
- 58:18output partitions
- 58:20on
- 58:21the first executor.
- 58:27And one output partition on the second
- 58:31executor. And that is what you see over
- 58:33here.
- 58:34Number one, number two, and
- 58:37number three.
- 58:39Yeah. So, this is how a shuffle hash
- 58:43join works.
- 58:44Okay, so let's go ahead with the lab for
- 58:47the shuffle hash join. And you'll find a
- 58:49lot of things to be similar. We do the
- 58:51import. We set
- 58:53spark.sql.adaptive.enabled
- 58:56basically the AQE to false and shuffle
- 58:59partitions to four. And we've read about
- 59:03this very important property, which is
- 59:05spark.sql.join.preferSortMergeJoin
- 59:09to be false. So, this is the property
- 59:12which enables Spark to choose shuffle
- 59:15hash join whenever it is viable, right?
- 59:19And the second important property is we
- 59:21have lowered the threshold for broadcast
- 59:24join from 10 MB to 1 MB, right? So, the
- 59:27default is 10 MB. We have lowered it to
- 59:311 MB. If we don't do this, our data set
- 59:34is less than 10 MB, but it is greater
- 59:36than 1 MB, right? So, it is
- 59:39automatically going to trigger a
- 59:40broadcast hash join, but we don't want
- 59:43that, right? So, if you remember the
- 59:45original definition,
- 59:46the situation or the condition should be
- 59:48such that
- 59:49it should be bigger, it should be uh
- 59:53higher for a broadcast hash join, but it
- 59:56also doesn't need the full machinery of
- 59:58a sort merge join. And that is why a
- 1:00:01shuffle hash join falls in between.
- 1:00:03Yeah? So, that is why we have lowered
- 1:00:05the threshold so Spark doesn't
- 1:00:08automatically do a broadcast hash join.
- 1:00:11Yeah? But also, when the shuffle
- 1:00:13happens, the data set is broken into
- 1:00:15small small partition, and those small
- 1:00:17partition, it is able to build a hash
- 1:00:20table on those small partition, right?
- 1:00:23So, I hope this makes sense. Now, we
- 1:00:25this is the normal drill that we've been
- 1:00:28repeating for quite some time now. Uh
- 1:00:31we basically run this, and we set the
- 1:00:34paths.
- 1:00:35We read both the data frame. We have a
- 1:00:38look at both of them, and then finally,
- 1:00:40we do a join between the two data sets.
- 1:00:44Yeah?
- 1:00:46Now, before running show, let me just go
- 1:00:48ahead and print the plan.
- 1:00:51Yeah? Let me copy this.
- 1:00:54Let's put
- 1:00:57a shuffle hash join title over here.
- 1:01:01Yeah. So, let's go ahead and read this
- 1:01:04from the bottom to the top.
- 1:01:06So, we first read the file, which is the
- 1:01:08customer data set.
- 1:01:10We simply filter out all of the null for
- 1:01:13customer ID, and then a normal shuffle
- 1:01:17happen, right? A normal shuffle happen.
- 1:01:20Similarly, for the second data set, the
- 1:01:22order data set, it is read, it is
- 1:01:24scanned, and then a normal shuffle
- 1:01:28happen, right? So, the three steps for
- 1:01:31both of them, they look identical. And
- 1:01:34this is normal, right? The hash is built
- 1:01:37after the shuffling, and that is what
- 1:01:39you see over here.
- 1:01:41Once both of them have been shuffled, a
- 1:01:43shuffle hash join takes place, and you
- 1:01:46see this is an inner join. The keyword
- 1:01:49to note over here is build right. So,
- 1:01:52build right basically specifies the
- 1:01:55table on which the hash map is to be
- 1:01:58built, right? And this is the right
- 1:02:00table, which is the customer data set.
- 1:02:02This is the left table. Yeah? So, uh
- 1:02:06hash table is built on the customer data
- 1:02:09set, and then a shuffle hash join is
- 1:02:12performed. And finally, the data uh a
- 1:02:15project operation is called, which
- 1:02:16selects the relevant columns. Yeah?
- 1:02:19Okay, let's go ahead and also have a
- 1:02:21look at the DAG. And for this, we are
- 1:02:24going to run this action over here, and
- 1:02:27this produces one Spark job, which is
- 1:02:29job number 40. Yeah? So, let's quickly
- 1:02:32open this up. If I look at job, this is
- 1:02:35basically job number 40, and it has
- 1:02:39three stages, right? It has three stages
- 1:02:42that you would find over here.
- 1:02:44But, let me go to the SQL data frame.
- 1:02:48Yeah, let's click over here.
- 1:02:51And actually, the structure looks very
- 1:02:53similar to the sort merge join. Yeah?
- 1:02:56So, so on the left-hand side, you see
- 1:02:59scanning both of the data sets, right?
- 1:03:02The this one on the right is the
- 1:03:04customer ID, and then on the left is the
- 1:03:06order ID, order order data set, right?
- 1:03:10Um then, you are going to filter on the
- 1:03:13customer ID not equal to null, and
- 1:03:15similarly for this one, something that
- 1:03:16we've seen in the plan already.
- 1:03:19And now there is going to be an exchange
- 1:03:21hash partitioning that we see through
- 1:03:23this step, the left and the right. And
- 1:03:26finally, there is going to be a shuffle
- 1:03:28hash join. But where do we actually see
- 1:03:31that the hash map is being built? So
- 1:03:33that is what you see over here, time to
- 1:03:36build hash map, which is 6 ms, right? So
- 1:03:40this is how
- 1:03:42it shows that a hash map was built and
- 1:03:44then it was probed to find the matches,
- 1:03:48right? And that is what is also
- 1:03:50specified by build right. And on the
- 1:03:52right side, you have the customer data
- 1:03:55set. And finally, the shuffle hash join
- 1:03:58is performed and it projects eight
- 1:04:00columns, which is the final output of
- 1:04:03the data set. This is how a DAG will
- 1:04:06look like for a shuffle hash join. Okay,
- 1:04:09so the final type of join that we're
- 1:04:12going to study is the broadcast nested
- 1:04:15loop join, right? So let's first
- 1:04:18understand when is it performed, when is
- 1:04:20it chosen by Spark
- 1:04:23over any other join, right?
- 1:04:25So the first point is a non-equi join
- 1:04:29condition, right? So when the join
- 1:04:31predicate contains no equality clause.
- 1:04:34So when you're doing a join, let's say
- 1:04:35you are doing DF1.x
- 1:04:38= DF
- 1:04:412.x.
- 1:04:43Instead of equal to, you have a non-equi
- 1:04:46join predicate, right? And what that
- 1:04:48means is you have something like equal
- 1:04:51to, greater than equal to, less than
- 1:04:53equal to, less than,
- 1:04:55not equal to, or something like a range,
- 1:04:59right? So for example, if you use
- 1:05:01between,
- 1:05:03right? So DF1.x
- 1:05:05between
- 1:05:07some value and DF2.x, right? Something
- 1:05:10like that. So basically, it has to be a
- 1:05:12non-equi join condition. It is either an
- 1:05:16inequality, a range, or a complex
- 1:05:19condition, right? So, in that case, a
- 1:05:22broadcast nested loop join will be
- 1:05:25chosen. Another case where a broadcast
- 1:05:28nested loop join is chosen is during a
- 1:05:32cross join, when there is no join
- 1:05:34condition at all, right? So, every row
- 1:05:36from one table must be paired with every
- 1:05:39row from the other table, right? So,
- 1:05:42these are the two conditions when Spark
- 1:05:45is going to choose a broadcast nested
- 1:05:48loop join.
- 1:05:49So, now, what exactly are the steps? The
- 1:05:52steps are twofold. The first one is the
- 1:05:54broadcast step. And this broadcast is
- 1:05:58identical to the broadcast hash join, in
- 1:06:00which the smaller table is serialized by
- 1:06:04the driver, and it is pushed in full to
- 1:06:07every executor, right? No shuffle of the
- 1:06:10large table.
- 1:06:11And this also brings me to one more
- 1:06:13condition, right? If this broadcast
- 1:06:16needs to happen, then
- 1:06:19if we have two tables,
- 1:06:21if we have two tables, then one of the
- 1:06:24table needs to be small to be able to be
- 1:06:26broadcasted, and the other table can be
- 1:06:30large, right? So, along with these two
- 1:06:32conditions, this is also a very
- 1:06:36important condition. Otherwise, the
- 1:06:38broadcast step cannot happen, right? So,
- 1:06:41that is why the first step is the
- 1:06:43broadcast step, where you take the
- 1:06:44smaller table, it's serialized by the
- 1:06:46driver, and it is pushed in full to
- 1:06:49every executor. Yeah? And there is no
- 1:06:52shuffle of the large table.
- 1:06:54And step number two,
- 1:06:57step number two is the nested loop,
- 1:06:59right? So, every executor runs a double
- 1:07:02loop over its partition of the large
- 1:07:04table, and the full broadcast table, and
- 1:07:07then it evaluates the join condition.
- 1:07:10So, it looks something like this. So,
- 1:07:11for
- 1:07:12row
- 1:07:14in large
- 1:07:16partition right?
- 1:07:20For row in large partition, and then
- 1:07:21again for row
- 1:07:23in broadcast table
- 1:07:28For row in broadcast table, yeah.
- 1:07:30And if the condition
- 1:07:34if the condition
- 1:07:36of this row
- 1:07:38which is the la L row and B row
- 1:07:42Let me name it like that.
- 1:07:44L row and B row
- 1:07:47they match, then we are going to emit
- 1:07:51the joined row.
- 1:07:54Right? Then we are going to call that a
- 1:07:57join has happened. Yeah. So, essentially
- 1:08:00there are two steps, which is the
- 1:08:01broadcast step, and then a nested loop
- 1:08:04step. The nested loop step is basically
- 1:08:06a
- 1:08:07two loop, a dual loop between the large
- 1:08:10and large partition and the broadcast
- 1:08:13table. And if the condition matches,
- 1:08:14then we emit the joined row. Now, let's
- 1:08:17understand this through visuals and an
- 1:08:20example. Okay, now let's take this
- 1:08:22example where I have a customer data set
- 1:08:26and I also have an order data set,
- 1:08:28right?
- 1:08:29The customer has a credit limit, and the
- 1:08:33orders has which customer
- 1:08:36spent what amount on a particular order,
- 1:08:39right? And the problem statement is that
- 1:08:42I want to figure out for each of the
- 1:08:44customers
- 1:08:45what are the orders in which the amount
- 1:08:48is greater than my credit limit, right?
- 1:08:51So, for each of the customers, I want to
- 1:08:53figure out all orders. It can be my
- 1:08:55order, it can be somebody else's order.
- 1:08:57The only thing I want to figure out
- 1:08:59which are the orders in which the amount
- 1:09:02of the order is greater than my credit
- 1:09:05limit, right? So, the join condition
- 1:09:08the join condition is this order.amount
- 1:09:13orders.amount
- 1:09:17should be greater than
- 1:09:19order.amount should be greater than
- 1:09:21customer's
- 1:09:25.credit limit.
- 1:09:29Yeah?
- 1:09:30And we see that in the join condition
- 1:09:32there is an inequality over here. And
- 1:09:36for that reason, we are going to choose,
- 1:09:39or rather Spark is going to choose a
- 1:09:42broadcast nested loop join.
- 1:09:46Right? So, this table, let's assume that
- 1:09:48the size of this table is 320 bytes. And
- 1:09:51this is much lesser than the 10
- 1:09:53megabytes default auto broadcast join
- 1:09:56threshold size. That means this table
- 1:09:59can be broadcasted
- 1:10:01this table can be broadcasted to both of
- 1:10:05the executors.
- 1:10:06And both of them are going to receive
- 1:10:08their individual copy of the customer's
- 1:10:12table.
- 1:10:13Right? So, both of them are going to
- 1:10:16receive an individual copy of the
- 1:10:19customer's table.
- 1:10:21Right? So, this is the first step, which
- 1:10:23is the broadcast step.
- 1:10:28The broadcast step, right? Now, let's
- 1:10:30move over to the next step, which is the
- 1:10:33nested loop join step.
- 1:10:36Nested loop step.
- 1:10:42So, in order for the nested loop join to
- 1:10:44happen,
- 1:10:46let's revisit this condition, right?
- 1:10:48Where we are saying order amount is
- 1:10:50greater than customer.credit limit. So,
- 1:10:52if we were to look at this row,
- 1:10:55this row, the order amount has to be
- 1:10:57compared with this row, this row, and
- 1:11:00this row.
- 1:11:01So, we have to figure out whether this
- 1:11:03order amount is it greater than this
- 1:11:05credit limit, which is yes, 250 is
- 1:11:08greater than 200. Then, it is going to
- 1:11:10join this and this row. Similarly, this
- 1:11:13row is going to be
- 1:11:15checked against this row,
- 1:11:17where 250 is it greater than 350? No,
- 1:11:21that means this row is not going to be
- 1:11:22joined, right? And for this one, 250 is
- 1:11:25it greater than 100? Yes, that means
- 1:11:28this row is going to be joined. So, in
- 1:11:29order to visualize that, in order to
- 1:11:32visualize that, what I've done is
- 1:11:34comparing this row with all of the three
- 1:11:37customers' credit limit, right? So, 250
- 1:11:40and
- 1:11:42where the credit limit was 200 and 100,
- 1:11:45these are the output, which is 200 and
- 1:11:48100 over here, C1
- 1:11:52and C3, right? C2 will not be the
- 1:11:55output, right?
- 1:11:57So, this row will join with this one,
- 1:11:58which is C1 and C3. Similarly, this row,
- 1:12:02which is
- 1:12:03this order number 102, will join or will
- 1:12:06not join with anything, actually,
- 1:12:08because this is the amount is 80, the
- 1:12:10credit limit is 200, 350, and 100.
- 1:12:14And similarly, this row will join with
- 1:12:16all of the rows, because 400 is greater
- 1:12:18than 200, 350, and 100, right? So, as a
- 1:12:23result, you're going to see
- 1:12:25two rows, which is one and
- 1:12:28one and two
- 1:12:32will be the output produced for this
- 1:12:33row.
- 1:12:34Zero for this one.
- 1:12:38And three for this one, right? One
- 1:12:42one, two, and three. So, as a result,
- 1:12:44five rows are going to be produced,
- 1:12:46right? And that is what you see over
- 1:12:48here, one
- 1:12:50one, two, three, three, four, and 5,
- 1:12:53right? And this is basically C1 will
- 1:12:56have two rows.
- 1:12:58C3 will have three rows.
- 1:13:00C1 is having two rows and C3 will have
- 1:13:03three rows, right? And the similar case
- 1:13:05is going to happen for this one as well
- 1:13:08where order 104, it is going to be
- 1:13:10matched 150.
- 1:13:13150, is it greater than 200? No, 150 is
- 1:13:16not greater than 350. 150 is greater
- 1:13:18than 100, so it will produce one row.
- 1:13:21And similarly, this is going to produce
- 1:13:23two rows.
- 1:13:25Right? So, you're going to see 2 + 1
- 1:13:28which is three rows as output over here.
- 1:13:34Three rows as output over here.
- 1:13:37Yeah.
- 1:13:382 + 1
- 1:13:40which is three rows over here. Yeah. So,
- 1:13:42this is how a broadcast nested loop join
- 1:13:45works. Okay, now let's get started with
- 1:13:48the final lab.
- 1:13:51Which is the lab for broadcast nested
- 1:13:54loop join.
- 1:13:55And again, let's go ahead and run the
- 1:13:57import disable AQE.
- 1:14:01Right? We read the orders data frame,
- 1:14:05right? And here we are going to do
- 1:14:06something a little different, right? So,
- 1:14:09let's read the order data frame. Let's
- 1:14:11display the orders data frame.
- 1:14:14And here I'm creating a tiers data
- 1:14:18frame, right? And this basically
- 1:14:20contains three columns which is base
- 1:14:22which is the tier economy, standard,
- 1:14:26premium, or enterprise.
- 1:14:28And it basically specify
- 1:14:32which order amount the range of order
- 1:14:34amounts that qualify for that tier,
- 1:14:38right?
- 1:14:39So, let's go ahead and run this as well.
- 1:14:44So, the meaning of this is basically if
- 1:14:46let's say I've placed an order
- 1:14:48and the value of that falls somewhere
- 1:14:51between 1 and 50,
- 1:14:53then it is considered to be an economy
- 1:14:55order.
- 1:14:56Right? And similarly between 51 to 150
- 1:14:59it's the standard order, premium order
- 1:15:01for 151 to 400, right?
- 1:15:05And finally uh between this range to be
- 1:15:08enterprises. Yeah.
- 1:15:10So, there are companies who would want
- 1:15:12to place certain benefits for each
- 1:15:14category or tier of orders. Yeah.
- 1:15:17So, this is how your data frame is going
- 1:15:19to look like. Now, for each of these
- 1:15:21orders, I want to classify in which of
- 1:15:25the tier is it going to fall, right? So,
- 1:15:28the join condition is going to look
- 1:15:30something like this.
- 1:15:31Right?
- 1:15:33So, here I am forcing a broadcast nested
- 1:15:36loop join.
- 1:15:37Uh actually not a broadcast nested loop
- 1:15:39join, but a broadcast step over here.
- 1:15:41Nested loop is something that will be
- 1:15:44automatically chosen because of this
- 1:15:47condition over here. This condition
- 1:15:49contains inequalities, right? So, that
- 1:15:51is why a nested loop join will be
- 1:15:54chosen. And here we are enforcing a
- 1:15:57broadcast of the smaller data set.
- 1:16:00Right? So, the join condition is simply
- 1:16:03uh simply this one which is amount
- 1:16:05should be greater than the minimum
- 1:16:07amount and it should be less than the
- 1:16:10maximum amount. Yeah.
- 1:16:12So, let's go ahead and this is an inner
- 1:16:14join. So, let's go ahead and run this.
- 1:16:17Let's do a show.
- 1:16:20This is going to trigger an action and
- 1:16:22hence is going to trigger Spark jobs.
- 1:16:25And therefore you see over here
- 1:16:28this amount which is 160
- 1:16:31and this fall in the premium tier
- 1:16:33because the minimum and the maximum
- 1:16:36range for any order to fall in the
- 1:16:39premium tier is 151 to 400, right? And
- 1:16:43similarly, we do for the other rows as
- 1:16:45well, right? Now, let's have a look at
- 1:16:49how the plan is going to look like. So,
- 1:16:52let's go ahead and run this.
- 1:16:54And now, let's copy this over broadcast
- 1:16:57nested loop join, and let's paste this
- 1:16:59over here.
- 1:17:01So, let's go through this one by one.
- 1:17:03This is a local table scan basically for
- 1:17:07the data frame that we created locally
- 1:17:11in the Spark session, yeah. And it does
- 1:17:15a filter of the not null condition
- 1:17:17because we're using two of the columns
- 1:17:20in the join condition, right? We're
- 1:17:22using minimum amount and the maximum
- 1:17:24amount. So, it is making sure that both
- 1:17:26of them is not null. And finally, it
- 1:17:30does a broadcast exchange with the
- 1:17:34identity broadcast mode, right? So, we
- 1:17:37understood that broadcast exchange
- 1:17:38simply means that it is going to take
- 1:17:40this data set and serialize it and send
- 1:17:43it over the network to the executor,
- 1:17:46right? Now, what does the identity
- 1:17:48broadcast mode mean? So, first of all,
- 1:17:51broadcast mode, as we discussed earlier,
- 1:17:54it tells you how the broadcast side is
- 1:17:57packaged before it is shipped to the
- 1:18:00executor, right? How the data set data
- 1:18:03set is packaged before it is shipped
- 1:18:05shipped to the executor. Last time we
- 1:18:08studied about the hash relation
- 1:18:10broadcast mode, in which it is packaged
- 1:18:13in the form of a hash map, and then
- 1:18:16broadcasted over to the executor. In the
- 1:18:20identity broadcast mode, this transform
- 1:18:23is actually an identity function, right?
- 1:18:26So, what it simply means is that it is
- 1:18:27going to take the data set as is as a
- 1:18:30plain array, and then send it over to
- 1:18:33the executor. So, there is no hashing,
- 1:18:36no keying, none of that is going to
- 1:18:38happen. It is simply going to take it
- 1:18:40and then pass it to the executors as an
- 1:18:43array. Yeah?
- 1:18:45Now, the second data set, which is the
- 1:18:47orders data set, which is read in, we
- 1:18:49simply do a not null over the amount
- 1:18:52because the amount is the one which is
- 1:18:55used for the orders data frame.
- 1:18:58And finally, we do a nested loop join
- 1:19:02over here. Yeah? So, build right build
- 1:19:05right over here. The right data set is
- 1:19:08this one over here. Build right simply
- 1:19:10means that this is the data set which is
- 1:19:14broadcasted, right? Which Spark holds in
- 1:19:16memory and is broadcasted. And
- 1:19:19therefore, an inner join happens with
- 1:19:21this condition. So, that is how the
- 1:19:23physical plan is going to look like for
- 1:19:26a broadcast nested loop join.
- 1:19:29Now, let's have a look at how the Spark
- 1:19:31UI looks for
- 1:19:33the broadcast nested loop join. And this
- 1:19:36is job number three and four. Let's open
- 1:19:39this up. Let's go to the jobs.
- 1:19:43And here you see job three and four,
- 1:19:44which is one and one two stages, right?
- 1:19:48Uh let's go to SQL data frame, this is
- 1:19:51for job three and four.
- 1:19:56Yeah?
- 1:19:56And you see this one, which is the first
- 1:20:00one is your scan CSV, right?
- 1:20:03Which is on the left-hand side, the
- 1:20:05orders data set, and on the right-hand
- 1:20:07side is the data set that we have
- 1:20:08created. Both of them go through this
- 1:20:10filter optimization that we've seen, and
- 1:20:13this is again a part of Spark own
- 1:20:15optimization strategy. And then we see a
- 1:20:18broadcast exchange using the identity
- 1:20:21broadcast mode. This is broadcasted over
- 1:20:25to the same location where you have the
- 1:20:29orders data set, right? And then there
- 1:20:31is a broadcast nested loop join using
- 1:20:34this condition that we have specified
- 1:20:36that the amount should be between the
- 1:20:39minimum and the maximum amount, right?
- 1:20:43And finally, you see that seven columns
- 1:20:45are projected as an output for the final
- 1:20:49data set. So, this is how the DAG is
- 1:20:51going to look like for a broadcast
- 1:20:54nested loop join. I hope this gave you
- 1:20:56an overview on the kind of joins we have
- 1:20:58in Apache Spark and how do all of them
- 1:21:01behave under the hood and behind the
- 1:21:04scenes. If you like this video, you are
- 1:21:06also going to love the Apache Spark
- 1:21:09playlist where we go in-depth into a lot
- 1:21:12of different topics in Apache Spark. And
- 1:21:14I also produce content on various other
- 1:21:17topics on data engineering and interview
- 1:21:19preparation. Please don't forget to like
- 1:21:21and share this video and subscribe to
- 1:21:23the channel. I will see you in the next
- 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.