# Dask on Ray + Ray Distributed Cluster - Workers not getting used?

**URL:** <https://discuss.ray.io/t/dask-on-ray-ray-distributed-cluster-workers-not-getting-used/845>\
**Category:** Ray Core\
**Created:** [February 11, 2021, 5:34pm UTC](https://discuss.ray.io/t/dask-on-ray-ray-distributed-cluster-workers-not-getting-used/845 "2021-02-11T17:34:51Z")\
**Posts on this page:** 10\
**Page:** 1

<div class="post-metadata">

**Author:** ![jennakwon06](https://avatars.discourse-cdn.com/v4/letter/j/b19c9b/32.png) [@jennakwon06](https://discuss.ray.io/u/jennakwon06)\
**Post date:** [February 11, 2021, 5:34pm UTC](https://discuss.ray.io/t/dask-on-ray-ray-distributed-cluster-workers-not-getting-used/845/1 "2021-02-11T17:34:51Z")

</div>

(Pleas feel free to re-categorize)

**Context**

The context of my work is that I am bench-marking Dask workloads (Petabytes of data) on Dask distributed backend vs Ray backend.

I’ve done the benchmark for Dask distributed backend. My workload would cost ~$1 million dollars (yes you’ve read that right) with Dask distributed backend - it’d be running 35 r5.24xlarge instances continuously for ~180 days.

I am hoping Ray backend can be better 🙂

So I’ve set up a Ray cluster with AWS EC2 (autoscaler).  
The workload is about a million Dask graphs. Each graph’s last node is a Dask delayed object, so I do the compute using `delayed_obj.compute(scheduler=ray_dask_get)`.

**Problem**

It doesn’t seem that any of my worker nodes are getting utilized. I am submitting 1000 graphs to the cluster at a time, so there should be a plenty of work to do.

I am wondering if Dask on Ray + Ray Distributed Cluster is supposed to work? I am looking at this line [ray/scheduler.py at master · ray-project/ray · GitHub](https://github.com/ray-project/ray/blob/master/python/ray/util/dask/scheduler.py#L71) . I have head/worker nodes with VCPU 96. Total CPU in the cluster should be 96 \* 35. But when I do `print(CPU_COUNT`), where CPU\_COUNT is from `dask.system`, I’ll get 96 (which makes sense). So does this mean thread pool of size 96 is getting used instead of pool of 96 \* 35?

---

<div class="post-metadata">

**Author:** ![Stephanie\_Wang](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/stephanie_wang/32/18_2.png) [@Stephanie\_Wang](https://discuss.ray.io/u/Stephanie_Wang)\
**Post date:** [February 11, 2021, 5:56pm UTC](https://discuss.ray.io/t/dask-on-ray-ray-distributed-cluster-workers-not-getting-used/845/2 "2021-02-11T17:56:36Z")

</div>

HI @jennakwon06, Dask Distributed and Dask-on-Ray won’t work well together, since they are different backends for Dask. You can default to using the Dask-on-Ray scheduler by setting it in the config at the beginning of your script:

```python
dask.config.set(scheduler=ray_dask_get)

```

Note that in this case, you should not create a Dask.distributed client, since according to the Dask [docs](https://distributed.dask.org/en/latest/client.html), this will set the Dask.distributed cluster as the default backend, like this:

```python
client = Client('scheduler:8786')

```

---

<div class="post-metadata">

**Author:** ![sangcho](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/sangcho/32/425_2.png) [@sangcho](https://discuss.ray.io/u/sangcho)\
**Post date:** [February 11, 2021, 6:06pm UTC](https://discuss.ray.io/t/dask-on-ray-ray-distributed-cluster-workers-not-getting-used/845/3 "2021-02-11T18:06:37Z")

</div>

Can you also try something like

```auto
@ray.remote
def f():
    time.sleep(10)

refs = [f.remote() for _ in range(10000)]

```

and if the cluster gets utilized?

---

<div class="post-metadata">

**Author:** ![Clark\_Zinzow](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/clark_zinzow/32/445_2.png) [@Clark\_Zinzow](https://discuss.ray.io/u/Clark_Zinzow)\
**Post date:** [February 11, 2021, 6:43pm UTC](https://discuss.ray.io/t/dask-on-ray-ray-distributed-cluster-workers-not-getting-used/845/4 "2021-02-11T18:43:45Z")

</div>

Just to clear something up, that threadpool that you linked to is only used by the Dask-on-Ray scheduler to submit Ray tasks, not to execute Ray tasks. Specifically, the Dask-on-Ray scheduler traverses the Dask graph and submits each Dask task as a Ray task to the Ray cluster, using that threadpool to submit tasks in parallel; the Ray tasks should be executing across all of the cores (and machines) that have been allocated to your Ray cluster.

---

<div class="post-metadata">

**Author:** ![jennakwon06](https://avatars.discourse-cdn.com/v4/letter/j/b19c9b/32.png) [@jennakwon06](https://discuss.ray.io/u/jennakwon06)\
**Post date:** [February 11, 2021, 9:25pm UTC](https://discuss.ray.io/t/dask-on-ray-ray-distributed-cluster-workers-not-getting-used/845/5 "2021-02-11T21:25:59Z")

</div>

Hi Stephanie! Thanks for replying. Maybe my original post was misleading. I edited it a bit so hopefully you’ll have more context.

When I am benchmarking Ray, I am not using Dask Distributed client.

---

<div class="post-metadata">

**Author:** ![jennakwon06](https://avatars.discourse-cdn.com/v4/letter/j/b19c9b/32.png) [@jennakwon06](https://discuss.ray.io/u/jennakwon06)\
**Post date:** [February 11, 2021, 9:59pm UTC](https://discuss.ray.io/t/dask-on-ray-ray-distributed-cluster-workers-not-getting-used/845/6 "2021-02-11T21:59:09Z")

</div>

Hey @sangcho! Great suggestion. So yes the cluster gets utilized while submitting a script with that code.

---

<div class="post-metadata">

**Author:** ![Stephanie\_Wang](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/stephanie_wang/32/18_2.png) [@Stephanie\_Wang](https://discuss.ray.io/u/Stephanie_Wang)\
**Post date:** [February 12, 2021, 2:21am UTC](https://discuss.ray.io/t/dask-on-ray-ray-distributed-cluster-workers-not-getting-used/845/7 "2021-02-12T02:21:03Z")

</div>

Yes, I misunderstood in my previous answer! This sounds like either there is a bug in the setup or a bug in Ray. Could you open a github issue for this? A reproducible script would be ideal so that a Ray developer can try it out too.

---

<div class="post-metadata">

**Author:** ![jennakwon06](https://avatars.discourse-cdn.com/v4/letter/j/b19c9b/32.png) [@jennakwon06](https://discuss.ray.io/u/jennakwon06)\
**Post date:** [February 12, 2021, 4:05pm UTC](https://discuss.ray.io/t/dask-on-ray-ray-distributed-cluster-workers-not-getting-used/845/8 "2021-02-12T16:05:51Z")

</div>

Yep sounds good! I’m working directly with Sang on it 🙂

---

<div class="post-metadata">

**Author:** ![sangcho](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/sangcho/32/425_2.png) [@sangcho](https://discuss.ray.io/u/sangcho)\
**Post date:** [February 12, 2021, 11:16pm UTC](https://discuss.ray.io/t/dask-on-ray-ray-distributed-cluster-workers-not-getting-used/845/9 "2021-02-12T23:16:31Z")

</div>

For the future references who would find similar issues, the problem was just that we don’t support Dask client’s `.compute` method! If you are calling `.compute` within Dask on Ray, it follows the same API behavior as regular Dask running in a cluster, which means `.compute` will be a blocking call, not generating async futures!

---

<div class="post-metadata">

**Author:** ![jennakwon06](https://avatars.discourse-cdn.com/v4/letter/j/b19c9b/32.png) [@jennakwon06](https://discuss.ray.io/u/jennakwon06)\
**Post date:** [February 14, 2021, 1:25am UTC](https://discuss.ray.io/t/dask-on-ray-ray-distributed-cluster-workers-not-getting-used/845/10 "2021-02-14T01:25:45Z")

</div>

Yes. To complement Sang’s comment -

I have a pattern where I submit `BATCH_SIZE` number of Dask graphs (collections) to the cluster at a time by calling `res_future = dask_client.compute(delayed_object)` for `BATCH_SIZE` number of times. I collect the `res_future` in a list of length `BATCH_SIZE` and use `dask.distributed.wait` on the list to wait until all graphs are done. I do this many times.

I was trying to do the same thing in Ray and I didn’t realize Ray didn’t return futures! So that was my bad.

The problem I am having right now is in equal work distribution though. When `BATCH_SIZE` number of Dask graphs are submitted to the cluster, it seems that all the tasks go to the head node first, then get distributed across the workers. I believe Sang is working on improving it for all nodes to get evenly full, to reduce pressure on distributed object storage.
