Presto: Interacting with petabytes of data at Facebook
facebook.com
facebook.com
Wow. As somebody who skimmed through the original Google Dremel paper and thought for a while about how one would go about implementing such an interactive system, that strikes me as an amazingly impressive timeline.
Do you think this would be a decent candidate for click stream analysis?
We have over a thousand employees in Facebook using it daily for anything you can imagine, so I recommend trying it and letting us know what you find useful and how it can be improved.
2. In data warehouse terms, would you fit it as a MOLAP, ROLAP or HOLAP kind of engine? I'm not sure whether all data is held in memory and whether aggregates are cached. Can you preprocess the dataset in batch mode (let's say overnight), repopulate aggregate caches for faster retrieval later on? (I mean that as something similar to SQL Server Analysis Services in MOLAP mode).
3. Can you compare it to Apache Drill ?
We have a pull request for partitioned hash aggregations, but currently all the groups must fit in memory limit specified by the "task.max-memory" configuration parameter (see http://prestodb.io/docs/current/installation/deployment.html for details).
Regarding approximate queries, we are working with the author of BlinkDB (http://blinkdb.org/) to add it to Presto. BlinkDB allows very fast approximate queries with bounded errors (an important requirement for statisticians / data scientists).
2) Presto is a traditional SQL engine, so it would be ROLAP. We don't yet have any support for building cubes or other specialized structures (though full materialized view support with rewrite is on the roadmap).
The Presto query engine is actually agnostic to the data source. Data is queried via pluggable connectors. We don't currently have any caching ready for production.
There is a very alpha quality native store that we plan to use soon for query acceleration. The idea is you create a materialized view against a Hive table which loads the data into the native store. The view then gets used transparently when you query the Hive table, so the user doesn't have to rewrite their queries or even know about it. All they see is that their query is 100x faster. (All this code is there today in the open source release but needs work to productionize it.)
We have a dashboard system at Facebook that uses Presto. For large tables, users typically setup pipelines in Hive that run daily to compute summary tables, then write the dashboard queries against these summary tables. In the future, we would like to be able to handle all of this within Presto as materialized views.
3) We're excited about the Drill project. They have some interesting ideas about integrating unstructured data processing (like arbitrary JSON documents) with standard SQL. However, last I looked they were still in early development, whereas Presto is in production at Facebook and is usable today. Please also see this comment: https://news.ycombinator.com/item?id=6684785
Additionally, in the long term, we want to enable Presto to be a completely standalone system that is not dependent on HDFS or the Hive metastore, while enabling next-generation features such as full transaction support, writable snapshots, tiered storage, etc.
In the case that you're younger, at Google, we offer something called Engineering Practicum internships. In this program, you get paired up with another freshman or sophomore, and work on an intern project together.
https://www.google.com/about/jobs/search/?src=Online/TOPs/NA...
Feel free to email me if you're interested in this. I'm doing some heavy data analysis and would be more than happy to host some interns next summer.
If you're a freshman or sophomore, apply!
For freshmen we recommend FBU: https://www.facebook.com/careers/university/fbu
Everyone else should check out regular internships or new grad positions: https://www.facebook.com/careers/university
As I understand it you just need to be going back to school at the end of your internship.
Unfortunately, we don't have any docs, so you'll just have to peruse the code. There's minimal server in the codebase to demonstrate usage of some of its features: https://github.com/airlift/airlift/tree/master/sample-server
We currently have a pull request open right now that will allow range predicates to be be pushed into the connectors. This will allow connectors the ability to implement range/skip scans.
The core presto engine does not take advantage of indexes right now.
We're also working on a example connector that can read from files/urls. We should have that code up soon.
One big advantage we have that speeds up development is the Facebook culture of moving fast and shipping often. Facebook employees are used to working with software as it's being built and refined. This kept us focused on the key subset of features that matter to users and getting close to realtime feedback from them. Our development cycle is typically one release and push to production per week.
Also, from the beginning Presto had to work with current Facebook infrastructure (100s of machines, 100s of petabytes), so we faced and solved all the associated scaling challenges up-front.
⒉⬠ In which DB do you store the text and media?
MySQL is used for storing textual user content like comments, etc. See https://www.facebook.com/MySQLatFacebook
Photos are stored using specialized systems: https://www.facebook.com/note.php?note_id=76191543919 http://www.stanford.edu/class/cs240/readings/haystack.pdf
And you can see some pictures of our new data center used for cold storage of older photos: http://readwrite.com/2013/10/16/facebook-prineville-cold-sto...
I'm missing something very basic here. The core idea of Presto seems to be to scan data from dumber systems -- Hive, non-relational stores, whatever -- and do SQL processing on it. So isn't its speed bounded by the speed of the underlying systems? I know you're working on various solutions to that, but what parts are actually in production today?
Also, very curious to know (from any Googlers browsing HN) if Dremel is still the state-of-the-art within Google, or if there is already a newer replacement.
Hadoop/Hive was not focused on speed but proving it is even possible to do queries reliably on such large datasets. Once this was achieved it immediately became clear the long waits for Hadoop/Hive batch jobs to finish made it impractical for many uses. Presto, Impala, Drill, RedShift, etc were all designed Primarily to address this problem and be much, much faster than Hadoop/Hive so that the data could be queried interactively.
All these new projects/products are in a very active competition to find the best way or ways to do this. You should compare RedShift to these other projects rather than Hadoop if speed is an issue for you. Each has it's pluses and minuses depending on the situation.
Impala, Drill, etc. avoid all those unnecessary reads and writes by implementing the querying logic directly, rather than by compiling to Map Reduce.
Shark [0] is an interesting counterpoint. It takes essentially the same approach as Hive but on Spark instead of Hadoop, and achieves similar or better performance than the more "direct" implementations.
For example as of today RedShift can hold a maximum of 256 terrabytes of compressed data while Facebook's Hadoop cluster was over 200 Petabytes in late 2012. RedShift only supports limited query and data types and a single index while Hadoop can theoretically handle arbitrary data processing. But if these constraints are acceptable then RedShift will likely be orders of magnitude faster in most cases.
Other projects/products will have different tradeoffs but they are almost always faster as this was almost always the primary goal.
Amazon Redshift enables you to start with as little
as a single 2TB XL node and scale up all the way to
a hundred 16TB 8XL nodes for 1.6PB of compressed user data.
from http://aws.amazon.com/redshift/features-and-benefits/The 16 node limit I was familiar with is not a hard limit. You can request more nodes.
http://aws.amazon.com/redshift/faqs/#0080
It would be interesting to see performance comparisons on these huge datasets. I would expect to see new and interesting problems at that scale.
[0] http://research.google.com/pubs/pub41344.html [1] http://shark.cs.berkeley.edu [2] http://blinkdb.org/
[0] http://research.microsoft.com/en-us/projects/dryad/ [1] http://research.microsoft.com/en-us/um/people/jrzhou/pub/Sco...
Redshift looks to be order-of-magnitude faster than Impala or Shark in all the test. Does this mean that once RedShift supports user-defined functions, there is no competing solution that is any match? (Unless you want to avoid using the cloud)
Rcfile or parquet would be a more interesting benchmark.
In the "What's next?" section, they say they want to re-do the Impala tests using Parquet, which is a columnar format based on the Dremel whitepaper (http://parquet.io/).
Parquet was a joint effort between Cloudera and Twitter, and now it's being developed by many other companies. You can use it with Hive, Pig, MapReduce, Cascading, Crunch and I think Apache Drill's first milestone has adopted it as a columnar format as well. Parquet also allows you to use your Avro or Thrift schema (soon Protobuffs) to write Parquet data, too.
It's a separate project in the ecosystem and has its own roadmap (https://github.com/Parquet/parquet-mr).
"Here from HackerNews? This was originally posted several months ago. Check back in two weeks for an updated benchmark including newer versions of Hive, Impala, and Shark."
I'm waiting in anticipation for that article!
If I have 5 clients at a time, and each one has a concept/idea every day that takes 20 seconds to explain to me (and isn't on the current iteration) then they are distracting me by not managing their time properly and calling me a dozen times a day to tell me about their thoughts. If it is a change to the current iteration, then it should have been discussed when we agreed on the current iteration's feature set. Either way, the call is not a result of well thought out time management.
As much as I like hearing new concepts and ideas, I also have to take attention away from a project that I'm working on in order to provide my full attention to the client calling me.
After the call is done, I also have to come back to the project at hand and hopefully I'm not working on something that requires that I retain a super complicated thought chain which may or may not have been lost in discussion with another client - especially in consideration that I'm not going to bill on other client project for the time that I've spent having been sidetracked and/or getting back to where I was before the call was made.
So, charging in $15 increments causes the client to actually manage their time with the same effectiveness that they would hope that I am managing mine.
Rounding down makes this more of a problem for me, not less of one. Now a 90 second call is at no charge, and I can get more than one of those in one hour - still at no charge - based on the suggested agreement.
- Standard ANSI SQL syntax, including all the basic features you'd expect from a SQL engine (aggregations, joins, etc) and other more advanced features like analytic window functions, common table expressions (WITH), approximate distinct counts and percentiles.
- It's extensible. The open source code base includes a connector for Hive, but we also have some custom connectors for internal data stores at Facebook. We're working on a connector for HBase, too.
- In comparison to Hive, it's very fast and efficient. For our workloads it's at least 10x more CPU-efficient. Airbnb is using it and has had a similar experience.
- Most importantly, Presto has been battle-tested. It's been in production at Facebook since January and it's used by 1,000 employees every day running 30,000 queries daily. We've hit every edge case you can imagine.
In general, you should be able to write your query in the simplest and most readable way, and Presto should execute it efficiently. We already have the start of an advanced optimizer that supports equality inference, full predicate move-around, etc. This means that you don't need to write redundant predicates everywhere as is required with some query engines.
Also, if you are familiar with PostgreSQL, you should feel right at home using Presto. When making decisions for things not covered by ANSI SQL, the first thing we look at is "what does PostgreSQL do".
That said, earlier this year, during a hackathon, we build a prototype connector that could split a query and push down the relevant parts to a distributed database that supports simple aggregations. It would be more work to clean this up and integrate, so if a lot of people are interested in this we can prioritize that.
The Hive data model supports complex types including arrays, maps and structs (nested tables), and those are used liberally. We also have a fair amount of data stored as JSON inside of string columns in Hive tables, though much of this is actually structured and could be better modeled using Hive's complex types.
Presto currently has limited support for Hive complex types. All complex types are converted to JSON at query time and can be accessed using the JSON functions: http://prestodb.io/docs/current/functions/json.html
— Lucius Fox, The Dark Knight