7 ms·
Moving product recommendations from Hadoop to Redshift saves us time and money
- rpedela 12y agoI chuckled when I read "We have a legacy data warehouse based in Hive and Elastic MapReduce, with backing data stores in S3.". I guess things have come full circle. It wasn't long ago that a relational database solution would have been "legacy".
- hyperliner 12y agoBut "legacy" is not only when technology gets old, but also when solutions get old. Maybe it was just a bad solution and they are moving it to a new solution, not necessarily a new platform.
- zenjzen 12y agoI did too, mainly because S3 + Hadoop isn't going to provide the locality and speed that HDFS would. I don't think it's a fair comparison.
- ergest 12y agoIs it me or are people switching to non-relational data warehouse architectures simply because it's en vogue? How many companies do you know that have enough data where a non-relational DW would actually make sense? I wonder, have we really pushed relational databases to their breaking point?
- bitL 12y agoIf you optimize for latency relational databases won't cut it.
- clubhi 12y agoThis might be the most inaccurate statement on the internet.
- bitL 12y agoOK, I should have written "if you optimize for read-access latency". Better?
- jpat 12y agoI'm the author of the article. At Monetate, we've chosen our data warehouses to maximize throughput, rather than minimize latency. That's where something like Redshift really shines, it's great a large bulk ingests and running large queries relatively quickly, but awful at running lots of small queries quickly. On our busiest day last year, we ingested over a quarter billion page views across all of our clients' websites. I'm sure someone has made MySQL scale to that volume, but for us Redshift has been working great for a relatively low price point.
- bitL 12y agoThank you for sharing your experience! It's always inspiring to read well-written articles as is yours!
- meritt 12y agoI've looked at and avoided doing anything serious with hdfs/mr for 6 years now. I'm glad some people are starting to realize that re-processing your entire dataset every single time you want to do something isn't very efficient. I'm still waiting for lightbulb moment where the usefulness of it really makes sense to me. Can anyone point me to a book or blog that discusses good uses of hadoop/map-reduce?
- jgrahamc 12y agoI'm waiting for the day people realize that materialized views in databases are awesome and decide to incorporate them into a framework.
- meritt 12y agoAt least if you're using Oracle they are, as it supports auto-refreshing. Postgres has only had them since 9.3 (and have to be manually refreshed). Meanwhile MySQL is still struggling with regular views.
- zenjzen 12y agoSimplistically speaking, you don't always have to do table scans. I run into this every day: "Let's use Hadoop and keep doing full table scans! It's scalable! We just add more machines!" Yeah, except continuing to scan all of your growing data each time you need it is inherently unscalable. :(
- personZ 12y agoI wonder, have we really pushed relational databases to their breaking point? The primary limitation of relational databases were traditionally that without expert level optimizations (which, realistically, data-focused organizations should have. But they very seldom do, especially in the start-up space), many queries would generate large numbers of effectively random IO. When you're rolling with magnetic drives, each drive offers maybe 60-150 IOPS, so this quickly becomes an enormous scaling problem. A large storage array offered maybe 2000 IOPS. Scaling becomes entirely about scaling IOPS, as CPU is seldom a limitation in databases. Add that many firms were starting on EC2 which not only gave you minimal memory, it offered absolutely miserable IOPS performance. Digg famously, and disastrously, solved this problem by essentially "denormalizing" every bit of data, enormously exploding the raw data they stored, but allowing for individual queries to be entirely localized, often served in a single, large IO: Instead of looking up all of your friends and finding the things they dug, the system would push every bit of data proactively to containers for every possible user. This is the model promoted by many advocates of alternative storage (e.g the advantage of MongoDb is always the "pull a single giant data bag versus pulling it together from various places"). If Kevin Rose dug something, it would update the "things my friends liked" containers for 40,000 or so of his friends, rather than having those 40,000 users check on-demand to see what each of their friends liked. But they did that right when flash storage was coming into the mainstream. A technology that offers, on simple, inexpensive cards, 100s of thousands to millions of IOPS. Add that RAM has exploded, such that servers with 256GB of memory are very affordable (that was enough to put the entire universe of Digg's data in memory, where of course random IO is in the tens to hundreds of millions). So now we're at a situation where having non-duplicated, highly relational database is often the highest performance, outside of all of its other advantages, because it fits in memory, and fits on economical flash storage. It has completely flipped the equation. http://www.commitstrip.com/en/2014/06/03/the-problem-is-not-the-tool-itself/ http://www.commitstrip.com/en/2014/06/03/the-problem-is-not-...
- tomphoolery 12y agoThe Digg thing is interesting...they pretty much took the complete opposite approach of Reddit, who basically store everything in two big SQL tables.
- onion2k 12y agoIt's often cheaper to use the "wrong" architecture than optimising the right architecture. I know I could have a CouchDB datastore searching a few GB with an afternoon of work. I imagine I could get MySQL fast enough with a few days of optimisation. In terms of time, which is by far the biggest cost in most development, CouchDB is the better option.
- personZ 12y agoIn terms of time, which is by far the biggest cost in most development, CouchDB is the better option. For a single, one-off utility, sure. For anything that you ever planned for production, that would be crazy. Just to be clear, the mentality that onion proposes (at least from my interpretation, though I apologize if I'm misunderstanding), usually justified under a gross misinterpretation of the "premature optimization" warning, is exactly how disaster implementations that end up failing or requiring enormous amounts of engineering time to try to triage and bandage into something usable.
- hcho 12y agoAt least for the startup world, it's about prioritisation of concerns. Will that disaster implementation take me to my next(or first) round of funding? If yes, I'll happily go with it. After that, I can throw money at the problem.
- deleted 12y ago[deleted]
- deleted 12y ago[deleted]
- sbov 12y agoDisaster recovery is easy to put off forever because you don't need it until you do. When it happens it can also kill off your company. I've been involved in companies that went 14 years without a disaster. Another company I was involved with had 2 in a span of 2 months, each taking between 2 and 3 days to recover from. Regardless of whether I need it or not, I sleep better at night knowing a decent plan is in place. Which means I can perform better during the day.
- darkxanthos 12y agoA lot are. I used to want to just because it was cool. Now that I actually do though my main use case is sifting through hundreds of GB of unstructured data. I use Hadoop to get the data into a structured form that I can then load into Redshift. It's awesome knowing that just about anything I throw at Redshift, it can handle.
- crdb 12y agoWe switched to Redshift for our data warehouse because it was MUCH cheaper, declarative, allowed us to retain the relational model, and abstracted away most of the admin. Very happy so far. ~5,000 tables, 7-10TB all maintained by one guy in his spare time. The workbench options weren't amazing so we made this for our BI team: https://github.com/zalora/redsift/ https://github.com/zalora/redsift/
- zenjzen 12y agoI have another question: Why are people still doing joins in this day and age? Big data + joins = teh suck. I'm a big fan of compressed denormalized data.
- lsb 12y agoIt's also unclear how many rows they're trying to do this on, and at what frequencies; that's the crux of what turns this from a small-to-medium-data problem, which you can easily solve on a large box with 10 lines of code, to a big data problem, which requires completely different tooling
- jpat 12y agoIn my testing of this query, I ran it against a time range that included over 40 million purchase lines, and our configuration of Redshift returned the result in ~6 minutes. That was much quicker than our legacy EMR implementation. Currently, we update our product recommendations nightly. However, the speed up we see here from this reimplementation may allow us to update product recommendations more frequently.
- zatkin 12y agoWasn't Hadoop the first of it's kind in Big Data?
- crb 12y agoHadoop was the first (major) open-source implementation of Google's MapReduce framework. http://research.google.com/archive/mapreduce.html http://research.google.com/archive/mapreduce.html In terms of data warehousing and near-real-time query over Big Data, Google's framework for that is called "Dremel", http://research.google.com/pubs/pub36632.html http://research.google.com/pubs/pub36632.html Google offer Dremel as a service known as BigQuery.
- mattj 12y agoI've gone through a similar transition (hive to redshift) in a very large scale data environment. Raw Hadoop / cascading is still very useful for more complicated workflows, but redshift is so vastly superior to hive it's not even funny. I thought I would miss adding my own UDFs, but this hasn't been an issue at all. I'm under the impression presto is a similar improvement, but I haven't spent any time with it. One huge advantage of redshift over hive: you can connect with plain old Postgres libraries, so you can build redshift results into your admin interfaces, one off scripts, and anywhere else you're fine trading a few seconds of latency for extra data.
- growse 12y agoI'm not surprised, given that my experiences with Hive are that it's extremely quirky and hardly ever the fastest way to do anything. Given the fact that people seem to be falling over themselves to reinvent better solutions to the same sorts of problems in the Hadoop space (see: Impala, Shark), I don't think I'm alone on that.
- endersshadow 12y agoJust as a quick note: You can use Postgres libraries because Redshift is a slightly modified version Postgres 8.1 under the covers. In fact, almost all massively-parallel-processing (MPP) databases are Postgres under the covers (including Microsoft's PDW). It really speaks to how impressive Postgres is at scaling. Even old releases, like 8.1!
- mattj 12y agoYup! My experience with redshift has actually made me curious to try out Postgres (I've always used MySQL before this). The stricter SQL dialect was a little odd at first, but I think I've become more comfortable with it over a few months.
- endersshadow 12y agoPostgres is an amazing database, and has some great features. Sadly, you don't get the full power of it in Redshift, but man, are some of the datatypes and functions just so useful, especially in a warehousing environment!
- alaiacano 12y agoYou should use something like tf-idf to normalize your cooccurance counts for your recommender, otherwise you'll just end up recommending the most globally popular products.
- monstrado 12y agoThese type of articles baffle me, you're comparing a high-performance analytical database to a batch-orientated SQL engine. The whole point behind these query engines on Hadoop (Hive, Presto, Impala, etc) is to separate the database from the query engine. With these engines you can project schemas over raw data in its original form, without having to load it into a table. With Redshift, or other similar analytical databases, you're forced to define a schema, and then load the data in row by row...bulk inserts are very slow in comparison to Hadoop technologies. Regardless, Hive in general should nver be used for interactive analytics, that's not what it's intended for. Where Hive shines is when you can dump 250TB of raw text data into a folder and then run a SQL query to extract useful information out of it. The extracted data could then be loaded into a RDBMS like RedShift for real-time reporting. With all that being said, if you want to run SQL queries on data in Hadoop at the speeds of Redshift, you should have used Impala with Parquet, which is known to be even faster than Redshift in many cases, and is based on the same technology Google uses (Dremel and F1). The benefits of keeping your data in Hadoop are enormous, not every problem can be solved using SQL. The same data you're querying with Impala could actually be used to do machine learning using Spark or Mahout. Maybe you want to start indexing one of your tables into Solr to provide search capabilities on a subset of your columns to your users...or maybe you want to use Giraph or Sparks' GraphX to do parallel graph computation. The data never moves, there's still only ONE copy of that data in Hadoop, and you can bring any kind of workload to it.
- bbrunner 12y agoRedshift is an especially limited SQL engine considering it doesn't support UDFs. It is wicked fast, but what you get in speed you lose in flexibility. Current (well, February, but fairly current) benchmarks[0] place Impala and Shark (SQL on top of Spark) within grasp of Redshift while pulling data from disk and, for certain workloads, on par or faster than Redshift. This is without using a columnar file format. Impala is impressive technology, but it does require you to run dedicated Impala daemons as it doesn't use map reduce under the hood. Shark is especially interesting, however, because it is fast AND build on top of spark, so you can run raw Spark jobs, SQL queries, graph processing and ML all on the same cluster. Shark currently uses Hive to generate it's query plans, but the Spark project is working on implementing it's own SQL engine called Catalyst[1] that promises to be a significant improvement. [0] https://amplab.cs.berkeley.edu/benchmark/ https://amplab.cs.berkeley.edu/benchmark/ [1] https://spark-summit.org/talk/armbrust-catalyst-a-query-optimization-framework-for-spark-and-shark/ https://spark-summit.org/talk/armbrust-catalyst-a-query-opti...
- callesgg 12y agoI generally love sql and use it in all My stuff more or less. The one thing I fucking hate is how hard it is to get the database to execute querys in a efficient way on anything that is more than a simple join and select.
- NorthernDaemon 12y agoThose who want on-premise Redshift should try Actian Matrix. Redshift is just Amazon's version of Actian Matrix (formerly Paraccel)
- alanctgardner3 12y agoAnother commenter pointed this out, but what you're trying to compute is cosine similarity, in which case you're missing the normalizing part in the denominator (the product of the magnitude of both vectors). In other words, two items which both occur frequently will score higher than two items which occur infrequently, but which co-occur higher than usual. This leads to a tendency to over-recommend popular items. When you were on EMR, you could have used Mahout's distributed collaborative filtering, which has the benefits of being correct, and requiring zero coding. Wikipedia explains here: http://en.wikipedia.org/wiki/Cosine_similarity http://en.wikipedia.org/wiki/Cosine_similarity
- jpat 12y agoThanks for the tip, I'll look into this more.
- shamney 12y agoaside from the example in the original article, what sort of questions are these tool used to answer?