We will go into this a bit in part 2 of the blog post, but the bottom line is that it doesn't look like GraphX is bottlenecked on barrier-sync latency in this computation. In fact, the iterations in GraphX are quite long and hardly use the network at all, so we're not sure if there's much fine-grained synchronization going on.
That said, leaving the implementation details aside, latency is definitely a big deal, but 10G does help there, too: the latency for sending a fixed-size message can be a lot lower on an idle 10G network than on an idle 1G network. If we're talking about very small synchronization messages, then maybe there isn't much of a gain (network stack overhead dominates), but techniques like our destination-oriented edge processing help reduce the need for very fine-grained synchronization (for this computation at least). The only barrier-synchronization necessary in our fast 10G implementation is at the point at which no more updates are to be sent by any worker (this only happens once per iteration).
You're quite right, however, that a lot of the work on network scheduling for big data computations mentioned in the NSDI paper operates at the coarse-grained level of some kind of 'flow' notion. This would indeed be very hard to disambiguate in our implementation (part 2 will show this in more detail); I'm not convinced that these algorithms would help timely dataflow at all.
Thanks, yes, I realised it wasn't really barrier-sync latency after I wrote the comment. :)
I remember that the NSDI paper actually made an Amdahl's-law-like argument (they give it a new name) and did something to the tune of "let's just eliminate time waiting on the network from the total runtime, which makes the network infinitely fast."
Coming back to the post: If it's CPU overhead, shouldn't Java be pretty competitive with C/C++/Rust for common computations? There might be a lot of other things going on that lower might affect how much one can squeeze from the CPU (GC/object sizes, time spent in reflection/serialisation, maybe?).
It would be great to look at (a) the number of instructions that Java and the Rust implementation execute, and (b) the instructions-per-cycle issued (or its inverse, the CPI) in both cases. If it's memory sync that's slowing down Java, then Java's CPI must be (edit) _higher_ than Rust's.
Yep; the NSDI paper is correct in that there's hardly any time spent waiting on the network (as our traces in part 2 will show). However, that is not to say that the network being faster cannot help: if computation and communication are perfectly overlapped, then "blocked time analysis" (term from the NSDI paper) would not show any potential improvement, but faster communication can still improve the overall runtime (e.g., by reducing busy polling, or crucial updates arriving sooner).
The CPI number investigation is quite a good idea -- we in fact already have these numbers for the Rust-based timely dataflow, but I'll have a look to see how hard it'd be to get them for GraphX/Spark.
That said, leaving the implementation details aside, latency is definitely a big deal, but 10G does help there, too: the latency for sending a fixed-size message can be a lot lower on an idle 10G network than on an idle 1G network. If we're talking about very small synchronization messages, then maybe there isn't much of a gain (network stack overhead dominates), but techniques like our destination-oriented edge processing help reduce the need for very fine-grained synchronization (for this computation at least). The only barrier-synchronization necessary in our fast 10G implementation is at the point at which no more updates are to be sent by any worker (this only happens once per iteration).
You're quite right, however, that a lot of the work on network scheduling for big data computations mentioned in the NSDI paper operates at the coarse-grained level of some kind of 'flow' notion. This would indeed be very hard to disambiguate in our implementation (part 2 will show this in more detail); I'm not convinced that these algorithms would help timely dataflow at all.