We've been testing 1.6 since before release, specifically SparkSQL and there are some big performance improvements in this release.
We're putting together a 3rd party benchmark I'll post to HN when we are done.
It sounds like many of the improvements are not available on PySpark yet, which is disappointing. (Notes say feature parity for MLIB, so I'll look into that) However, the notes sound promising.
To be honest PySpark and SparkR are always going to be 2nd category citizens (because of the serialization/pickling between the two environments). Databricks shows nice graphs, saying they are equivalent for DataFrames, however those count only for built-in functions that basically translate code into execution plan for Catalyst. For anything bespoke (UDFs, custom Transformers/Estimators) you're better off using Scala.
This is true when you compare the performance vs Java/Scala, but if you compare it with other tools that are native in Python, it is not really much worse. For examples, Pandas operations that use custom UDFs are substantially slower than the native operations.
That said, as part of Project Tungsten, we have some ideas about a batch columnar format that can be shared by Python, R, Scala and Java, and that should be able to eliminate most of the inefficiency in serialization across process boundaries.
That's fair; mostly I'd prefer keeping everything in Python due to the other data implications outside of Spark (e.g bespoke data cleanup with Python syntax) and I'd prefer to keep the stack smaller.
I think most of the improvements are indeed available in Python, including better memory management, improved Parquet performance, and the many algorithms. The main two things that are not yet available are the Dataset API and the streaming state management stuff.
I'm pretty sure based on previous comments you've made that groupBy was one of the things you'd rather eliminate from the RDD api, because of the performance impact compared to reduceByKey (which is almost always what people should be using instead).
Are you at all worried about confusion if groupBy now performs ok on datasets, but not on rdds?
Despite our attempts at warning people, a lot of users still use groupByKey in RDDs. Hopefully over time this won't be a problem as the engine should be able to figure out more intelligently and do the proper rewrite (of course, we won't be able to do it 100%).
Many people blindly point to the docs to say "don't use groupBy, prefer reduce because it's faster..." Are there better examples that illustrate the fundamental differences between the two operations? Surely there is still a need for both operations
Reduce can perform reductions on locally on each machine before shuffling the data. This decreases the memory as well as the network overhead.
If you need all the elements for a given key - e.g. to display them to a user or save them to a DB, perhaps you should use groupBy. If you're going to perform some form of a reduce after that though, it's likely sub-optimal.
In the section titled "Future Directions for Spark Streaming" there is a paragraph about Event time and out-of-order data and Backpressure. This would blow my mind to be able to use; this is a real pain currently.
Great news, Spark is awesome. Only problem is that I now need to review my Spark material for an eBook that I released a month ago and update the examples to work on version 1.6, if required.