What you will be able to do
- Write a DataFrame join with the right
onargument (a column name, a list of names, or a Column expression) and the righthowstring - Predict the rows returned by inner, left and cross joins
- Join on multiple keys and avoid ending up with duplicate key columns
- Mark a small DataFrame for a broadcast join with
broadcast(), and explain whatspark.sql.autoBroadcastJoinThresholdcontrols - Choose between
union,distinct(), SQLUNION/UNION ALLandunionByNamewhen stacking DataFrames
Key concept
join(other, on, how) — Every DataFrame join is one call with three parts: the right-side DataFrame, an on argument that says how rows match, and a how string that picks the join type. Inner is the default. If on is a column name or a list of names, the join is an equi-join on columns that exist on both sides.
1.The anatomy of DataFrame.join
Every join in this lesson goes through DataFrame.join(other, on=None, how=None), apart from crossJoin, which has its own method. other is the right side of the join. on can take four forms: a single column name as a string, a list of column names, a join expression built from Columns, or a list of Column expressions. how names the join type. If you leave it out, you get an inner join.
The form of on matters more than it looks. If you pass a string or a list of strings, the named columns must exist in both DataFrames, and Spark runs an equi-join on them. If you pass a Column expression such as df.name == df2.name, you can write any boolean condition, but both sides' columns stay in the result.
import pyspark.sql.functions as sf
from pyspark.sql import Row
df = spark.createDataFrame([Row(name="Alice", age=2), Row(name="Bob", age=5)])
df2 = spark.createDataFrame([Row(name="Tom", height=80), Row(name="Bob", height=85)])
df.join(df2, "name").show()
# +----+---+------+
# |name|age|height|
# +----+---+------+
# | Bob| 5| 85|
# +----+---+------+The how parameter accepts a fixed set of strings, and several are synonyms. left, leftouter and left_outer all mean the same join, and so do outer, full, fullouter and full_outer. The list also includes cross, semi/leftsemi/left_semi and anti/leftanti/left_anti. On the exam, a spelling like left_outer is a valid choice, not a typo.
Checkpoint 1 of 9· Check yourself
A candidate writes df.join(df2, "name") and passes no third argument. Which join type does Spark run?
The how parameter defaults to inner. That is why Alice and Tom, who each appear on only one side, are missing from the output above.
“default inner. Must be one of: inner, cross, outer, full, fullouter, full_outer, left, leftouter, left_outer”Source: docs.databricks.com
Sources1
2.Inner join vs left join: which rows survive
Inner and left joins differ in what happens to a row that has no match. Take an employee table with six people in departments 1–6 and a department table that lists only departments 1, 2 and 3.
-- Use employee and department tables to demonstrate left join.
> SELECT id, name, employee.deptno, deptname
FROM employee
LEFT JOIN department ON employee.deptno = department.deptno;
105 Chloe 5 NULL
103 Paul 3 Engineering
101 John 1 Marketing
102 Lisa 2 Sales
104 Evan 4 NULL
106 Amy 6 NULLChange LEFT JOIN to INNER JOIN and only Paul, John and Lisa come back, because an inner join returns only rows with matching values in both tables. The DataFrame API follows the same rules. how="inner" (or no how) is the inner join, and how="left" is the left outer join. The left side is the DataFrame you call .join on, and the right side is other.
Join type (how / SQL) | Unmatched left rows | Unmatched right rows | Rows in this example |
|---|---|---|---|
inner / INNER JOIN | Dropped | Dropped | 3 |
left / LEFT [OUTER] JOIN | Kept, with NULLs on the right | Dropped | 6 |
full / FULL [OUTER] JOIN | Kept, with NULLs | Kept, with NULLs | 6 |
cross / CROSS JOIN | No condition: every pairing | No condition: every pairing | 18 |
Checkpoint 2 of 9· Check yourself
You need a list of all customers. For each one, show their order total if they have one, and NULL if they don't. customers is the DataFrame you call .join on. Which how value is correct?
A left join keeps every row of the left DataFrame and fills NULL where the right side has no match. An inner join would drop customers who have no orders.
“Returns all values from the left table reference and the matched values from the right table reference”Source: docs.databricks.com
Checkpoint 3 of 9· Exam question
```python employees = spark.createDataFrame( [(1, "Ana", 10), (2, "Ben", 20), (3, "Cara", None)], ["emp_id", "name", "dept_id"] ) departments = spark.createDataFrame( [(10, "Sales"), (20, "Engineering"), (30, "Marketing")], ["dept_id", "dept_name"] ) result = employees.join(departments, on="dept_id", how="inner") ``` How many rows does `result` contain, and why?
Correct answer: B — Two rows, because Cara's null `dept_id` matches nothing in `departments`, and the unmatched Marketing row is dropped too.
- A. This describes left-join behavior, not inner-join behavior. An inner join only keeps rows where the join key is present on both sides, so a null `dept_id` and an unmatched department row are both excluded.
- B. This is correct. Cara's `dept_id` is null, and null never satisfies an equality join condition, so her row is dropped; Marketing has no matching employee, so it is dropped too, leaving only Ana-Sales and Ben-Engineering.
- C. Inner joins never keep an unmatched left row with nulls filled in for the right side — that padding behavior is specific to outer joins, not inner joins.
- D. Marketing would only survive the join if some employee row had `dept_id = 30`, which none do; an inner join drops departments with no matching employees.
- E. A cartesian-style multiplication of unmatched keys against every row on the other side only happens in a cross join, not an equality-based inner join.
Sources2
3.Joining on multiple keys and avoiding duplicate key columns
To join on several keys, pass a list. With a list of column names such as ["a", "b"], every named column must exist on both sides, and Spark runs an equi-join on all of them. A row matches only when every key is equal. SQL does the same thing with USING (a, b), which the SQL reference defines as an ON clause that joins the equality tests with AND, as in ON first.a = second.a AND first.b = second.b. For conditions that aren't plain equality, pass Column expressions instead, either as one expression combined with & or as a list of Columns.
joined = df.join(df2, df.name == df2.name, "outer").sort(sf.desc(df.name))
joined.show()
# +-----+----+----+------+
# | name| age|name|height|
# +-----+----+----+------+
# | Bob| 5| Bob| 85|
# |Alice| 2|NULL| NULL|
# | NULL|NULL| Tom| 80|
# +-----+----+----+------+Joining by name (a string or list of strings in the DataFrame API, or USING/NATURAL in SQL) collapses each key into a single output column, placed first. An expression join keeps the key from each side. That is usually why a later select("name") fails as ambiguous. When you can't avoid an expression join, for example in a self-join, give each side an alias with df.alias("a") and refer to columns as "a.name" and "b.name".
Checkpoint 4 of 9· Check yourself
orders and shipments both have region and order_id columns. Which call joins on both keys and leaves one region and one order_id column in the result?
A list of column names runs an equi-join on every listed column and keeps one copy of each key. The expression version matches the same rows but keeps both sides' region and order_id.
“SELECT * will only show one occurrence for each of the columns used to match first”Source: docs.databricks.com
4.Broadcast joins: sending the small side to every worker
Joins often pair a large fact table with a small lookup table. A broadcast join copies the whole small table to all worker nodes, so each worker can join its part of the large table locally. A Databricks engineering post puts it this way: broadcast hash join is usually the fastest strategy when one side fits comfortably in memory.
You can trigger it in two ways. The explicit way is pyspark.sql.functions.broadcast(df), which marks a DataFrame as small enough to broadcast. You then join as usual. The automatic way is the setting spark.sql.autoBroadcastJoinThreshold, which defaults to 10485760 bytes (10 MB). Spark plans a broadcast hash join on its own when a relation's estimated size is below this threshold. Setting it to -1 turns automatic broadcasting off. A BROADCAST join hint tells Spark to use that strategy for the relation you mark.
Checkpoint 5 of 9· Fill the gap
Which function marks df_small for a broadcast join?
from pyspark.sql import functions as dbf
df = spark.createDataFrame([1, 2, 3, 3, 4], "int")
df_small = spark.range(3)
df_b = dbf. ? (df_small)
df.join(df_b, df.value == df_small.id).show()broadcast() in pyspark.sql.functions returns the DataFrame marked for a broadcast join. cache and persist only control storage and say nothing about the join strategy.
Checkpoint 6 of 9· Exam question
A developer runs `df_small.join(df_large, "key")` without any join hints, and Spark's optimizer automatically converts it into a broadcast hash join. Under the default configuration, what determines whether Spark performs this automatic conversion?
Correct answer: D — Spark estimates `df_small`'s size and compares it against `spark.sql.autoBroadcastJoinThreshold`, which defaults to 10 MB.
- A. `spark.sql.shuffle.partitions` controls how many partitions a shuffle produces after a wide transformation; it plays no role in deciding whether a table is small enough to broadcast.
- B. Join order in the code has no bearing on which side gets broadcast — the optimizer picks based on estimated size, and either side could be chosen.
- C. Automatic broadcast selection is based on estimated data size in bytes, not row count, since a narrow table with millions of tiny rows could still be small while a wide table with a few rows could be large.
- D. This is correct. Spark's catalyst optimizer compares the estimated byte size of each side against `spark.sql.autoBroadcastJoinThreshold`, whose default is 10485760 bytes (10 MB), and automatically broadcasts a side that falls under it.
- E. Automatic broadcast conversion is a cost-based optimizer decision that happens regardless of whether `broadcast()` has been imported; the import only matters for forcing a broadcast explicitly.
5.Cross joins: the Cartesian product
A cross join has no join condition. It pairs every row on the left with every row on the right, so the result has left × right rows. In the DataFrame API you write it as df.crossJoin(other). In SQL you write CROSS JOIN. The same employee/department tables (6 rows and 3 rows) produce 18 rows.
df.crossJoin(df2.select("height")).select("age", "name", "height"
).orderBy("age", "name", "height").show()
# +---+-----+------+
# |age| name|height|
# +---+-----+------+
# | 14| Tom| 80|
# | 14| Tom| 85|
# | 16| Bob| 80|
# | 16| Bob| 85|
# | 23|Alice| 80|
# | 23|Alice| 85|
# +---+-----+------+You can also get a cross join by accident. The SQL reference warns that if you leave out the join criteria, any join type behaves as a CROSS JOIN. Because the output grows with the product of the two inputs, Databricks' join-performance guidance says cross joins are expensive and should be removed from workloads that need low latency or frequent recomputation.
Checkpoint 7 of 9· Check yourself
A SQL query reads SELECT * FROM a INNER JOIN b with no ON or USING clause. What does it return?
Without join criteria, every join type behaves as a cross join. You get every combination of rows. Matching on same-named columns is what NATURAL does, and that keyword is missing here.
“If you omit the join_criteria the semantic of any join_type becomes that of a CROSS JOIN.”Source: docs.databricks.com
6.Stacking rows: union, UNION ALL and unionByName
Joins add columns side by side. Unions stack rows, and the main question is what happens to duplicate rows.
df1 = spark.createDataFrame([(1, 'A'), (2, 'B'), (3, 'C')], ['id', 'value'])
df2 = spark.createDataFrame([(3, 'C'), (4, 'D')], ['id', 'value'])
df3 = df1.union(df2).distinct().sort("id")
df3.show()
# +---+-----+
# | id|value|
# +---+-----+
# | 1| A|
# | 2| B|
# | 3| C|
# | 4| D|
# +---+-----+SQL works the other way round. In SQL, UNION without a qualifier means UNION DISTINCT and removes duplicates. You write UNION ALL to keep them. Both sides must have the same number of columns with compatible types, otherwise you get NUM_COLUMNS_MISMATCH or INCOMPATIBLE_COLUMN_TYPE. So the DataFrame method union() gives you SQL's UNION ALL result. A note on scope: the sources for this lesson don't document a separate DataFrame unionAll method. They list only union and unionByName. What this lesson can teach is the union-all behaviour, which keeps duplicates.
union() matches columns by position, not by name. If one DataFrame has id, value and the other has value, id, union() silently puts each value under the wrong column. unionByName matches columns by name instead. With allowMissingColumns=True, it also accepts columns that exist on only one side and fills the gaps with NULL.
df1 = spark.createDataFrame([[1, 2, 3]], ["col0", "col1", "col2"])
df2 = spark.createDataFrame([[4, 5, 6]], ["col1", "col2", "col0"])
df1.unionByName(df2).show()
# +----+----+----+
# |col0|col1|col2|
# +----+----+----+
# | 1| 2| 3|
# | 6| 4| 5|
# +----+----+----+Checkpoint 8 of 9· Match them up
Match each operation to how it treats duplicates and columns
Tap a term, then the definition that fits it.
The DataFrame union() and SQL UNION ALL both keep duplicates. SQL's plain UNION deduplicates, and unionByName is the only one that matches columns by name.
“If ALL is specified duplicate rows are preserved. If DISTINCT is specified the result does not contain any duplicate rows. This is the default.”Source: docs.databricks.com
Checkpoint 9 of 9· Exam question
A pipeline joins a 500 GB `transactions` DataFrame with a 2 MB `country_codes` lookup DataFrame. `spark.sql.autoBroadcastJoinThreshold` has been set to `-1` in the cluster's Spark config, disabling automatic broadcast selection, but the developer still wants this specific join to avoid a shuffle. Which line correctly forces `country_codes` to be broadcast? ```python from pyspark.sql.functions import broadcast result = transactions.join( ____________, on="country_id", how="inner" ) ```
Correct answer: A — `broadcast(country_codes)`
- A. This is correct. Wrapping the small lookup DataFrame in `pyspark.sql.functions.broadcast()` marks it for a broadcast hash join, which works even when the automatic threshold has been disabled.
- B. `broadcast` is a module-level function in `pyspark.sql.functions`, not a DataFrame method, so calling `.broadcast()` on a DataFrame instance raises an `AttributeError`.
- C. Collapsing a DataFrame to one partition changes its parallelism but does not tell the join planner to skip the shuffle; the join would still execute as a shuffle-based join.
- D. Passing the string `"shuffle"` to `.hint()` requests a sort-merge/shuffle join strategy, which is the opposite of what avoiding a shuffle requires here.
- E. Broadcasting the 500 GB `transactions` DataFrame instead of the 2 MB lookup table would try to replicate a huge dataset to every executor, which is the wrong side to mark and would likely exhaust executor memory.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.
DataFrame.union()removes duplicate rows the way SQLUNIONdoes.Why is that wrong?union()keeps duplicates, so it behaves like SQLUNION ALL. Add.distinct()if you want them removed.2.
union()lines columns up by name, so the column order of the two DataFrames doesn't matter.Why is that wrong?union()matches columns by position. When the column orders differ, useunionByName.3.A SQL join with no
ONorUSINGclause fails to run.Why is that wrong?Without join criteria, any join type behaves as a cross join and returns the Cartesian product, which can be very large.
Covered in Cross joins: the Cartesian product
4.Spark only runs a broadcast join when you call
broadcast().Why is that wrong?Spark also broadcasts on its own when a table's estimated size is below
spark.sql.autoBroadcastJoinThreshold, which defaults to 10 MB. Setting the threshold to -1 turns this off.Covered in Broadcast joins: sending the small side to every worker
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“the column(s) must exist on both sides, and this performs an equi-join.”
↩︎ The anatomy of DataFrame.join“default inner. Must be one of: inner, cross, outer, full, fullouter, full_outer, left, leftouter, left_outer”
↩︎ The anatomy of DataFrame.join“the column(s) must exist on both sides, and this performs an equi-join.”
↩︎ Joining on multiple keys and avoiding duplicate key columns“a string for the join column name, a list of column names, a join expression (Column), or a list of Columns.”
↩︎ Key concept - 2.
“Returns the rows that have matching values in both table references.”
↩︎ Inner join vs left join: which rows survive“Returns all values from the left table reference and the matched values from the right table reference”
↩︎ Inner join vs left join: which rows survive“Matches the rows by comparing equality for list of columns column_name which must exist in both relations.”
↩︎ Joining on multiple keys and avoiding duplicate key columns“SELECT * will only show one occurrence for each of the columns used to match first”
↩︎ Joining on multiple keys and avoiding duplicate key columns“If you omit the join_criteria the semantic of any join_type becomes that of a CROSS JOIN.”
↩︎ Cross joins: the Cartesian product“If you omit the join_criteria the semantic of any join_type becomes that of a CROSS JOIN.”
↩︎ Exam trap 3 - 3.
“Marks a DataFrame as small enough for use in broadcast joins.”
↩︎ Broadcast joins: sending the small side to every worker - 4.https://spark.apache.org/docs/latest/sql-performance-tuning.htmlSecondary source
“Configures the maximum size in bytes for a table that will be broadcast to all worker nodes when performing a join.”
↩︎ Broadcast joins: sending the small side to every worker“By setting this value to -1, broadcasting can be disabled.”
↩︎ Broadcast joins: sending the small side to every worker“instruct Spark to use the hinted strategy on each specified relation when joining them with another relation”
↩︎ Broadcast joins: sending the small side to every worker“Configures the maximum size in bytes for a table that will be broadcast to all worker nodes when performing a join.”
↩︎ Exam trap 4 - 5.https://www.databricks.com/blog/2020/05/29/adaptive-query-execution-speeding-up-spark-sql-at-runtime.htmlSecondary source
“broadcast hash join is usually the most performant if one side of the join can fit well in memory.”
↩︎ Broadcast joins: sending the small side to every worker“Spark plans a broadcast hash join if the estimated size of a join relation is lower than the broadcast-size threshold.”
↩︎ Broadcast joins: sending the small side to every worker - 6.
“Returns the cartesian product with another DataFrame.”
↩︎ Cross joins: the Cartesian product - 7.
“Cross joins are expensive.”
↩︎ Cross joins: the Cartesian product - 8.
“with no automatic deduplication of elements.”
↩︎ Stacking rows: union, UNION ALL and unionByName“The method resolves columns by position (not by name), following the standard behavior in SQL.”
↩︎ Stacking rows: union, UNION ALL and unionByName“with no automatic deduplication of elements.”
↩︎ Exam trap 1“The method resolves columns by position (not by name), following the standard behavior in SQL.”
↩︎ Exam trap 2 - 9.https://docs.databricks.com/aws/en/sql/language-manual/sql-ref-syntax-qry-select-setopsOfficial docs
“If ALL is specified duplicate rows are preserved. If DISTINCT is specified the result does not contain any duplicate rows. This is the default.”
↩︎ Stacking rows: union, UNION ALL and unionByName“Both subqueries must have the same number of columns and share a least common type for each respective column.”
↩︎ Stacking rows: union, UNION ALL and unionByName - 10.
“resolving columns by name (rather than position). When allowMissingColumns is True, missing columns will be filled with null.”
↩︎ Stacking rows: union, UNION ALL and unionByName