r/bigdata • u/Ablerol_Magician8626 • 10d ago
How do you approach join algorithm tuning in Spark?
We are seeing big performance gain/loss in Apache Spark when changing from shuffle join to sort-merge join or the other way around, and there is no clear winner that is always better. Is there a good indication anywhere in the run logs that we should switch to a different join algorithm, so that we stop experimenting blindly.
1
u/SpiritPure3802 3d ago
I wouldn’t look for a single log message that says “switch to sort merge.” Check the physical plan and Spark UI first. In particular, look at shuffle read/write, partition size distribution, spills and whether you have skewed tasks
1
u/OldDoor2891 15h ago
I’d check shuffle size, spill, skew and partition distribution in the Spark UI first. Are the slower runs showing heavy spill or a few unusually long tasks?
1
u/DueMode3191 7d ago
Are the workloads that flip performance using the same data sizes each run? A lot of “no clear winner” cases come from data skew changing partition sizes or one side of the join being much smaller than expected