Is it network bandwidth or latency?
I suspect it's latency: If you're bottlenecked on latency, the barrier-synchronised nature of many jobs (due to shuffles) lowers network utilisation to the extent that many of the smart network scheduling algorithms the NSDI paper refers to don't work at all.
If it's latency, it also makes sense that a framework that's closer to bare-metal (a highly tuned implementation) can get squeeze more utilisation on a cluster, lowering end to end job times. I wonder if the JVM intrinsically prevents some hardware-specific optimisations due to its memory model.