dask / dask/distributed

KeyboardInterrupt kills workers in LocalCluster

Open
#1,925 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

```python
In [1]: from distributed import Client
cl
In [2]: client = Client(silence_logs=False, n_workers=1)

distributed.scheduler - INFO - Clear task state
distributed.scheduler - INFO - Scheduler at: tcp://127.0.0.1:42699
distributed.scheduler - INFO - bokeh at: 127.0.0.1:8787
distributed.nanny - INFO - Start Nanny at: 'tcp://127.0.0.1:38145'
distributed.worker - INFO - Start worker at: tcp://127.0.0.1:34721
distributed.worker - INFO - Listening to: tcp://127.0.0.1:34721
distributed.worker - INFO - nanny at: 127.0.0.1:38145
distributed.worker - INFO - bokeh at: 127.0.0.1:34265
distributed.worker - INFO - Waiting to connect to: tcp://127.0.0.1:42699
distributed.worker - INFO - -------------------------------------------------
distributed.worker - INFO - Threads: 4
distributed.worker - INFO - Memory: 16.68 GB
distributed.worker - INFO - Local Directory: /home/mrocklin/dask-worker-space/worker-4rlf_r1g
distributed.worker - INFO - -------------------------------------------------
distributed.scheduler - INFO - Register tcp://127.0.0.1:34721
distributed.worker - INFO - Registered to: tcp://127.0.0.1:42699
distributed.worker - INFO - -------------------------------------------------
distributed.scheduler - INFO - Starting worker compute stream, tcp://127.0.0.1:34721
distributed.scheduler - INFO - Receive client connection: Client-cde49f30-43d6-11e8-8b47-9b9dc1e1c1fe

In [3]:
^[[24;1R
In [3]: import time

In [4]: future = client.submit(time.sleep, 20)

In [5]: future.result()
^Cdistributed.scheduler - INFO - Worker 'tcp://127.0.0.1:34721' failed from closed comm: in : Stream is closed
distributed.scheduler - INFO - Remove worker tcp://127.0.0.1:34721
distributed.scheduler - INFO - Lost all workers
distributed.nanny - WARNING - Restarting worker
---------------------------------------------------------------------------
KeyboardInterrupt Traceback (most recent call last)
in ()
----> 1 future.result()

~/workspace/distributed/distributed/client.py in result(self, timeout)
167 # shorten error traceback
168 result = self.client.sync(self._result, callback_timeout=timeout,
--> 169 raiseit=False)
170 if self.status == 'error':
171 six.reraise(*result)

~/workspace/distributed/distributed/client.py in sync(self, func, *args, **kwargs)
613 return future
614 else:
--> 615 return sync(self.loop, func, *args, **kwargs)
616
617 def __repr__(self):

~/workspace/distributed/distributed/utils.py in sync(loop, func, *args, **kwargs)
249 else:
250 while not e.is_set():
--> 251 e.wait(10)
252 if error[0]:
253 six.reraise(*error[0])

~/Software/anaconda/lib/python3.6/threading.py in wait(self, timeout)
549 signaled = self._flag
550 if not signaled:
--> 551 signaled = self._cond.wait(timeout)
552 return signaled
553

~/Software/anaconda/lib/python3.6/threading.py in wait(self, timeout)
297 else:
298 if timeout > 0:
--> 299 gotit = waiter.acquire(True, timeout)
300 else:
301 gotit = waiter.acquire(False)

KeyboardInterrupt:

In [6]: distributed.worker - INFO - Start worker at: tcp://127.0.0.1:39105
distributed.worker - INFO - Listening to: tcp://127.0.0.1:39105
distributed.worker - INFO - nanny at: 127.0.0.1:38145
distributed.worker - INFO - bokeh at: 127.0.0.1:41421
distributed.worker - INFO - Waiting to connect to: tcp://127.0.0.1:42699
distributed.worker - INFO - -------------------------------------------------
distributed.worker - INFO - Threads: 4
distributed.worker - INFO - Memory: 16.68 GB
distributed.worker - INFO - Local Directory: /home/mrocklin/dask-worker-space/worker-knglf610
distributed.worker - INFO - -------------------------------------------------
distributed.scheduler - INFO - Register tcp://127.0.0.1:39105
distributed.worker - INFO - Registered to: tcp://127.0.0.1:42699
distributed.worker - INFO - -------------------------------------------------
distributed.scheduler - INFO - Starting worker compute stream, tcp://127.0.0.1:39105
```

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.