Error transitioning task from 'processing' to 'memory' when using annotations
- Dominant language
- Python
- Stars
- 1.7k
- Forks
- 778
- Avg merge
- 2h 50m
- Merged PRs (30d)
- 3
Description
**What happened**:
Tasks which are currently persisting cause annotated tasks to fail.
**What you expected to happen**:
Annotated and non-annotated tasks should run concurrently.
**Minimal Complete Verifiable Example**:
_Optional: Create a new environment with latest Dask_
```console
$ conda create -n daskwait -c conda-forge dask distributed ipython -y
$ conda activate daskwait
$ ipython
```
```python
import dask
import dask.array as da
from dask.distributed import LocalCluster, Client, Scheduler, Nanny, wait
# Create a scheduler
scheduler = await Scheduler()
# Launch a worker with no resources
nanny = await Nanny(scheduler.address)
# Launch a worker with 1 FOO resources
foo_nanny = await Nanny(scheduler.address, resources={"FOO": 1})
# Connect a client
client = await Client(scheduler.address, asynchronous=True)
# Create some data
arr = da.random.random((10_000, 10_000), chunks=(1000, 1000)).persist()
# Use the annotated workers for compute
with dask.annotate(resources={'FOO': 1}):
result = arr.mean().persist()
await wait(result)
```
```python-traceback
distributed.scheduler - ERROR - 'FOO'
Traceback (most recent call last):
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2845, in transition_processing_memory
_remove_from_processing(self, ts)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 8021, in _remove_from_processing
state.release_resources(ts, ws)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 3479, in release_resources
ws._used_resources[r] -= required
KeyError: 'FOO'
distributed.scheduler - ERROR - Error transitioning "('random_sample-7038b7ac935ddd67754c4d9b22877b64', 5, 0)" from 'processing' to 'memory'
```
Full traceback
```python-traceback
distributed.scheduler - INFO - Clear task state
distributed.scheduler - INFO - Scheduler at: tcp://10.51.100.43:41979
distributed.scheduler - INFO - dashboard at: :8787
distributed.nanny - INFO - Start Nanny at: 'tcp://10.51.100.43:42819'
distributed.worker - INFO - Start worker at: tcp://10.51.100.43:43427
distributed.worker - INFO - Listening to: tcp://10.51.100.43:43427
distributed.worker - INFO - dashboard at: 10.51.100.43:33969
distributed.worker - INFO - Waiting to connect to: tcp://10.51.100.43:41979
distributed.worker - INFO - -------------------------------------------------
distributed.worker - INFO - Threads: 12
distributed.worker - INFO - Memory: 93.04 GiB
distributed.worker - INFO - Local Directory: /home/jtomlinson/Scratch/dask-worker-space/worker-4eyrd3zl
distributed.worker - INFO - -------------------------------------------------
distributed.scheduler - INFO - Register worker
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.51.100.43:43427
distributed.worker - INFO - Registered to: tcp://10.51.100.43:41979
distributed.worker - INFO - -------------------------------------------------
distributed.core - INFO - Starting established connection
distributed.core - INFO - Starting established connection
distributed.nanny - INFO - Start Nanny at: 'tcp://10.51.100.43:38997'
distributed.worker - INFO - Start worker at: tcp://10.51.100.43:33137
distributed.worker - INFO - Listening to: tcp://10.51.100.43:33137
distributed.worker - INFO - dashboard at: 10.51.100.43:35407
distributed.worker - INFO - Waiting to connect to: tcp://10.51.100.43:41979
distributed.worker - INFO - -------------------------------------------------
distributed.worker - INFO - Threads: 12
distributed.worker - INFO - Memory: 93.04 GiB
distributed.worker - INFO - Local Directory: /home/jtomlinson/Scratch/dask-worker-space/worker-wpq_e4ft
distributed.worker - INFO - -------------------------------------------------
distributed.scheduler - INFO - Register worker
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.51.100.43:33137
distributed.worker - INFO - Registered to: tcp://10.51.100.43:41979
distributed.worker - INFO - -------------------------------------------------
distributed.core - INFO - Starting established connection
distributed.core - INFO - Starting established connection
distributed.scheduler - INFO - Receive client connection: Client-9e96453c-a160-11ec-940c-80e82ccdc37c
distributed.core - INFO - Starting established connection
dask.array
distributed.scheduler - ERROR - 'FOO'
Traceback (most recent call last):
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2845, in transition_processing_memory
_remove_from_processing(self, ts)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 8021, in _remove_from_processing
state.release_resources(ts, ws)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 3479, in release_resources
ws._used_resources[r] -= required
KeyError: 'FOO'
distributed.scheduler - ERROR - Error transitioning "('random_sample-7038b7ac935ddd67754c4d9b22877b64', 5, 0)" from 'processing' to 'memory'
Traceback (most recent call last):
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2277, in _transition
recommendations, client_msgs, worker_msgs = func(key, *args, **kwargs)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2845, in transition_processing_memory
_remove_from_processing(self, ts)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 8021, in _remove_from_processing
state.release_resources(ts, ws)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 3479, in release_resources
ws._used_resources[r] -= required
KeyError: 'FOO'
distributed.core - ERROR - 'FOO'
Traceback (most recent call last):
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/core.py", line 586, in handle_stream
handler(**merge(extra, msg))
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 5466, in handle_task_finished
r: tuple = self.stimulus_task_finished(key=key, worker=worker, **msg)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 4905, in stimulus_task_finished
r: tuple = parent._transition(key, "memory", worker=worker, **kwargs)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2277, in _transition
recommendations, client_msgs, worker_msgs = func(key, *args, **kwargs)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2845, in transition_processing_memory
_remove_from_processing(self, ts)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 8021, in _remove_from_processing
state.release_resources(ts, ws)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 3479, in release_resources
ws._used_resources[r] -= required
KeyError: 'FOO'
distributed.scheduler - INFO - Remove worker
distributed.core - INFO - Removing comms to tcp://10.51.100.43:43427
distributed.utils - ERROR - 'FOO'
Traceback (most recent call last):
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/utils.py", line 681, in log_errors
yield
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 4516, in add_worker
await self.handle_worker(comm=comm, worker=address)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 5607, in handle_worker
await self.handle_stream(comm=comm, extra={"worker": worker})
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/core.py", line 586, in handle_stream
handler(**merge(extra, msg))
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 5466, in handle_task_finished
r: tuple = self.stimulus_task_finished(key=key, worker=worker, **msg)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 4905, in stimulus_task_finished
r: tuple = parent._transition(key, "memory", worker=worker, **kwargs)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2277, in _transition
recommendations, client_msgs, worker_msgs = func(key, *args, **kwargs)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2845, in transition_processing_memory
_remove_from_processing(self, ts)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 8021, in _remove_from_processing
state.release_resources(ts, ws)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 3479, in release_resources
ws._used_resources[r] -= required
KeyError: 'FOO'
distributed.core - ERROR - Exception while handling op register-worker
Traceback (most recent call last):
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/core.py", line 520, in handle_comm
result = await result
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 4516, in add_worker
await self.handle_worker(comm=comm, worker=address)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 5607, in handle_worker
await self.handle_stream(comm=comm, extra={"worker": worker})
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/core.py", line 586, in handle_stream
handler(**merge(extra, msg))
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 5466, in handle_task_finished
r: tuple = self.stimulus_task_finished(key=key, worker=worker, **msg)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 4905, in stimulus_task_finished
r: tuple = parent._transition(key, "memory", worker=worker, **kwargs)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2277, in _transition
recommendations, client_msgs, worker_msgs = func(key, *args, **kwargs)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2845, in transition_processing_memory
_remove_from_processing(self, ts)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 8021, in _remove_from_processing
state.release_resources(ts, ws)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 3479, in release_resources
ws._used_resources[r] -= required
KeyError: 'FOO'
tornado.application - ERROR - Exception in callback functools.partial(. at 0x7fc6d1304280>, exception=KeyError('FOO')>)
Traceback (most recent call last):
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/tornado/ioloop.py", line 741, in _run_callback
ret = callback()
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/tornado/tcpserver.py", line 331, in
gen.convert_yielded(future), lambda f: f.result()
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/comm/tcp.py", line 530, in _handle_stream
await self.comm_handler(comm)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/core.py", line 520, in handle_comm
result = await result
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 4516, in add_worker
await self.handle_worker(comm=comm, worker=address)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 5607, in handle_worker
await self.handle_stream(comm=comm, extra={"worker": worker})
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/core.py", line 586, in handle_stream
handler(**merge(extra, msg))
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 5466, in handle_task_finished
r: tuple = self.stimulus_task_finished(key=key, worker=worker, **msg)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 4905, in stimulus_task_finished
r: tuple = parent._transition(key, "memory", worker=worker, **kwargs)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2277, in _transition
recommendations, client_msgs, worker_msgs = func(key, *args, **kwargs)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2845, in transition_processing_memory
_remove_from_processing(self, ts)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 8021, in _remove_from_processing
state.release_resources(ts, ws)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 3479, in release_resources
ws._used_resources[r] -= required
KeyError: 'FOO'
distributed.worker - INFO - Connection to scheduler broken. Reconnecting...
distributed.batched - INFO - Batched Comm Closed Scheduler local=tcp://10.51.100.43:45970 remote=tcp://10.51.100.43:41979>
Traceback (most recent call last):
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/batched.py", line 93, in _background_send
nbytes = yield self.comm.write(
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/tornado/gen.py", line 762, in run
value = future.result()
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/comm/tcp.py", line 247, in write
raise CommClosedError()
distributed.comm.core.CommClosedError
distributed.worker - INFO - -------------------------------------------------
distributed.scheduler - INFO - Unexpected worker completed task. Expected: None, Got: , Key: ('random_sample-7038b7ac935ddd67754c4d9b22877b64', 5, 0)
distributed.scheduler - ERROR - 'NoneType' object has no attribute 'address'
Traceback (most recent call last):
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2810, in transition_processing_memory
worker_msgs[ts._processing_on.address] = [
AttributeError: 'NoneType' object has no attribute 'address'
distributed.scheduler - ERROR - Error transitioning "('random_sample-7038b7ac935ddd67754c4d9b22877b64', 5, 0)" from 'processing' to 'memory'
Traceback (most recent call last):
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2277, in _transition
recommendations, client_msgs, worker_msgs = func(key, *args, **kwargs)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2810, in transition_processing_memory
worker_msgs[ts._processing_on.address] = [
AttributeError: 'NoneType' object has no attribute 'address'
distributed.utils - ERROR - 'NoneType' object has no attribute 'address'
Traceback (most recent call last):
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/utils.py", line 681, in log_errors
yield
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 4450, in add_worker
t: tuple = parent._transition(
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2277, in _transition
recommendations, client_msgs, worker_msgs = func(key, *args, **kwargs)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2810, in transition_processing_memory
worker_msgs[ts._processing_on.address] = [
AttributeError: 'NoneType' object has no attribute 'address'
distributed.core - ERROR - Exception while handling op register-worker
Traceback (most recent call last):
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/core.py", line 520, in handle_comm
result = await result
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 4450, in add_worker
t: tuple = parent._transition(
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2277, in _transition
recommendations, client_msgs, worker_msgs = func(key, *args, **kwargs)
File "/home/jtomlinson/miniconda3/envs/rapids-22.02/lib/python3.8/site-packages/distributed/scheduler.py", line 2810, in transition_processing_memory
worker_msgs[ts._processing_on.address] = [
AttributeError: 'NoneType' object has no attribute 'address'
```
**Anything else we need to know?**:
If I add a `wait` to the data generation before moving to the annotated section things work. But this means that we have a synchronisation point here that isn't ideal.
```python
...
# Create some data
arr = da.random.random((10_000, 10_000), chunks=(1000, 1000)).persist()
await wait(arr)
# Use the annotated workers for compute
with dask.annotate(resources={'FOO': 1}):
result = arr.mean().persist()
await wait(arr)
```
**Environment**:
- Dask version: 2022.02.1
- Python version: Tried 3.8 and 3.10 on multiple machines
- Operating System: Linux (Ubuntu)
- Install method (conda, pip, source): conda
Cluster Dump State:
Contributor guide
Assessment
This issue has not been assessed yet.