Update serialization doc (#8381)

* update serialization doc
This commit is contained in:
Siyuan (Ryans) Zhuang
2020-05-12 16:47:00 -07:00
committed by GitHub
parent 96f4d82cc3
commit ab278071ac
+26 -28
View File
@@ -13,33 +13,24 @@ Each node has its own object store. When data is put into the object store, it d
Overview
--------
Objects that are serialized for transfer among Ray processes go through three stages:
**1. Serialize directly**: Below is the set of Python objects that Ray can serialize using ``memcpy``:
1. Primitive types: ints, floats, longs, bools, strings, unicode, and numpy arrays.
2. Any list, dictionary, or tuple whose elements can be serialized by Ray.
**2. ``__dict__`` serialization**: If a direct usage is not possible, Ray will recursively extract the objects ``__dict__`` and serialize that directly. This behavior is not correct in all cases.
**3. Cloudpickle**: Ray falls back to ``cloudpickle`` as a final attempt for serialization. This may be slow.
Ray has decided to use a customed `Pickle protocol version 5 <https://www.python.org/dev/peps/pep-0574/>`_ backport to replace the original PyArrow serializer. This gets rid of several previous limitations (e.g. cannot serialize recursive objects).
Ray is currently compatible with Pickle protocol version 5, while Ray supports serialization of a wilder range of objects (e.g. lambda & nested functions, dynamic classes) with the support of cloudpickle.
Numpy Arrays
------------
Ray optimizes for numpy arrays by using the `Apache Arrow`_ data format.
Ray optimizes for numpy arrays by using Pickle protocol 5 with out-of-band data.
The numpy array is stored as a read-only object, and all Ray workers on the same node can read the numpy array in the object store without copying (zero-copy reads). Each numpy array object in the worker process holds a pointer to the relevant array held in shared memory. Any writes to the read-only object will require the user to first copy it into the local process memory.
.. tip:: You can often avoid serialization issues by using only native types (e.g., numpy arrays or lists/dicts of numpy arrays and other primitive types), or by using Actors hold objects that cannot be serialized.
Serialization notes and limitations
-----------------------------------
Serialization notes
-------------------
- Ray currently handles certain patterns incorrectly, according to Python
semantics. For example, a list that contains two copies of the same list will
be serialized as if the two lists were distinct.
- Ray is currently using Pickle protocol version 5. The default pickle protocol used by most python distributions is protocol 3. Protocol 4 & 5 are more efficient than protocol 3 for larger objects.
- Ray may create extra copies of simple native objects (e.g. list, and this is also the default behavior of Pickle Protocol 4 & 5), but recursive objects are treated carefully without any issues:
.. code-block:: python
@@ -48,28 +39,27 @@ Serialization notes and limitations
l3 = ray.get(ray.put(l2))
assert l2[0] is l2[1]
assert not l3[0] is l3[1]
- For reasons similar to the above example, we also do not currently handle
objects that recursively contain themselves (this may be common in graph-like
data structures).
.. code-block:: python
assert l3[0] is l3[1] # will raise AssertionError for protocol 4 & 5, but not protocol 3
l = []
l.append(l)
# Try to put this list that recursively contains itself in the object store.
ray.put(l)
ray.put(l) # ok
This will throw an exception with a message like the following.
- For non-native objects, Ray will always keep a single copy even it is referred multiple times in an object:
.. code-block:: bash
.. code-block:: python
This object exceeds the maximum recursion depth. It may contain itself recursively.
import numpy as np
obj = [np.zeros(42)] * 99
l = ray.get(ray.put(obj))
assert l[0] is l[1] # no problem!
- Whenever possible, use numpy arrays or Python collections of numpy arrays for maximum performance.
- Lock objects are mostly unserializable, because copying a lock is meaningless and could cause serious concurrency problems. You may have to come up with a workaround if your object contains a lock.
Last resort: Custom Serialization
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
@@ -117,3 +107,11 @@ On Linux, it is possible to increase the write throughput of the Plasma object s
.. _`Apache Arrow`: https://arrow.apache.org/
Known Issues
------------
Users could experience memory leak when using certain python3.8 & 3.9 versions. This is due to `a bug in python's pickle module <https://bugs.python.org/issue39492>`_.
This issue has been solved for Python 3.8.2rc1, Python 3.9.0 alpha 4 or late versions.