Deadlock on submiting tasks by chunks
- 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
Assessment
This issue has not been assessed yet.