Bucketing pre-shuffles a table on write: bucketBy(8, "customer_id") hashes rows into eight buckets and records the layout in the catalog, so joining two tables bucketed alike on the key needs no shuffle.
import contextlib, io, os, time
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1") # force a shuffle-based join
for name, df in [("lines_b", lines), ("customers_b", customers)]:
(df.write.mode("overwrite").bucketBy(8, "customer_id").sortBy("customer_id")
.option("path", f"{os.getcwd()}/out/bucketed/{name}").saveAsTable(name))
def run(label, df):
buf = io.StringIO()
with contextlib.redirect_stdout(buf):
df.explain()
t0 = time.perf_counter()
df.write.format("noop").mode("overwrite").save()
print(f"{label:<9} exchanges: {buf.getvalue().count('Exchange')} "
f"sorts: {buf.getvalue().count('+- Sort')} {time.perf_counter() - t0:.2f}s")
for _ in range(2): # second round is warm
run("plain", lines.join(customers, "customer_id"))
run("bucketed", spark.table("lines_b").join(spark.table("customers_b"), "customer_id"))Output
plain exchanges: 2 sorts: 2 9.86s bucketed exchanges: 0 sorts: 2 5.22s plain exchanges: 2 sorts: 2 6.24s bucketed exchanges: 0 sorts: 2 3.06s
Both Exchange operators disappeared and the warm join ran 1.7 to 2 times as fast over two runs; the sorts stayed because four writer tasks left several files per bucket. Bucketing needs saveAsTable and the same bucket count on both sides, and pays only for tables joined on the same key again and again.