distributed.worker - ERROR - ('waiting', 'executing') OR distributed.protocol.core - CRITICAL - Failed to deserialize
- Dominant language
- Python
- Stars
- 1.7k
- Forks
- 778
- Avg merge
- 2h 50m
- Merged PRs (30d)
- 3
Description
**What happened**:
The MCV example below works on dask+distributed 2021.3.0, but an error is reported. The same example does *not* work on later versions, but instead the code either stops with an error (2021.3.1 and 2021.4.0), or just hangs (2021.4.1 and later versions). On 2021.7.0, the code works if Pandas DataFrame has 28 million rows, but it does not work for 29 million rows.
**What you expected to happen**:
The code should create a folder called ```output``` with 240 CSV files without reporting any errors.
**Minimal Complete Verifiable Example**:
```python
import multiprocessing.popen_spawn_win32
from dask.distributed import Client
import dask.dataframe as dd
import pandas as pd
from shutil import rmtree
import numpy as np
client = Client()
for Nrows in [28e6, 29e6]:
# create a Pandas DataFrame
df = pd.DataFrame(np.random.randn(int(Nrows), 2), columns=list('AB'), index=pd.util.testing.rands_array(16, int(Nrows), dtype='O'))
# create a Dask DataFrame
ddf = dd.from_pandas(df, npartitions=240)
# dump the DataFrame to CSV
ddf.to_csv('output_%d' % int(Nrows))
print('Success writing %d rows.' % int(Nrows))
rmtree('output_%d' % int(Nrows))
```
When using dask=2021.3.0 and distributed=2021.3.0, the following error is reported, but the ```output``` folder is created with 240 CSV files
```
distributed.worker - ERROR - ('waiting', 'executing')
Traceback (most recent call last):
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\worker.py", line 2689, in execute
await self.ensure_computing()
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\worker.py", line 2561, in ensure_computing
self.transition(ts, "executing")
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\worker.py", line 1553, in transition
func = self._transitions[start, finish]
KeyError: ('waiting', 'executing')
tornado.application - ERROR - Exception in callback functools.partial(>, exception=KeyError(('waiting', 'executing'))>)
Traceback (most recent call last):
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\tornado\ioloop.py", line 741, in _run_callback
ret = callback()
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\tornado\ioloop.py", line 765, in _discard_future_result
future.result()
```
With dask+distributed 2021.3.1 or 2021.4.0, the code creates the ```output``` folder without any CSV files, it crashes, and a different error is reported
```
distributed.protocol.core - CRITICAL - Failed to deserialize
Traceback (most recent call last):
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\protocol\core.py", line 104, in loads
return msgpack.loads(
File "msgpack/_unpacker.pyx", line 195, in msgpack._cmsgpack.unpackb
ValueError: 2917822586 exceeds max_bin_len(2147483647)
distributed.core - ERROR - 2917822586 exceeds max_bin_len(2147483647)
Traceback (most recent call last):
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\core.py", line 555, in handle_stream
msgs = await comm.read()
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\comm\tcp.py", line 218, in read
msg = await from_frames(
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\comm\utils.py", line 80, in from_frames
res = _from_frames()
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\comm\utils.py", line 63, in _from_frames
return protocol.loads(
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\protocol\core.py", line 104, in loads
return msgpack.loads(
File "msgpack/_unpacker.pyx", line 195, in msgpack._cmsgpack.unpackb
ValueError: 2917822586 exceeds max_bin_len(2147483647)
distributed.core - ERROR - Exception while handling op register-client
Traceback (most recent call last):
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\core.py", line 501, in handle_comm
result = await result
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\scheduler.py", line 4732, in add_client
await self.handle_stream(comm=comm, extra={"client": client})
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\core.py", line 555, in handle_stream
msgs = await comm.read()
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\comm\tcp.py", line 218, in read
msg = await from_frames(
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\comm\utils.py", line 80, in from_frames
res = _from_frames()
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\comm\utils.py", line 63, in _from_frames
return protocol.loads(
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\protocol\core.py", line 104, in loads
return msgpack.loads(
File "msgpack/_unpacker.pyx", line 195, in msgpack._cmsgpack.unpackb
ValueError: 2917822586 exceeds max_bin_len(2147483647)
tornado.application - ERROR - Exception in callback functools.partial(. at 0x000001D8C9476940>, exception=ValueError('2917822586 exceeds max_bin_len(2147483647)')>)
Traceback (most recent call last):
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\tornado\ioloop.py", line 741, in _run_callback
ret = callback()
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\tornado\tcpserver.py", line 331, in
gen.convert_yielded(future), lambda f: f.result()
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\comm\tcp.py", line 493, in _handle_stream
await self.comm_handler(comm)
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\core.py", line 501, in handle_comm
result = await result
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\scheduler.py", line 4732, in add_client
await self.handle_stream(comm=comm, extra={"client": client})
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\core.py", line 555, in handle_stream
msgs = await comm.read()
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\comm\tcp.py", line 218, in read
msg = await from_frames(
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\comm\utils.py", line 80, in from_frames
res = _from_frames()
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\comm\utils.py", line 63, in _from_frames
return protocol.loads(
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\protocol\core.py", line 104, in loads
return msgpack.loads(
File "msgpack/_unpacker.pyx", line 195, in msgpack._cmsgpack.unpackb
ValueError: 2917822586 exceeds max_bin_len(2147483647)
Traceback (most recent call last):
File "", line 9, in
ddf.to_csv('output')
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\dask\dataframe\core.py", line 1465, in to_csv
return to_csv(self, filename, **kwargs)
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\dask\dataframe\io\csv.py", line 864, in to_csv
delayed(values).compute(**compute_kwargs)
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\dask\base.py", line 283, in compute
(result,) = compute(self, traverse=False, **kwargs)
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\dask\base.py", line 565, in compute
results = schedule(dsk, keys, **kwargs)
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\client.py", line 2665, in get
results = self.gather(packed, asynchronous=asynchronous, direct=direct)
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\client.py", line 1974, in gather
return self.sync(
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\client.py", line 844, in sync
return sync(
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\utils.py", line 353, in sync
raise exc.with_traceback(tb)
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\utils.py", line 336, in f
result[0] = yield future
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\tornado\gen.py", line 762, in run
value = future.result()
File "C:\Users\--\Anaconda3\envs\work2\lib\site-packages\distributed\client.py", line 1840, in _gather
raise exc
CancelledError: list-1db0b54c-c12d-4a3f-a18b-b39ea304e229
```
With 2021.4.1 and later versions the code simply hangs at some point and the ```output``` folder does not contain all of the 240 CSV files.
**Anything else we need to know?**:
I have reported similar issues as https://github.com/dask/dask/issues/7848 and #4990.
**Environment**:
- Dask version: works on 2021.3.0 but an error is reported; 2021.3.1 and 2021.4.0 the code crashes; >= 2021.4.1 the code hangs
- Python version: 3.9.5
- Operating System: Windows Server 2016 64 bits
- Install method (conda, pip, source): conda
Contributor guide
Assessment
This issue has not been assessed yet.