3 ms·
Large Scale Distributed Deep Learning on Hadoop Clusters
- duggan 11y agoBoth this and Twitter Engineering's recent post[1] on HDFS make me wonder whether HDFS is something a team would reach for in 2015. I'm starting to read into the technologies in this area (i.e., I have not used much of the Hadoop stack yet), and I haven't found a fundamental reason why one would not base their batch processing on S3 (or your object store of choice). Existing software appears to make assumptions about the storage medium being a local hard drive. Much of the challenge of HDFS appears to be around scaling the NameNode, and provisioning capacity. S3 dispenses with these issues, and the only cost appears to be throughput. If software like Spark was modified to have a much more native approach to S3, could HDFS be dispensed with entirely? [1] https://blog.twitter.com/2015/hadoop-filesystem-at-twitter https://blog.twitter.com/2015/hadoop-filesystem-at-twitter
- TallGuyShort 11y agoI find Spark almost as easy to use on S3 as it is on HDFS. I mean you add your keys to the config (or use an instance that already has an IAM profile), and instead of typing hdfs://namenode:8020/ you type s3a://bucket/. Done. I haven't tried running Spark entirely without HDFS, so I don't know how well it works to have application log history, etc. go to S3. But certainly for some use cases HDFS doesn't bring anything to the table. If you need stronger consistency guarantees or the higher throughput is a big deal (e.g. very large MapReduce jobs or sequences of them), then there are still teams for which heavier use of HDFS makes total sense.
- seiji 11y agoHDFS has never been a good idea or implementation, it was just first and rode the fad wave to wide deployment (also see: mysql, php, javascript, ...). If it were stable and scalable to begin with, there wouldn't be 5 different billion dollar companies offering professional services and customized "fixed hadoop" distributions against the openawfulsource versions. (from ~3 years ago too: https://news.ycombinator.com/item?id=4298580 https://news.ycombinator.com/item?id=4298580)
- duggan 11y agoThis is largely what I'm trying to avoid. Looking at the Hadoop toolchain raises more than a few alarm bells, but I don't want to dismiss it too quickly since a lot of people have clearly invested a lot of time into it.
- bradhe 11y agoTraditionally, using HDFS solved the data locality problem. Networks have gotten MUCH faster, though.
- jbooth 11y agoAt every bigger company, teams wind up shipping a 128MB jar assembly to their 128MB data block, annihilating any locality gains.
- thinkmassive 11y agoS3 support is built into Apache Hadoop [1]. Azure Blob Storage & OpenStack Swift are also included (see Hadoop Compatible File Systems section in the left sidebar). A properly provisioned NameNode can support hundreds of millions of blocks. Like any high performance system, it won't do this out of the box and does require some tuning based on your use case. [1] http://hadoop.apache.org/docs/current/hadoop-aws/tools/hadoop-aws/index.html http://hadoop.apache.org/docs/current/hadoop-aws/tools/hadoo...
- duggan 11y ago> Like any high performance system, it won't do this out of the box and does require some tuning based on your use case. Absolutely, I'm just trying to discern whether it can reasonably be considered a "second choice" after S3 today. I think what I'm finding out here is that the choice is highly dependent on the workload. I suppose there isn't anything terribly surprising in that.
- gopalv 11y ago> I haven't found a fundamental reason why one would not base their batch processing on S3 (or your object store of choice) Naming Consistency for things like list after create. I had to debug through data loss issues with an object store[1] since they treat paths as a key lookup instead of being hierarchical - because rename("/tmp/x", "/tmp/y") has race conditions, unlike what's usually expected of a filesystem. I recommend reading on about S3mper (fi) from Netflix. [1] - https://cloud.google.com/storage/docs/gsutil/commands/mv?hl=en#non-atomic-operation https://cloud.google.com/storage/docs/gsutil/commands/mv?hl=...
- squarecog 11y agoS3 has a number of challenges of its own, for example, file access is quite expensive (performance-wise). The Netflix engineering team has a ton of experience running data processing "in the cloud", and have their share of war stories. The scaling challenges the Twitter blog post addressed (disclaimer: I work at Twitter, was one of the first Hadoop people at the company, am now doing slightly different things) happen at fairly extreme scale ranges. The NN design leaves much to be desired, but it does work just fine in the vast majority of cases. The scaling challenges we are talking about involve thousands of nodes and hundreds of petabytes of data. This is not what one normally designs for, even in a "big data" system. Take this into consideration when exploring alternatives. Do you have a solid way of managing a few hundred petabytes and a few hundred million objects in S3? Does GlusterFS work at that scale?