3 ms·
It's hardware. They have have 52 servers, each with 16 disks which I assume are capable of I/O at 100MB/s. That's an aggregate disk I/O rate of 83.2 GB/s. Th
by dadkins 16y ago
It's hardware. They have have 52 servers, each with 16 disks which I assume are capable of I/O at 100MB/s. That's an aggregate disk I/O rate of 83.2 GB/s. Their network switch provides 10Gbps connectivity between each pair, which implies they can move data between nodes at 65 GB/s. With that hardware setup they can read 1TB from disk and distribute in around 15 seconds. They can write sorted data back to disk in 12 seconds. That leaves a little over 30 seconds for each node to sort 19.2 GB in memory. But each node has 8 cores, so each core has to sort 2.4 GB in 30 seconds. I've seen Burstsort chew through that amount of data in half the time.
- dadkins 16y agoI should add that it's the massive network switch that makes this all possible. You can scale disk I/O by adding disks, and you can scale in-memory sorting by adding cores and memory. But it's awfully expensive to buy your way out of a bottleneck in your network when everyone needs to talk to each other at full speed.
- timr 16y agoThat's a blunt pronouncement. It sounds very much like software improvements played a role: "We asked ourselves, ‘What does it mean to build a balanced system where we are not wasting any system resources in carrying out high end computation?’” said Vahdat. “If you are idling your processors or not using all your RAM, you’re burning energy and losing efficiency.” From personal experience, the hardest part of dealing with super-high-end hardware is utilizing the capacity that it gives you. You don't just throw a larger computer at a program written in Blub for a smaller system, and gain instant world-record performance. If nothing else, you usually have to re-architect your code to match the assumptions of the new hardware. Also, the I/O numbers you're assuming are probably very high: you never get 100% utilization of a network interface -- 70% is more realistic. And for a problem like this, you're bounded not by the product of the bandwidth of the pairwise interconnects (i.e. (10Gbps * 52) / 8 == 65GB/s), but the bandwidth of any single interconnect, because you can get to a point where you're transferring most of the data across a few connections. Also, latency matters a lot, because a small connection latency can easily throttle your overall bandwidth when you're moving lots of small chunks of data. Realistically, these guys probably had something more like 7Gbps / 8 = ~.88 GB/s bandwidth between nodes, not counting latency. That means that it would take ~1163s to transfer 1TB between any two nodes -- which means that that naive solution (distributed merge sort where the intermediates are transferred between nodes) is already ruled out. So, like I said, we're back to software improvements, if only to utilize the hardware available.