5 ms·
Saying things like "Going beyond Hadoop" is very misleading. Virtually all of the Hadoop vendors out there, whether it be Cloudera, Hortonworks, MapR commercial
by monstrado 12y ago
Saying things like "Going beyond Hadoop" is very misleading. Virtually all of the Hadoop vendors out there, whether it be Cloudera, Hortonworks, MapR commercially support Spark as a computation framework for Hadoop, and some already have customers who are already using Spark on HDFS to power mission critical applications. People still pretend that Hadoop is some batch-orientated system with a distributed file system, but they couldn't be farther from the truth. Hadoop is a movement, an evolution of how data is to be analyzed in today's world.
The fact is, the vast majority of people who use spark (or will use spark), will rely on Hadoop for a lot of underlying technologies, such as, YARN (resource management) or HDFS (distributed filesystem / in-memory caching). Further, if you think Spark is somehow the "end all be all" computation framework, you're living in a fantasy world. The best part of Hadoop is that depending on your use case, you can bring a multitude of applications to your data, whether it's Spark, Tez, MapReduce, HBase, Impala, Drill, Presto, Tajo, Accumulo, ...the list goes on, and continues to evolve. Spark is in no way replacing Hadoop, it's only strengthening it.
- rs_atl 12y agoYou have a point, but it's also a fact that lots of people use "Hadoop" interchangeably with "MapReduce". And Spark can in fact replace the Hadoop infrastructure entirely, as it's not a component of that ecosystem. Just because the various Hadoop vendors also support Spark only validates the point that there's a need to "go beyond Hadoop".
- monstrado 12y agoJust because people use Hadoop and MapReduce interchangeably doesn't make it correct. I would love to hear how you think Spark can replace Hadoop, because that is an astonishingly inaccurate statement. Which part of Spark reliably distributes data? Which part of Spark handles enterprise level security? Which part of Spark can coordinate resources in multi-tenant environments? The answer is none of them, it relies on Hadoop for that.
- x0x0 12y agoum, you sound like a vendor. hadoop does mean map-reduce + hdfs as the common usage by the majority of devs + admins. Claiming hadoop is now some distribution of tools is fine, but that's simply not what the common usage is. It remains to be seen if yarn will carry the day or no; my suspicion is that many people are essentially going to be running spark on hdfs. I don't see much use for yarn unless you need to balance hadoop and yarn, and weren't their claims that yarn was going to support eg mpi style computation that didn't pan out?
- colin_mccabe 12y agoum, you sound like a vendor. hadoop does mean map-reduce + hdfs as the common usage by the majority of devs + admins. Claiming hadoop is now some distribution of tools is fine, but that's simply not what the common usage is. I am a Hadoop developer, and I can tell you that Hadoop does not mean "map-reduce + hdfs". That's also not what people are installing when they install Cloudera's distribution of Hadoop, Hortonworks' distribution of Hadoop, or even Intel's distribution of Hadoop (which is being discontinued in favor of adopting Cloudera's). This is more old information from 2008, being replayed as current. YARN even lives in the Hadoop source code repository, it's hard to get more "Hadoop" than that. Spark has its own repo, but it uses many classes from Hadoop like InputFormat, etc. It remains to be seen if yarn will carry the day or no; my suspicion is that many people are essentially going to be running spark on hdfs. I don't see much use for yarn unless you need to balance hadoop and yarn, and weren't their claims that yarn was going to support eg mpi style computation that didn't pan out? Much confusion. Much sadness. You run YARN (or its close competitor, Mesos) because you want to have multiple jobs going on in the same cluster at once. You need things like per-user queues, job control, reserving CPU and memory resources. The jobs going on at once may be multiple MapReduce jobs, or they may be multiple Spark jobs. Even Databricks, which employs many of the early Spark developers, doesn't ship a product that runs Spark in standalone mode. They run on Mesos.
- srean 12y agoYou seem to be very knowledgeable about Hadoop. Could you tell us a bit about the technical reasons why its performance is so god-awfully underwhelming. Is it a wrong choice of algorithms, internal data structure, architecture, design. I have used Google's implementation of mapreduce and then current version of Hadoop. Hadoop's performance was a piece of crap in comparison, and yes the the bottleneck was shuffle. With Java, I would have expected a 20~30% hit not a slowdown by 4~6 times and using roughly that many times more memory. Weirdly enough, I have heard that Python streaming on Hadoop outperforms the Java Hadoop, but I have not tried this myself. At that time UIUC's sector would outperform Hadoop too, and that was an university project maintained by a grad student. Finally has this changed with Hadoop 2.* And in case you have experienced Googles's mapreduce and Hadoop (which seems likely) I would like to know more.
- x0x0 12y agoYahoo's hadoop has a variety of issues. I'd rather not say company names, but this comes from working at a company that had hadoop in production by 2007 and last I heard was running 10+pb datastores with 10k+ hadoop cores. Yahoo's hadoop writes all output (from mappers and reducers) to local disks which are typically spinning rust. This is an anti-pattern: you run multiple mappers per box, typically 1.1-1.3 per core, and you have multiple jobs running simultaneously. You then perform a disk based merge sort, and then every reducer must connect to every mapper to get the data associated with its keys. The second you get a hot disk, whole jobs get choked because of the api invariants: all mappers complete before you can finish (or perhaps start) sort, all sorts complete before reducers, all reducers finish before the next mapper pass. Further, avoiding having every reducer have to have high throughput network to every mapper is key; the hadoop design makes it difficult to optimize over network topology like higher within-rack bandwidth. There are far better designs implemented: never write to local disk from the mappers, instead writing to a global dfs then performing the sort in the dfs. I worked for a company that took hadoop 0.17 and made it roughly 10x more performant (and not on toy datasets: on petabyte scale datasets) with roughly 10 engineer-years with this design. I think it's also fair to say that the quality of engineers on hadoop, particularly when they worked at yahoo, was... subpar. There was a ton of low hanging fruit. Also, yahoo's workload wasn't, or so I've heard, highly cpu bound. They where using hadoop at least in part just for hdfs capacity. So if you aren't cpu constrained, a lot of optimizations don't matter. Better companies have to run at say 95%+ cpu or their cfo won't authorize more boxes. You can also go after things like optimizing your compressors / decompressors, which are probably by far the hottest bits of code. Get icc, break out vtune, hire a really good optimizing dev to go after even small gains, etc. Finally, the hdfs design is pretty shitty, on two fronts: the namenode gets slow (choking everyone), and the disk-based format is missing lots of things to enable optimizations like certain types of append, or indexing. Indexing is incredibly powerful for things like join, enabling a map-side join instead of having to read all your data once, sort by join keys, then reread in your reducers. Another problem with hdfs design is capacity adding; hdfs has shitty rebalancing tools. The typical problem is this: hot data is new data. So if you add capacity to hdfs, your old capacity is say 80% full and your new capacity is 0% full so where does all new data go? To the new capacity. Which is quickly 100% saturated and choking your jobs. Properly rebalancing a dfs is a hard problem on it's own. Further, hdfs' use of triple replication (well, you can actually set the replication level) instead of better algorithms like erasure encoding does no favors to performance as well. edit: by the way, never writing to local disk and instead going from mappers straight to dfs then performing a rough/fine grained sort in the dfs is a lot like google's design. You do trade network i/o for disk i/o, but that's often easier to scale and, as mentioned, easier to optimize over your network topology. Finally, you can do multi-pass sorts partially locally by running sorters on your disk nodes.