dask / dask/distributed

Deadlock on submiting tasks by chunks

Open
#3,718 1 comment 0 reactions 0 assignees View on GitHub
Dominant language
Python
Stars
1.7k
Forks
778
Avg merge
2h 50m
Merged PRs (30d)
3

Description

Environments:
* python 2.7 (anaconda2), dask 1.2.2, distributed 1.28.1
* python 3.8, dask 2.14.0, distributed 2.14.0
---

I try to schedule tasks one by one when prev task completed or new worker added to scheduler. The reason is what workers appears at runtime, and dask sends all tasks to first found worker even if resources restriction is set (this probably different question).

Logic is:

1. Submit single task to idle worker.
2. Wait for new worker or ``forgotten`` tasks state via ``SchedulerPlugin``.
3. Go to step 1 if there are remaining tasks.

Deadlock happens on second task scheduling.

I use workers in docker containers but reproduced the issue with ``LocalCluster``:

```python
import logging
from distributed import LocalCluster, Client
from distributed.diagnostics.plugin import SchedulerPlugin
from tornado import gen
from tornado.concurrent import Future
from tornado.locks import Event

logging.basicConfig(level=logging.INFO, format='%(message)s')
logger = logging.getLogger()

class TaskSchedulingPlugin(SchedulerPlugin):
def __init__(self, client, func):
super(TaskSchedulingPlugin, self).__init__()
self.client = client
self.func = func
self._num_tasks = self._num_remaining_tasks = 0
self._results = self._all_scheduled = None

def transition(self, key, start, finish, *args, **kwargs):
logger.info('%s: %s -> %s', key, start, finish)
if finish == 'forgotten' and self._num_remaining_tasks:
self.client.loop.add_callback(self._schedule_next)

def add_worker(self, scheduler=None, worker=None, **kwargs):
logger.info('add_worker(): %s', worker)
self.client.loop.add_callback(self._schedule_next)

@gen.coroutine
def _wait(self, f):
yield self._all_scheduled.wait()
logger.info('All tasks scheduled.')

results = yield self._results
logger.info('All tasks completed.')

f.set_result(results)

@gen.coroutine
def _schedule_next(self):
if not self._num_remaining_tasks:
logger.info('_schedule_next(): no remaining tasks')
return

num_idle_workers = len(self.client.cluster.scheduler.idle)
logger.info('_schedule_next(): idle workers: %s', num_idle_workers)

for _ in range(num_idle_workers):
idx = self._num_tasks - self._num_remaining_tasks
self._num_remaining_tasks = self._num_remaining_tasks - 1

self._results.append(self.client.submit(self.func, idx))
logger.info('_schedule_next(): scheduled task %s', idx)

if not self._num_remaining_tasks:
self._all_scheduled.set()
break

yield

def schedule(self, num_tasks):
self._num_tasks = self._num_remaining_tasks = num_tasks
self._results = []
self._all_scheduled = Event()
f = Future()

self.client.loop.add_callback(self._wait, f)
self.client.loop.add_callback(self._schedule_next)
return f

def remote_func(index):
return index * 10

def main(num_workers, num_tasks):
cluster = LocalCluster(n_workers=num_workers, threads_per_worker=1)

with Client(cluster) as client:
# Setup task scheduler.
plugin = TaskSchedulingPlugin(client, remote_func)
cluster.scheduler.add_plugin(plugin)

# Start tasks scheduling.
@gen.coroutine
def schedule():
r = yield plugin.schedule(num_tasks)
raise gen.Return(r)
results = client.sync(schedule)

# Display results.
for i, res in enumerate(results, start=1):
logger.info('RESULT[%d]: %s', i, res)

if __name__ == '__main__':
main(num_workers=1, num_tasks=2)
```

Output for 2 workers and 2 tasks ``main(num_workers=2, num_tasks=2)``:
```
_schedule_next(): idle workers: 2
_schedule_next(): scheduled task 0
_schedule_next(): scheduled task 1
All tasks scheduled.
remote_func-74047a4aa1432b007504297e023c6dcd: released -> waiting
remote_func-74047a4aa1432b007504297e023c6dcd: waiting -> processing
remote_func-8b85c9df1ea8f27ac84a510b00354111: released -> waiting
remote_func-8b85c9df1ea8f27ac84a510b00354111: waiting -> processing
remote_func-74047a4aa1432b007504297e023c6dcd: processing -> memory
remote_func-8b85c9df1ea8f27ac84a510b00354111: processing -> memory
All tasks completed.
RESULT[1]: 0
RESULT[2]: 10
remote_func-74047a4aa1432b007504297e023c6dcd: memory -> forgotten
remote_func-8b85c9df1ea8f27ac84a510b00354111: memory -> forgotten
```

Output for 1 workers and 2 tasks ``main(num_workers=1, num_tasks=2)`` (deadlock):
```
_schedule_next(): idle workers: 1
_schedule_next(): scheduled task 0
remote_func-74047a4aa1432b007504297e023c6dcd: released -> waiting
remote_func-74047a4aa1432b007504297e023c6dcd: waiting -> processing
remote_func-74047a4aa1432b007504297e023c6dcd: processing -> memory
```

Here I expected processing like:
```
_schedule_next(): idle workers: 1
_schedule_next(): scheduled task 0
remote_func-74047a4aa1432b007504297e023c6dcd: released -> waiting
remote_func-74047a4aa1432b007504297e023c6dcd: waiting -> processing
remote_func-74047a4aa1432b007504297e023c6dcd: processing -> memory
remote_func-74047a4aa1432b007504297e023c6dcd: memory -> forgotten
_schedule_next(): idle workers: 1
_schedule_next(): scheduled task 1
All tasks scheduled.
remote_func-8b85c9df1ea8f27ac84a510b00354111: released -> waiting
remote_func-8b85c9df1ea8f27ac84a510b00354111: waiting -> processing
remote_func-8b85c9df1ea8f27ac84a510b00354111: processing -> memory
All tasks completed.
remote_func-8b85c9df1ea8f27ac84a510b00354111: memory -> forgotten
...
```

Can I send tasks like that or daks is not supposed to schedule this way?

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.