4 ms·
Won't this have serious implications for Spark, Hadoop, and other frameworks that distribute workloads across multiple JVM instances?
by sewercake 8y ago
Won't this have serious implications for Spark, Hadoop, and other frameworks that distribute workloads across multiple JVM instances?
- zeroxfe 8y agoAFAIK, Hadoop uses protocol buffers for message passing, not Java serialization.
- yzmtf2008 8y agoHadoop uses Avro actually, but the point still stands :)
- erik_seaberg 8y agoHadoop expects your keys and values to implement Writable and serialize themselves (a lot of these are actually hand-written expecting the instance to get reused for each input tuple). There's optional and fairly clumsy glue that makes Avro work in a key or value.
- sitkack 8y agoLast time I dived into details, the communication protocol for HDFS was serialized java objects using the built in serialization mechanism. Edit, they switched to Protocol Buffers in 0.23 https://wiki.apache.org/hadoop/HadoopRpc https://wiki.apache.org/hadoop/HadoopRpc
- slaymaker1907 8y agoSpark uses its own serialization system which while similar to the built in serialization, is designed to give much better performance.
- Scea91 8y agoOnly partly. Spark's RDD API which is still used quite heavily requires external serializer. Most of the time you would use Kryo, but it does not fork for objects larger than 2 GB or sometimes for custom classes that are not explicitly registered as Kryo serializable. In these cases, the changes would break existing code. However, who knows when will Oracle decide to remove the Java serialization API. I expect it will take a few years and the situation will be different on the Spark side then.
- kgoutham93 8y agoI think wildfly also uses serialization to store session state when scaling horizontally.
- imtringued 8y agoThe JPA spec requires entity IDs to be serializable and tomcat requires session variables to be serializable too. This is going to break a lot of code.
- anonacct37 8y agoSo to be clear most "code" shipping frameworks in jvm land use jar files. Big data systems are not using native serialization for either code or data.
- saryant 8y agoThat's not quite true. Both Akka and Flink use Java serialization in many key aspects.