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)