Spark as a Compiler: Joining a Billion Rows per Second on a Laptop
databricks.com
databricks.com
For me, is particularly interesting reading the Spark achievements. I was part of the similar Hive effort (the Stinger initiative [1]) and I contributed some parts of the Hive vectorized execution [2]. I see the same solution that applied to Hive now applies to Spark:
- move to a columnar, highly compressed storage format (Parquet, for Hive it was ORC)
- implement a vectorized execution engine
- code generation instead of plan interpretation. This is particularly interesting for me because for Hive this was discussed then and actually not adopted (ORC and vectorized execution had, justifiably, bigger priority).
Looking at the numbers presented in OP, it looks very nice. Aggregates, Filters, Sort, Scan ('decoding') show big improvement (I would expected these, is exactly what vectorized execution is best at). I like that Hash-Join also shows significant improvement, is obvious their implementation is better than the HIVE-4850 I did, of which I'm not too proud. The SM/SMB join is not affected, no surprise there.
I would like to see a separation of how much of the improvement comes from vectorization vs. how much from code generation. I get the feeling that the way they did it these cannot be separated. I think there is no vectorized plan/operators to compare against the code generation, they implemented both simultaneously. I'm speculating, but I guess the new whole-stage code generation it generates vectorized code, so there is no vectorized execution w/o code generation.
All in all, congrats to the DataBricks team. This will have a big impact.
[0] http://oai.cwi.nl/oai/asset/16497/16497B.pdf [1] http://hortonworks.com/blog/100x-faster-hive/ [2] https://issues.apache.org/jira/browse/HIVE-4160
/**
* WholeStageCodegen compile a subtree of plans that support codegen together into single Java
* function.
*
* Here is the call graph of to generate Java source (plan A support codegen, but plan B does not):
*
* WholeStageCodegen Plan A FakeInput Plan B
* =========================================================================
*
* -> execute()
* |
* doExecute() ---------> inputRDDs() -------> inputRDDs() ------> execute()
* |
* +-----------------> produce()
* |
* doProduce() -------> produce()
* |
* doProduce()
* |
* doConsume() <--------- consume()
* |
* doConsume() <-------- consume()
*
* SparkPlan A should override doProduce() and doConsume().
*
* doCodeGen() will create a CodeGenContext, which will hold a list of variables for input,
* used to generated code for BoundReference.
*/
[0] https://issues.apache.org/jira/browse/SPARK-12795
[1] https://github.com/apache/spark/blob/0e70fd61b4bc92bd744fc44...disclaimer: I worked on whole-stage codegen and other performance related stuff.
True, since MonetDB came of out of the Netherlands around 1993, but KDB+ came out in 1998, well over 8 years before C-Store.
K4 the language used in KDB is 230 times faster than Spark/shark and uses 0.2GB of RAM vs. 50GB of RAM for Spark/shark, yet no mention in the article. It seems a strange omission for such a sensational sounding title [1]. I don't understand why big data startups don't try and remake the success of KDB instead of reinventing bits and pieces of the same tech with a result in slower DB operations and more RAM usage.
I beat a databricks stack (hand tuned by one of the authors of big-DF) running on a large cluster using one jd node by factors of "a lot." And K is faster than J.
Hadoop is also about distributed computing, but if the process requires 250x as much RAM, you're going to have a very high TCO regardless of how cheap RAM or servers can be.
J is an opensource APL-derived language, but different in some important ways, but it is not as fast as KDB+/Q/K. J is truly an array-based language, whereas K is list based, so it has more in common with Lisp in that singular respect.
Kerf is a new language being worked on by one of the creators of the Kona language, an opensource version of the K3 language [1].
It seems these type of articles ignore non-opensource solutions even if they are more efficient in many ways including TCO. How many man-years need to be spent to try and duplicate an existing solution, and still not be a better solution?
Since it has not been mentioned yet, let me point out that Barry Jay's Fish should be a part of this conversation. Sadly its quite hard to find. I don't know why those pages have disappeared.
http://quant.stackexchange.com/questions/3156/is-there-any-t...
$25k for 2 core setup is a little crazy. They should definitely add a free tier if they want more adoption and more people investigating their tech. Spark is open source and free.
As I replied below, I am just surprised that similar technology has not been created by the opensource community that is closer in performance to KDB+.
$25k may seem absurd to you, but if it delivers results 230x faster using 1/250th of the RAM, decision makers in business see the TCO (total cost of ownership) as being worth it. It all depends on your needs and how they are met.
Jd is the commercial database for the opensource language J, that is also fast, and I believe it costs less than KDB, so there are other faster options.
If free as in opensource is your criteria, and not actual cost-benefits analysis, you only have those free choices, but in business time is money, so the 'free' option would be more costly.
I think a prime example of this is I was using some very basic windowing functions and due to the data shuffling (The data wasn't naturally partitioned) it seemed to be very buggy and not very clear why stuff was failing. I ended up rewriting the same section of code using hive and it had both better performance and didn't seem to have any odd failures. I realize this stuff will improve but I'm still skeptical.
I did hit issues w/ multiple joins and shuffling though. Have you not hit issues w/ shuffling?
I was using Spark 1.5.1 for the record.
That is a good metric for a project.
But let's be clear here Spark is not a 'real' open source project. It's a dictatorship run by Databricks (in the nicest possible way). There are 414 outstanding pull requests some of which came from the likes of Intel which added subquery support over a year ago. Never got merged.
The best metric for any open source project is how easily can you contribute and improve the product.
I showed up on the mailing list, said "hey, here's some stuff that was useful for me, let me know if it's useful for anyone else." A committer said, "cool, send a PR". Couple of rounds of code review and it was merged. Later on, bigger features required a lot more discussion, but as long as I was willing to follow up, I got my changes in.
Spark has on the order of 1,000 contributors. I'd call that a "real" open source project.
ODBC works via a JDBC bridge which though they are often commercial and vary in quality.
[0] https://www.amazon.com/Transaction-Processing-Concepts-Techn...
I also loved the MOOC by Prof. Jennifer Widom[1]
[0] http://pages.cs.wisc.edu/~anhai/courses/764-sp07-anhai/datam...
https://github.com/postgres/postgres/tree/master/src/backend...
Sqlite is another good one to read, and play around with.
http://wwwlgis.informatik.uni-kl.de/archiv/wwwdvs.informatik...
You could easily turn it into a database by managing a set of flat files.
Depends on the enterprise database, is my point.
(I'm just curious...used to work the internal SAP HANA dev team)