6 ms·
Spark Breaks Previous Large-Scale Sort Record
- ddlatham 12y agoMost recent results I can see to compare to (Google, Yahoo, Quantcast): https://www.quantcast.com/inside-quantcast/2013/12/petabyte-sort/ https://www.quantcast.com/inside-quantcast/2013/12/petabyte-...
- showerst 12y agoFor the curious, the (max) price of those instances is $6.82/hr, so 206 * 6.82 * (23/60) = $538.55 --If they did it with non-reserved instances in US East. If they used reserved instances in USEast, it drops to $181. Obviously there are lots of costs involved beside the final perfect run, but it's an interesting ballpark.
- sp332 12y agoYou have to put spaces around your * 's to keep HN from italicizing everything.
- showerst 12y agoOops, edited. Thanks!
- gphil 12y agoOne of the big positives of Spark is that its architecture is amenable to having workers run on spot instances, which are even cheaper than reserved instances.
- discardorama 12y agoIt's interesting, but not earth-shattering. The "10x fewer nodes" means nothing; how powerful are the new nodes? What's the network? Do you use SSDs? etc. etc. They also tuned their code to this specific problem: "Exploiting Cache Locality: In the sort benchmark, each record is 100 bytes, where the sort key is the first 10 bytes. As we were profiling our sort program, we noticed the cache miss rate was high, because each comparison required an object pointer lookup that was random..... Combining TimSort with our new layout to exploit cache locality, the CPU time for sorting was reduced by a factor of 5." I would love to see MR and Spark compete on the exact same hardware configuration.
- saryant 12y agoThe article says exactly what they ran on. EC2 i2.8xlarge instances which have 32 cores, 800GB SSD and 244GB RAM.
- discardorama 12y agoI read that. But how does that compare with the nodes they're comparing against ("10x fewer nodes")?
- rxin 12y agoThe old entry had 10Gb/s <full-duplex> (40 nodes/rack 160Gbps rack to spine. 2.5:1 subscription), 64GB of RAM, and 12 x 3TB SATA. The network part is probably the most important one here, and both have comparable network.
- discardorama 12y agoSince each node was handling 500GB of data (roughly), I think the disk speed may have been a more critical factor since each node had 244GB of memory. Their nodes used SSDs; the older nodes used spinning rust. The seek times alone will be a killer.
- panarky 12y agoThe 100 terabyte benchmark used 206 Spark nodes, compared with 2100 Hadoop nodes. Going up to 1 petabyte, the Hadoop comparison adds more nodes, 3800, while the Spark benchmark actually reduced the number of nodes to 190. Does Spark scale well beyond ~200 nodes, or does the network become the bottleneck? In any case, it's an impressive result considering that they didn't use Spark's in-memory cache.
- rxin 12y agoIt is mainly the cost of getting nodes from EC2 at that point. It becomes hard to get a huge number of i2.8xl instances. Spark runs fine on thousands of nodes.
- Lanzaa 12y agoI believe the network had become a bottleneck. As per the article: > [O]ur Spark cluster was able to sustain ... 1.1 GB/s/node network activity during the reduce phase, saturating the 10Gbps link available on these machines. If the network is the bottleneck it makes sense to reduce the number of nodes to reduce the network communications.
- rxin 12y agoThe job is actually very linearly scalable. i.e. running it on 200 nodes roughly doubles the throughput of 100 nodes.
- rxin 12y agoThanks for sharing this. I'm the author of this blog post. Free free to ask me anything.
- AustinBGibbons 12y agoHi Reynold! Do you have numbers / intuition for how previous versions of spark would have run? I'm upgrading (soon) from spark 0.8 to spark 1.1 and am curious to see the performance gains (especially w.r.t. shuffles)
- rxin 12y agoHi Austin, We haven't tested Spark 0.8 at this scale. In general Spark is advancing at a rapid rate that 1.1 is very very different from 0.8.
- chad_walters 12y agoYour post mentions "single root IO virtualization" as a factor in maximizing network performance. I am wondering what the impact of this was in your sorting. Do you have data for runs where you didn't enable this?
- rxin 12y agoIt was part of the enhanced networking. Without enhanced networking, we were getting about 600MB/s, vs 1.1GB/s with.
- ch 12y agoCurious if using the Sparrow scheduler would have been a net gain/loss to this type of work load?
- rxin 12y agoIt would help a little bit (maybe a few percent), but not much because the scheduling latency was relatively low for these tasks (the largest scheduling delay was ~10 secs, whereas each task takes minutes).
- chubot 12y agoFWIW, in 2011, Google wrote that they achieved a PB sort in 33 minutes on 8000 computers, vs. 234 minutes on 190 computers with 6080 cores reported by Spark here. http://googleresearch.blogspot.com/2011/09/sorting-petabytes-with-mapreduce-next.html http://googleresearch.blogspot.com/2011/09/sorting-petabytes...
- deeviant 12y agoI'm not sure why you list Google as using "8000 computers" and Spark using "190 computers with 6080 cores". Using two different metrics for two like things seems like there is some sort of implication there. Were Google's machines single-cored?
- chubot 12y agoI'm just writing down exactly what they reported. They used different metrics. Certainly it would be interesting to have an apples to apples comparison. But the computers aren't the only thing that is relevant -- we also need to know about the networking hardware.
- coldcode 12y agoNo matter what the circumstances, sorting 100 TB or 1 PB of anything is impressive, much less doing it during the time it takes me to eat lunch.
- metronius 12y agoWhat change the biggest difference in performance between Spark and MapReduce?
- gtrubetskoy 12y agoThe strength of Hadoop isn't so much speed but that it's been around and there is a pretty impressive and fairly mature set of projects that comprises the Hadoop ecosystem, from Yarn to Hive, etc. There are still many issues to resolve, and this evolution will continue for decades to come. The TB sort benchmark is pretty useless to me - I am much more concerned with stability, a vibrant community (which means people, the software they write and institutions using Hadoop in production). Last time I tinkered with Spark (this was over a year ago) it was so buggy, next to useless, but perhaps things have changed. Still - the idea that there is some sort of a revolutionary new approach that is paradigm-shifting and is way better than anything before should be viewed with extreme skepticism. The problem of distributed computing is not a simple one. I remember tinkering with the Linux kernel back in the mid nineties, and 20 years later it still has ways to go to improve. Twenty years from now it might or might not be Hadoop that is the tool for this sort of thing, we don't know, but I will not take seriously anything or anyone who claims that the "next best thing" is here in 2014.
- metronius 12y ago1. Cloudera left M/R for Spark, Mahout left M/R for Spark. Spark community will be huge soon. 2. Yes, Spark was/is buggy. 3. For me Spark is really paradigm shift, next generation framework compared to M/R
- gtrubetskoy 12y agoHadoop != M/R, FWIW. M/R support is left in Yarn for backwards compatibility mostly. If by M/R you mean Hadoop - Cloudera has done no such thing, their largest customer base is Hadoop. As to "paradigm shift", we're so early in this that I don't think there even is a paradigm to shift.
- metronius 12y agoI mean M/R by M/R ;) http://vision.cloudera.com/mapreduce-spark/ http://vision.cloudera.com/mapreduce-spark/
- xxcode 12y agoDoes this mean that Spark is the new God. If this is the case, then Databricks will be the next Cloudera. Cloudera is probably a 10B+ company. Good job
- gane5h 12y agoGoing on a tangent here: this benchmark highlights the difficulty of sorting in general. Sorts are necessary for computing percentiles (such as the median.) In practical applications, an approximate algorithm such as t-digest should suffice. You can return results in seconds as opposed to "chest thumping" benchmarks to prove a point. :) I wrote a post on this: http://www.silota.com/site-search-blog/approximate-median-computation-big-data/ http://www.silota.com/site-search-blog/approximate-median-co...
- sonoffett 12y agoPerhaps I misunderstand your comment, but you actually don't need to sort to compute a median (see O(n) median of medians algorithm [1]). [1] http://en.wikipedia.org/wiki/Median_of_medians http://en.wikipedia.org/wiki/Median_of_medians
- vinay_ys 12y agoWhere can I find the source code & instructions on how to reproduce this benchmark?