Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Join two very large tables

PySpark · Joins & Data Layout

Join two very large tables

Hardpyspark-64
scenariojoinsbucketingskewlarge-tables

Question

How would you join two tables of 2 TB each efficiently?

Solution

There is no magic trick for joining two 2 TB tables. A sort-merge join has to shuffle both sides, so the goal is to shuffle less data, shuffle it evenly, and avoid doing it repeatedly.

Shrink the inputs first

  • Filter rows before the join, and in a way that pushes down to the files (partition and Parquet filters). If you only need the last 90 days, do not join the whole history.
  • Select only the columns you need. Wide rows make the shuffle heavier.
  • Pre-aggregate if the join is only used to compute totals. Joining 2 TB to 2 TB and then grouping is much worse than grouping each side to one row per key first, when the logic allows it.

Make the join correct and even

  • Make sure the join keys have the same type on both sides. A string against an integer forces casts that can break pruning and give wrong matches.
  • Check for skew with a group-by count on the key. Handle hot keys with AQE skew handling or salting.
  • Check for NULL keys and drop the ones that cannot match.
  • Set enough shuffle partitions. 2 TB shuffled at 128 MB per partition is about 16,000 partitions. AQE helps tune this at runtime.

If the same join runs every day

Pay the shuffle cost once. Write both tables bucketed by the join key into the same number of buckets, using bucketBy(256, "customer_id").sortBy("customer_id").saveAsTable(...). A join between two tables bucketed the same way on the same key can skip the shuffle and sort. It only works with metastore tables, and it needs matching bucket counts.

Consider the job shape

  • Join only the new data. If yesterday's 2 TB is unchanged, join just today's new partition against the history, not everything every time.
  • Split the work by date range and run slices separately if the cluster is too small to hold the shuffle.

How to answer

Say that you would look at the plan and the stage metrics before changing anything, and then walk through the list: reduce data, fix skew, tune partitions, and consider bucketing if the join repeats. This shows you know the cost is in the shuffle.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext