Arrow version compatibility


I believe that Arrow Plasma is no longer maintained under the Apache Arrow project now. So, I am wondering what is Arrow version supported by Ray currently and is the Plasma project managed internally by the Ray project now?

Ray moved onto the cloudpickle due to some robustness in generic python serialization. That says, Ray currently doesn’t use Arrow for serialization. We might be using arrow for implementing the higher level abstraction though. Plasma store has been back-ported to Ray in order to add new features (such as object spilling) and performance optimization.

So, is Ray using Arrow format still? If so, is there a compatibility version compatibility?

It doesn’t use the Arrow format for serialization I believe. cc @suquark in case I am wrong

Yes, we are not using Arrow format for serialization now.