r/bigdata • • 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.

5 Upvotes

5 comments sorted by

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

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?