3 ms·
Most "big data" distributed compute frameworks that come to mind are written in a JVM language, so the focus on Rust is interesting. So then, would Rust be bet
by cozos 7y ago
Most "big data" distributed compute frameworks that come to mind are written in a JVM language, so the focus on Rust is interesting.
So then, would Rust be better than a JVM language for a distributed compute framework like Apache Spark?
Based on what others said in this thread, these are the primary arguments for Rust:
1. JVM GC overhead
2. JVM GC pauses
3. JVM memory overhead.
4. Native code (i.e. Rust) has better raw performance than a JVM language
My take on it:
(1) I believe Spark basically wrote its own memory management layer with Unsafe that let's it bypass the GC [0], so for Dataframe/SQL we might be ok here. Hopefully value types are coming to Java/Scala soon.
(2) Majority of Apache Spark use-cases are batch right? In this case who cares about a little stop-the-world pause here and there, as long as we're optimizing the GC for throughput. I recognize that streaming is also a thing, so maybe a non-GC language like Rust is better suited for latency sensitive streaming workloads. Perhaps the Shenandoah GC would be of help here.
(3) What's the memory overhead of a JVM process, 100-200 MB? That doesn't seem too bad to me when clusters these days have terabytes of memory.
(4) I wonder how much of an impact performance improvements from Rust will have over Spark's optimized code generation [1], which basically converts your code into array loops that utilize cache locality, loop unrolling, and simd. I imagine that most of the gains to be had from a Rust rewrite would come from these "bare metal' techniques, so it might the case that Spark already has that going for it...
Having said that, I can't think of any reasons why a compute engine on Rust is a bad idea. Developer productivity and ecosystem perhaps?
[0] https://databricks.com/blog/2015/04/28/project-tungsten-bringing-spark-closer-to-bare-metal.html https://databricks.com/blog/2015/04/28/project-tungsten-brin...
[1] https://databricks.com/blog/2016/05/23/apache-spark-as-a-compiler-joining-a-billion-rows-per-second-on-a-laptop.html https://databricks.com/blog/2016/05/23/apache-spark-as-a-com...
- andygrove 7y agoSome good points. Some incredible engineering has gone into Spark to work around the fact that it runs on the JVM. Memory overhead of Spark particularly (not just JVM) is very high. In some cases close to 100x more memory than equivalent query execution with DataFusion [0]. Also you might be interested to see my original blog post with some of my thoughts on this [1]. [0] https://andygrove.io/2019/04/datafusion-0.13.0-benchmarks/ https://andygrove.io/2019/04/datafusion-0.13.0-benchmarks/ [1] https://andygrove.io/2018/01/rust-is-for-big-data/ https://andygrove.io/2018/01/rust-is-for-big-data/
- cozos 7y agoInsightful blog posts! IMO a better memory comparison would be between a Spark executor and a DataFusion ... container I guess (i.e. graphing query time vs spark.executor.memory). This would give you a better idea of memory TCO on a cluster.
- kod 7y agoOne of the creators of Spark apparently thinks Rust is worth pursuing: https://www.weld.rs/ https://www.weld.rs/
- ADefenestrator 7y agoI'm not too familiar with JVM internals or Spark, but I know with Cassandra at least there's a cost to off-heap memory. You gain in GC, but eventually you have to move that data in and out of the JVM's heap. Even for batch processing, long GCs can be bad. It's not just processing the batch that stops, but the whole world. Anything trying to listen for more data, keep track of time elapsed, etc is going to run across more problems. It can also expose some race conditions that would normally be so unlikely that you'd never hit them. Static JVM memory overhead isn't bad at all, probably even under your 100-200MB guess. The lack of compact primitives adds some additional proportional overhead, but the biggest factor is just the extra "slack" space needed for good GC performance. Depending on the circumstances and requirements that could be 30-400%.