PySpark courseLesson 4 of 6
PySpark course · Lesson 4 of 6
PySpark Joins and Join Strategy
Write correct PySpark joins (inner, left, anti, semi), avoid duplicate-column and fan-out bugs, and understand when Spark broadcasts or sort-merges a join.
On this page
Join semantics in PySpark match SQL (see SQL joins). Two things are specific to Spark: how you express the join to avoid ambiguous columns, and which physical strategy Spark uses, which decides cost.
Setup
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.master("local[2]").appName("example").getOrCreate()
orders = spark.createDataFrame(
[(10, 1, 50), (11, 1, 70), (12, 2, 20), (13, 9, 99)],
["order_id", "customer_id", "amount"],
)
customers = spark.createDataFrame([(1, "Asha"), (2, "Ben"), (3, "Chen")], ["customer_id", "name"])
Join on a column name, not an expression, when names match
joined = orders.join(customers, on="customer_id", how="inner")
print(joined.columns)
['customer_id', 'order_id', 'amount', 'name']
Passing the column name keeps a single customer_id column. Joining with an expression such as orders.customer_id == customers.customer_id keeps both columns, and a later select("customer_id") fails as ambiguous.
Join types
print(orders.join(customers, "customer_id", "left").count()) # 4: all orders kept
print(customers.join(orders, "customer_id", "left_anti").collect()) # customers with no orders
print(customers.join(orders, "customer_id", "left_semi").count()) # customers with at least one order
4
[Row(customer_id=3, name='Chen')]
2
left_antireturns left rows with no match: the efficient way to find missing records.left_semireturns left rows with a match, without adding right-side columns or duplicating rows when there are several matches.
Watch for fan-out
If the right side has several rows per key, the left row repeats. Before a join, check that the key is unique where you expect it:
dupes = customers.groupBy("customer_id").count().filter("count > 1")
print(dupes.count())
0
Join strategies
| Strategy | When Spark uses it | Cost |
|---|---|---|
| Broadcast hash join | One side is small (below spark.sql.autoBroadcastJoinThreshold, 10 MB by default) or hinted |
No shuffle of the large side |
| Sort-merge join | Both sides large, equi-join | Both sides shuffled and sorted |
| Shuffle hash join | Chosen in some cases, for example by AQE or a hint | Shuffle, no sort |
small_join = orders.join(F.broadcast(customers), "customer_id")
small_join.explain() # look for BroadcastHashJoin
With Adaptive Query Execution (on by default), Spark can switch a planned sort-merge join to a broadcast join at run time once it sees that one side is actually small.
Common mistakes
- Expression joins on same-named columns, then ambiguous column errors.
- Using
leftjoin plusfilter(isNull)instead ofleft_anti. - Broadcasting a table that is not small.
- Not checking key uniqueness and silently multiplying rows.
Interview relevance
Expect “broadcast versus sort-merge”, “how do you find records missing from another table?” (anti join) and “why did my row count grow after a join?”.
Key takeaway
Join on column names, use semi and anti joins for existence checks, verify key uniqueness, and let small tables be broadcast.
Progress is saved in this browser only. No account needed.