Add documentation for Dask-on-Ray (#10452)

* Dask on ray

* Dask on ray

* line

* Add links
This commit is contained in:
Stephanie Wang
2020-09-01 08:13:04 -07:00
committed by GitHub
parent d584a4e5c4
commit 9c2c952262
5 changed files with 51 additions and 1 deletions
+43
View File
@@ -0,0 +1,43 @@
Dask on Ray
===========
Ray offers an experimental scheduler backend for Dask.
With this plugin, you can use familiar Dask APIs such as Dask DataFrames, and the computation will be executed by the Ray system.
The Ray plugin can be used with any Dask `.compute() <https://docs.dask.org/en/latest/api.html#dask.compute>`__ call.
Note that for execution on a Ray cluster, you should *not* use the `Dask.distributed <https://distributed.dask.org/en/latest/quickstart.html>`__ client.
Just follow the instructions for :ref:`using Ray on a cluster <using-ray-on-a-cluster>` to modify the ``ray.init()`` call.
Here's an example:
.. code-block:: python
import ray
from ray.experimental.dask import ray_dask_get
import dask.delayed
from time import sleep
# Start Ray.
# Tip: If you're connecting to an existing cluster, use ray.init(address="auto").
ray.init()
def inc(x):
sleep(1)
return x + 1
def add(x, y):
sleep(1)
return x + y
x = dask.delayed(inc)(1)
y = dask.delayed(inc)(2)
z = dask.delayed(add)(x, y)
# The Dask scheduler submits the recorded task graph to Ray.
z.compute(scheduler=ray_dask_get)
Why use this feature?
1. If you'd like to use Dask and Ray libraries in the same application.
2. To take advantage of Ray-specific features such as the :ref:`cluster launcher <ref-automatic-cluster>` and :ref:`shared-memory store <memory>`.
Note that Dask-on-Ray is an ongoing project and is not expected to achieve the same performance as using Ray directly.
+1
View File
@@ -207,6 +207,7 @@ Academic Papers
joblib.rst
iter.rst
pandas_on_ray.rst
dask-on-ray.rst
.. toctree::
:hidden:
+2
View File
@@ -1,3 +1,5 @@
.. _memory:
Memory Management
=================
+3 -1
View File
@@ -26,6 +26,8 @@ On top of **Ray Core** are several libraries for solving problems in machine lea
Ray also has a number of other community contributed libraries:
- :doc:`../dask-on-ray`
- `Mars on Ray <https://github.com/mars-project/mars/pull/1508>`__
- :doc:`../pandas_on_ray`
- :doc:`../joblib`
- :doc:`../multiprocessing`
- :doc:`../multiprocessing`
+2
View File
@@ -43,6 +43,8 @@ To check if Ray is initialized, you can call ``ray.is_initialized()``:
See the `Configuration <configure.html>`__ documentation for the various ways to configure Ray.
.. _using-ray-on-a-cluster:
Using Ray on a cluster
----------------------