Yes, of course, you should try to build your application so queries and transactions are very short. That solves a great many problems.
But sometimes you can't. Sometimes you just have to do joins across large tables. There just isn't any other way. Your query is going to run for 5, 10, maybe 30 seconds. That's a huge burden on your server, and it will slow down the other tenants. In that case, the only answer is to distribute the queries across other boxes.
I agree with the advice in general -- delay moving to a distributed system as long as possible -- but sometimes you just can't.