Apache Spark: The Next Big Data Thing?
blog.mikiobraun.de
blog.mikiobraun.de
Spark's primitives (especially the RDD abstraction) make it feel like a DSL for distributed ML algorithms - for example, check out the implementation of distributed alternating least squares matrix factorization at [2] - ~100 lines of Scala code once stripped of the Java interoperability boilerplate.
On the weekend, I took a crack at implementing some distributed ML algorithms using Alternating Direction of Multipliers (ADMM) in Spark, following [3]. In a few hours and ~400 lines of Scala code, it was possible to implement distributed versions of L^1 regularized logistic regression, ridge regression, and SVMs [4].
It's a very impressive framework, and very much empowering.
[1]: http://spark.incubator.apache.org/docs/latest/mllib-guide.ht...
[2]: https://github.com/apache/incubator-spark/blob/fdaabdc673875...
[3]: http://www.stanford.edu/~boyd/papers/pdf/admm_distr_stats.pd...
[4]: https://github.com/ajtulloch/admmlrspark/tree/master/src/mai...
Not sure why the author thinks there's no built-in support for iterations. The support is there, native scala/java. But you might have to collect to the driver/master to check for convergence, for example. Their simple logistic regression example doesn't check for convergence, but hard codes the number of iterations. http://spark.incubator.apache.org/examples.html
Some quirks that I've learned. When doing groupbys or other join operations, it helps to specify the number of partitions -- actually this is true in general, but more so for join-based operations, otherwise you can run into memory/gc issues. The equivalent is of course, setting the number of reducers to something which ensure you won't run out of heap space.
Because spark serializes the closures which transform your data, if you don't cache (at least via disk, using persist), then when iterating over an RDD or re-using it, you'll just waste cycles.
Beyond that, as an experienced scala developer, I've found that Spark feels incredibly natural. And being able to run things on the repl cannot be appreciated enough.
Specifically, "faster batch + in-memory" is basically a simple patch on the batch mode problem but does not really address continuous flow real-time models. At large scales it is quite difficult to get robustly efficient behavior out of a parallel system that is batching things on 5 second intervals if the data flow is anything but trivial.
For geospatial and graph analysis, Spark appears to retain the same limitation of Hadoop in that it cannot deal with data models and operations without an a priori optimal partitioning function. Static hash and range partitioning won't cut it, particularly if the streaming data sources are actually real-time. The ability to generate uniform partitioning of complex data models with inherently unpredictable data distributions is critical to parallelizing some important analysis types but there is no obvious support for such mechanisms in Spark.
One could look at Spark as Hadoop done right, more or less.
Here are the papers for GraphLab and GraphX...
GraphLab: http://graphlab.org/home/publications/
GraphX: A Resilient Distributed Graph System on Spark (https://amplab.cs.berkeley.edu/publication/graphx-grades/)
See also "Introduction to GraphX - Presented by Joseph Gonzalez, Reynold Xin - UC Berkeley AmpLab 2013" (http://www.youtube.com/watch?v=mKEn9C5bRck)
My first impressions are that there seems to be a fair amount of abstraction leakage. Some things that compile and look valid fail at runtime - e.g. referring to other RDDs from within a filter predicate or map function. Other More complex jobs cause the nodes to fail and it gets into infinite loops of restarting the nodes, replaying the job, and them dieing again.
I hope once I get a better understanding of what is going on underneath I will understand what is going on here.
More complex jobs cause the nodes to fail and it gets
into infinite loops of restarting the nodes, replaying
the job, and them dieing again.
This is probably a bug in your code. Debugging these cluster applications does take some getting used to. You'll want to look at the stderr output of the failing executor, and you'll probably see that it's dying due to some kind of exception. You can do this by visiting port 8080 of the master node over HTTP, i.e. http://mymaster:8080. Feel free to email the Spark users list if you have any questions: https://spark.incubator.apache.org/mailing-lists.htmlIt's true that you do have to understand the programming model and some details of how it's implemented to use Spark effectively. However, any abstraction that was "pure" and perfectly non-leaky would necessarily sacrifice some performance and transparency to achieve that goal. Spark aims to be both high-level and high-performance.
Full disclosure: I'm on the Spark team at the UC Berkeley AMPLab.
[1] http://ampcamp.berkeley.edu/big-data-mini-course/
[2] http://ampcamp.berkeley.edu/big-data-mini-course/launching-a...
Feel free to suggest new features, or contribute to the project yourself.