ML Engineer MasterClass (October) | 4 seats left

Spark · Broadcast joins
AmazonAmazon Analytics
Amazon · Tune Spark Workloads

Broadcast joins

Use small lookup tables without assuming every join should broadcast.

Step 1 of 6 · Learn

Broadcast a customer lookup

broadcast marks a small join side for broadcasting. Spark can distribute it to executors so the larger side avoids a join-key shuffle. This still has transfer and memory costs. Unsupported join types or other planning constraints can prevent the requested strategy. Our customer key is unique, so enriching an order does not multiply it.

Lesson reference: PySpark and Scala

Broadcast a customer lookup

broadcast marks a small join side for broadcasting. Spark can distribute it to executors so the larger side avoids a join-key shuffle. This still has transfer and memory costs. Unsupported join types or other planning constraints can prevent the requested strategy. Our customer key is unique, so enriching an order does not multiply it.

PySpark example

from pyspark.sql.functions import col, broadcast, count, sum

result = orders.filter(col("total") >= 0).join(
    broadcast(customers.select("customer_id", "customer_name")), "customer_id", "inner"
).select("order_id", "customer_name", "total").orderBy("order_id")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._

val result = orders.filter(col("total") >= 0).join(
    broadcast(customers.select("customer_id", "customer_name")), Seq("customer_id"), "inner"
).select("order_id", "customer_name", "total").orderBy("order_id")
result.show(false)

Shrink the broadcast side

Filter a lookup before broadcasting when the query only needs a subset. An inner join then retains only lines whose product matches that subset. The category filter changes the result; it is not merely a performance hint. Do not force-broadcast a large table because a small exercise happens to fit.

PySpark example

from pyspark.sql.functions import col, broadcast, count, sum

result = order_items.join(
    broadcast(products.filter(col("category") == "Electronics")), "product_id", "inner"
).select("order_item_id", "product_name", "quantity").orderBy("order_item_id")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._

val result = order_items.join(
    broadcast(products.filter(col("category") === "Electronics")), Seq("product_id"), "inner"
).select("order_item_id", "product_name", "quantity").orderBy("order_item_id")
result.show(false)

Broadcast membership keys

A left-semi join returns left-side rows with a match; a left-anti join returns those without one. Broadcasting the small distinct order-customer key list lets this query test membership without adding right-side columns. Fran has no orders, so semi and anti partition the customer list differently.

PySpark example

from pyspark.sql.functions import col, broadcast, count, sum

result = customers.join(
    broadcast(orders.select("customer_id").distinct()), "customer_id", "left_semi"
).select("customer_id", "customer_name").orderBy("customer_id")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._

val result = customers.join(
    broadcast(orders.select("customer_id").distinct()), Seq("customer_id"), "left_semi"
).select("customer_id", "customer_name").orderBy("customer_id")
result.show(false)
example.pyPySpark
1. Small lookupFilter and project the lookup before broadcasting.
2. Broadcast exchangeDistribute the small side to executors; it must fit in memory.
3. Probe locallyMatch the other side against the lookup. Final display sorting is separate.
Source orders6 rows
order_idcustomer_idstatustotalitem_count
1001101Delivered89.52
1002102Shipped1493
1003101Cancelled351
1004103Delivered219.994
1005104Delivered49.991
1006105Shipped1202
Second input customers6 rows
customer_idcustomer_namecity
101AriSeattle
102BoAustin
103CamChicago
104DeeBoston
105EliDenver
106FranPortland
Broadcast a customer lookup
Result6 rows
order_idcustomer_nametotal
1001Ari89.5
1002Bo149
1003Ari35
1004Cam219.99
1005Dee49.99
1006Eli120
The example enriches all six orders; customer 101 supplies Ari for two orders.

The example is loaded in the editor. Run it as written, then try a small change.

Runs on the Spark backend. First startup may take a moment.

Run your code to see the result.