4 ms·
Announcing Spark 1.6
- minimaxir 11y agoIt 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.
- rxin 11y agoI 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.
- mziel 11y agoTo 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.
- rxin 11y agoThis 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.
- mziel 11y agoThat sounds very interesting. Is there any ticket, where I can follow the progress on the batch columnar format you mentioned? Btw, I was critical about the issue above, but I do love Spark, using it on a daily basis. :)
- rxin 11y agoI just created a JIRA ticket tracking this: https://issues.apache.org/jira/browse/SPARK-12635 https://issues.apache.org/jira/browse/SPARK-12635 Thanks for the reminder!
- minimaxir 11y agoThat'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.
- xjlin0 11y agoYes that's right. Team: could you make more API available through PySPark, please?
- mziel 11y agoNew statistics/machine learning algorithms are always welcome, but the big plus for productionizing is ML pipelines persistence. That said, during Spark Summit Databricks guys themselves were most excited about Dataset API. Looking forward to giving it a try.
- kod 11y agoSo the detailed post on datasets at https://databricks.com/blog/2016/01/04/introducing-spark-datasets.html https://databricks.com/blog/2016/01/04/introducing-spark-dat... uses groupBy 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?
- rxin 11y agoDespite 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%).
- gshayban 11y agoMany 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
- IvanVergiliev 11y agoReduce 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.
- dastbe 11y agodatabricks has a page that describes the pitfalls: https://databricks.gitbooks.io/databricks-spark-knowledge-base/content/best_practices/prefer_reducebykey_over_groupbykey.html https://databricks.gitbooks.io/databricks-spark-knowledge-ba... I don't know if the OutOfMemory exception can still occur in recent versions of Spark, but the performance impact of groupByKey is very real.
- mb22 11y agoWe'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.
- mark_l_watson 11y agoGreat 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.
- _laf 11y agoThis post: https://databricks.com/blog/2015/07/30/diving-into-spark-streamings-execution-model.html https://databricks.com/blog/2015/07/30/diving-into-spark-str... 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.
- peterstjohn 11y agoBackpressure is in 1.5+ by setting spark.streaming.backpressure.enabled=true (https://spark.apache.org/docs/latest/streaming-programming-guide.html#requirements https://spark.apache.org/docs/latest/streaming-programming-g...). Like you, I'm looking forward to the out-of-order data support.
- huula 11y agoLove Spark! great works, folks!
- ranjeet_hacker 11y agoExcited about dataset api and ML pipeline.
- DannoHung 11y agoAny support for nearest neighbor joins? They are very important for aligning events in time series data sets.