dask / dask/distributed

Workers idle even though there's queued work

Open
#4,501 11 comments 1 reaction 0 assignees View on GitHub
needs info
Dominant language
Python
Stars
1.7k
Forks
778
Avg merge
2h 50m
Merged PRs (30d)
3

Description

**What happened**: We have a large-ish cluster (about 100 nodes) and recently when we submit a lot of jobs (in the thousands) we notice that about 60% of the cluster is idle. Generally, a job will spawn about 20 downstream sub-jobs; these are submitted from inside the worker which will call secede / rejoin while it waits on those jobs. I'm fairly certain this use of secede / rejoin is related as you can see in the reproduction below.

**What you expected to happen**: The cluster uses all available resources

**Minimal Complete Verifiable Example**:

This requires running a scheduler, a worker with two procs, and then submitting jobs. Bear with me while I show all the pieces:

This is how I create the environment:

```bash
#!/bin/bash
python3 -m venv env
source env/bin/activate
pip install "dask[distributed,diagnostics,delayed]==2020.12.0"
```

...and this is the python file with the jobs and such. You can see the two operations are:
1) Submit a child job to wait for X seconds, and wait on that job. We call secede / rejoin while waiting.
2) Instantly print out a message

```python
import time
import dask.distributed
import sys

def long_running_job(seconds):
print(f"Doing the actual work ({seconds}s)")
time.sleep(seconds)
print(f"Finished working ({seconds}s)")

def root_job(seconds):
client = dask.distributed.Client(address="tcp://127.0.0.1:8786")
futures = client.map(long_running_job, [seconds])
print(f"Submitted long runing job ({seconds}s); seceding while we wait")
dask.distributed.secede()
client.gather(futures)
print(f"Job done ({seconds}s); rejoining")
dask.distributed.rejoin()

def other_job(message):
print(message)

if __name__ == "__main__":
client = dask.distributed.Client(address="tcp://127.0.0.1:8786")
if sys.argv[1] == "wait":
future = client.submit(root_job, int(sys.argv[2]))
dask.distributed.fire_and_forget(future)
elif sys.argv[1] == "message":
future = client.submit(other_job, sys.argv[2])
dask.distributed.fire_and_forget(future)
```

Finally, this is the script that will submit jobs that will show the issue we're running into:

```bash
#!/bin/bash

# start a scheduler in one terminal:
# $ dask-scheduler
# ...and a worker in another:
# $ dask-worker tcp://10.0.2.15:8786 --nthreads 1 --nprocs 2
# then run the below:

python stuff.py wait 120
sleep 1
python stuff.py wait 60
sleep 1
python stuff.py message instantaneous-job
```

This script submits a long job, a shorter job, and then just an instantaneous job to show that there's a scheduling problem. When the jobs are submitted the worker will print out:

```
Submitted long runing job (120s); seceding while we wait
Doing the actual work (120s)
Submitted long runing job (60s); seceding while we wait
Doing the actual work (60s)
Finished working (60s)
Job done (60s); rejoining
Finished working (120s)
instantaneous-job
Job done (120s); rejoining
```

The problem is on the line that says `Job done (60s); rejoining`. At this point there's one idle worker that *could* be running the instantaneous job but it doesn't -- instead it waits on the 120s job. After the 120s job is done (about a minute later) that instantaneous job finally runs. Hence the worker is idle for about a minute.

**Anything else we need to know?**:
Sorry for the length; I don't think I can cut it down any more. If the problem isn't clear let me know and I'll see if I can explain better.

**Environment**:

- Dask version: 2020.12.0
- Python version: 3.6.9
- Operating System: Ubuntu 20.04
- Install method (conda, pip, source): pip

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.