Fast Python Serialization with Ray and Apache Arrow
ray-project.github.io
ray-project.github.io
"Ray is a flexible, high-performance distributed execution framework." is certainly a nice tagline, but I'm pretty sure there are other projects in this domain - what's the USP ? who's using it and for what ?
More broadly it is useful for many parallel/distributed Python applications where low latency (~1ms) and high throughput of tasks are a requirement.
(1) Python single threaded performance: Here, most of the libraries we are using are implemented in C++ (like numpy, TensorFlow, Cython to speed up the code, etc.). Ray is orthogonal to that.
(2) Python parallel performance: Here Python is mostly problematic because of its lack of support for threading (the GIL is one problem here); we handle this problem by using multiple processes and shared memory throughout. Efficient serialization makes this feasible.
The core of Ray is implemented in C++, so performance is not an issue for that; also all of the serialization is implemented in C++.
Did they consider CapnProto?
In my limited experience with efficient serdes for analytics, CapnProto would meet 1 - 4.
Arrow is an in-memory format and CapnProto would afford the ondisc (memmap) heavy lifting with (very) little overhead. My CapnProto/LMDB datastructures allow me to write better than C++ code in Python with comparable runtime performance.
The only issue I have with CapnProto is getting it to run on Windows 7 x64, but it works out of the box on Linux.
My workaround for the failures of pickle is to use pathos. The API for the multiprocessing module is almost identical to the default multiprocessing module which makes it easy to use as a drop in replacement. Pathos uses Dill to serialise, which can serialise many more things, but you have to pay the price in performance. Although again, you can't swap out Dill for the serialiser of your choice.
There are too many bits that are no customisable like that (amongst other things).
Not the author, but I am assuming it breaks two of their requirements:
> 4. ... it should not require reading the entire serialized object
> 5. ... It should be language independent
I am curious about their thoughts about CapnProto though: https://news.ycombinator.com/item?id=15485082
Concerning dill, we have been using it for serializing function and classes (and then switched to cloudpickle because it supports some Python functionality better and the community around it is very active responsive); cloudpickle/dill are great in that they support a very wide variety of Python objects, especially concerning code (functions, lambdas, classes); it is less ideal for large data, because there is no zero copy mechanism, the format is not standardized and serialization/deserialization can be slow, sometimes it is slower than pickle. We do fall back to cloudpickle for objects we don't support like Python lambdas. We also use it to serialize class definitions, so the data associated to classes is serialized using the solution presented above and the code/methods are serialized using cloudpickle. This combines the advantages of both solutions.