8 ms·
Nobody ever got fired for buying a cluster
- gruseom 13y agoOff-topic, but since younger HNers may not know what the title is riffing on, https://www.google.com/search?q=nobody%20got%20fired%20for%20buying%20ibm https://www.google.com/search?q=nobody%20got%20fired%20for%2...
- samspenc 13y agoNice try, but a few thoughts from someone who now spends a lot of time in Hadoop/MapReduce. I will admit it took me a while to warm up to the whole concept, but I'm now so familiar with it that I'm not able to think of compute before Hadoop. I'm able to fire off pretty intensive MapReduce jobs on an Amazon Elastic MapReduce cluster with many nodes for a fraction of the price mentioned in the post (less than $100). While I can imagine I could repurpose all my MapReduce/Hadoop code to run on a single box - especially since Amazon does offer several high-memory instances today - I would be loathe to. The MapReduce framework provides a really nice framework that lets me horizontally scale out compute, rather than vertically, and that is really handy at terabyte-data volumes (data warehousing and large-data analytics.)
- wmf 13y agoIt sounds like you're using Hadoop correctly, which is fine. But a lot of people are using "big data" that isn't very big (<1TB) and crunching it with small clusters that are less powerful than a single server due to the massive overhead of Hadoop.
- msellout 13y agoIt's less about the number of bytes than the number of records produced by the map step(s). Sometimes 10gb input data will produce many billions of records to reduce (if you're looking at combinations of things). Basically, if your computation would never exceed memory on a single machine, then it is more processor efficient to code a more simple multi processing method and run on a single box than to code a map reduce on a cluster. But what if you're not sure of the input data size? Processors are cheap. Engineers are expensive. Code the thing once for map reduce and you don't have to worry about making the transition later.
- wmf 13y agoI agree that you should code your analytics once. I think the lesson from this work is that the market could benefit from a Hadoop-compatible but non-clustered runtime which should be easier to run and 10X faster.
- dude_abides 13y agoMapReduce is handy at terabyte-data volumes, but how often do we run jobs with input size > 1 TB? According to the paper: at least two analytics production clusters (at Microsoft and Yahoo) have median job input sizes under 14 GB, and 90% of jobs on a Facebook cluster have input sizes under 100 GB The other point the authors make is that DRAM costs are following Moore's law and terabyte-data workloads should "soon" be cost-feasible on single servers with DRAM.
- chubot 13y agoThe thing is that MapReduce is still a very good programming paradigm for a single box. Machines have 8-64 cores these days -- you don't want to write multi-threaded code every time you want to do an analysis. So you can write MapReduce, but use a multicore framework instead of a cluster framework (hadoop). The unfortunate thing is that there is no popular open source implementation of a multicore mapreduce, so people use Hadoop, which is wasteful on small data sets, as mentioned. But the great part about it is that you will use the same application code for both. I fully expect in 5 years or so that people will be running multicore mapreduce jobs on 100 or 1000 core boxes.
- scott_s 13y agoNot all algorithms map well to MapReduce, which is the authors' main point. They explain one such example in the paper (section 3).
- eropple 13y ago> The unfortunate thing is that there is no popular open source implementation of a multicore mapreduce Not popular, I'd agree, but I have had a lot of success with one-off Akka projects. My mappers and reducers are usually under 10 lines of Scala (more if I'm stuck writing Java, obviously).
- kyzyl 13y agoThere is Disco[1], which is a MR framework written in Python and Erlang. It's open source and pretty awesome, and if I'm not mistaken it will leverage multi-core processors, no need for a cluster. [1] http://discoproject.org/about http://discoproject.org/about
- shmerl 13y agoIt's a tradeoff - clear simplicity for limited logic. I.e. not every problem that requires distributed computation fits well with map reduce, yet many attempt to fit the problem into it to begin with, instead of trying to shape a right distributed solution for particular problem.
- nazka 13y agoDo you have any website, tutorial, pdf, etc. for me to reduce the latency on Hadoop?
- deleted 13y ago[deleted]
- marshray 13y ago> a single “big memory” (192 GB) server we are using has the performance capability of approximately 14 standard (12 GB) servers. That's 168 GB RAM total, within a single upgrade unit of their 192 GB single server, suggesting the problem is dominated by RAM. If I had $6638 + 2640 = $9278 to spend on computing hardware from NewEgg, how about: 1 at $1000 of HP ProLiant DL360e Gen8 Rack Server System Intel Xeon E5-2403 1.8GHz 4C/4T 4GB http://www.newegg.com/Product/Product.aspx?Item=N82E16859107943 http://h10010.www1.hp.com/wwpc/us/en/sm/WF06a/15351-15351-3328412-241644-241475-5249570.html?dnr=1 (12 DIMM slots) 4 at $70 of Kingston 8GB 240-Pin DDR3 SDRAM ECC Registered DDR3 1333 Server Memory Model KVR13LR9S4/8 http://www.newegg.com/Product/Product.aspx?Item=N82E16820239540 1 at $54 of Seagate Barracuda ST250DM000 250GB 7200 RPM 16MB Cache SATA 6.0Gb/s 3.5" Internal Hard Drive http://www.newegg.com/Product/Product.aspx?Item=N82E16822148765 $1334 ea server, 7 servers = $9338 So we could get 7 of these low-end name-brand 16 GB servers for the same money to give us 224 GB RAM. > MR++ runs on 27 servers whereas the standalone configurations are a single server running a single-threaded implementation Sure, nothing will beat a single system at message-passing algorithms when the entire graph fits in main memory. But when the dataset outgrows that (and it will), we can triple the RAM in the empty slots, or add more servers in units of $1334 instead of having to rewrite your whole analysis.
- wmf 13y agoYou can't compare old pricing against new. Now you can get 384 GB for under $9K and 768 GB would be maybe $15K.
- marshray 13y agoOh, the paper is dated January 2013 but I see now something about "prices converted to US$ as of 8 August 2012". I'm confused. I think the link changed to a different version of the paper. The numbers I went by aren't even in there any more!
- wmf 13y agoYeah, I convinced the moderators to update the link to a newer version of the paper. But keep in mind that a paper published in 2013 probably describes research performed in 2012 on hardware purchased in 2011 (if they're lucky). Getting back to hardware, the premium for 4-socket servers appears to be much smaller than it used to be, so now you can get a ton of DIMM slots for a reasonable price.
- Theory5 13y ago"... There IS the occasional savage beating, and more than their fair share of suicides. But that has "statistical clustering" all over it. "
- zdw 13y agototally offtopic typographic comment - anyone else seeing the alt-font fi ligature in the title? It's quasi-bolded on my machined (Safari, OS X), so it sticks out like a sore thumb.
- xfs 13y agoSeconded. Chromium, Linux.
- PebblesRox 13y agoI noticed that too and was confused.
- mturmon 13y agoYes. Chrome 27, OS X. I spent a couple of seconds wondering if it was some typographical quirk inserted on purpose by someone.
- sauravc 13y agoDidn't notice anything. You got a good eye.
- andrewcooke 13y agoi saw it, but i don't understand why it's happening. the text doesn't have any special markup so why doesn't the font display it as "fi" (two letters) if it doesn't have the correct ligature? is it a bug in the font? edit: if i reduce the style to a single font, wf_segoe-ui_light, then it still appears that way. so i guess it's an issue with the font, not the browser.
- wazoox 13y agoYes, same thing, Firefox 20, Linux.
- mappu 13y agoThe ligature appears for me, but the font only seems to mess up in firefox. Which is weird, because i would have expected Firefox and IE to have the same behaviour owing to DirectWrite..(?) http://imgur.com/a/HNH5H http://imgur.com/a/HNH5H From top to bottom: Chrome 27/beta, Firefox 16, IE 10 on Windows 7 x64.
- rbanffy 13y agoHaving just recovered from a 72+ hour outage (about a dozen Hyper-V based hosts on our hosting provider) caused by a single (large) machine attached to a single (large) storage appliance, I think I'll pass on this idea of just scaling up instead of scaling horizontally. edit: "pass on" (thanks, marshray)
- deleted 13y ago[deleted]
- corresation 13y agoFalse dichotomy. This paper is primarily focused on programming models, and the truth that single scale-up machines grossly exceed all but the rarest of edge cases (what people call "big data" is often very little data). If you have a virtualization cluster of course you need redundancy for all components, and that is just basic fundamentals.
- marshray 13y ago"At 100 GB, scale-up still provides the best performance/$, but the 16-node cluster is close at 88% of the performance/$ of scale-up." If these message-passing graph algorithms are representative of a "bad fit" to the parallel map-reduce model, I'd say a 12% penalty is not a bad price to pay at all in return for all the benefits of the parallel cluster in other cases.
- sauravc 13y agoReminds of this article: http://blog.wavii.com/2011/12/29/your-mileage-may-vary/ http://blog.wavii.com/2011/12/29/your-mileage-may-vary/
- sauravc 13y agoSpecifically this image: http://wavii.files.wordpress.com/2011/12/hadoop_too_big.jpg?w=129&h=250 http://wavii.files.wordpress.com/2011/12/hadoop_too_big.jpg?...
- jamesaguilar 13y agoI'm surprised they didn't select any tasks that couldn't easily fit on a single machine (largest input set was <200GB). That said, if the scale up performance can be improved without compromising scale out, why not?
- saosebastiao 13y agoI've done a handful of big-memory workloads on the JVM, and I have never seen it not choke on Allocation/GC for anything above 300gb of memory. Does this paper address this limitation?
- wmf 13y agoI think they ran one JVM per core. If you're willing to spend money, large heaps are supposed to be solved by Azul.
- marshray 13y agoYes, section 3.4: "Heap size - By default each Hadoop map and reduce task is run in a JVM with a 200 MB heap within which they allocate buffers for in-memory data. When the buffers are full, data is spilled to storage, adding overheads. We note that 200 MB per task leaves substantial amounts of memory unused on modern servers. By increasing the heap size for each JVM (and hence the working memory for each task), we improve performance. However too large a heap size causes garbage collection overheads, and wastes memory that could be used for other purposes (such as a RAMdisk). For the scale-out configurations, we found the optimal heap size for each job through trial and error. For the scale-up configuration we set a heap size of 4 GB per mapper/reducer task (where the maximum number of tasks is set to the number of processors) for all jobs."