That remaining hot key needs explicit salting, or pull it into its own branch. We used 16 salt buckets for any key above 0.5% of volume.
Processing ~500GB daily events joined to a 2M row dimension on user_id. One join key has ~8% of traffic.
df = events.join(dim, "user_id", "left")I tried broadcast(dim) but driver memory spiked. Salting sounds right but I haven't implemented it in production before. What pattern do you use?
That remaining hot key needs explicit salting, or pull it into its own branch. We used 16 salt buckets for any key above 0.5% of volume.
Idempotent writes with merge keys saved us during backfills.
Broadcast only helps when the small table truly fits in memory across tasks. 2M rows with wide columns can still be too large. Check spark.sql.autoBroadcastJoinThreshold and the actual serialized size in the UI.
Fair, dim is around 1.8GB serialized. Too big to broadcast reliably.
Sign in to reply.
© 2026 Lakebench, operated by Hunnurji Rao. Bengaluru, Karnataka, India.
No cluster. No install. Just the tab.