.map followed by .result() and/or client.gather() results in a crash if workers are remote

Viewed 185

I have a need to call a function which does complicated things; and which needs to be iterated over a few hundred thousand times (think monte carlo simulation). I'm here because it didn't work with dask in my setup (first time user). So, i reduced the code to barebones (see below) to determine whether its the human, the environment, or both :-). Hoping the gurus here could help me out?

code

def do(x):
    y = x**2.0 
    return y
client = Client("tcp://10.61.68.34:8786")
x=[]
x = client.map(do, range(0, 1799))
r =[]
for i in x:
    r.append(i.result())
print(r)
del x

SETUP
Scheduler: ip1, vmware vm, windows os, 128 cores, 200GB RAM, py 3.9, dask 2021.9.1
Worker: same as above

Network: intra- node (within 1 single node), Inbound rule allows incoming TCP connections to 8786, 8787. No specific rule for connections to ephemeral port is required because scheduler will communicate back to workers using the worker's ephemeral port used for TCP SYN packet.

launching scheduler using CLI:

dask-scheduler.exe

launching worker using CLI:

dask-worker.exe --nprocs 100 --nthreads 1 tcp://<<ip1..>>:8786

dask.yml has:

distributed:
  version: 2

  adaptive:
    interval: 1s
    target-duration: 15m

  admin:
    tick:
      interval: 20ms 
      limit: 300s 
  
  comm:
    retry:
      count: 5 
      offload: 100MiB
    
    zstd:
      level: 15 
      threads:0   # need to respect Python's GIL as not all delayed code handles dask objects
  
  deploy:
    lost-worker-timeout: 30s  
    cluster-repr-interval: 750ms 

  logging:
      distributed: debug
      distributed.client: debug
      bokeh: critical
      tornado: critical
      tornado.application: error

  scheduler:
    bandwidth: 10000000000    # 10 Gb/s even though LAN is capable of 100Gbps
    default-data-size: 100kiB
    events-cleanup-delay: 300s
    idle-timeout: 1h    
  
  worker:
    connections:            # Maximum concurrent connections for data
      outgoing: 300          # This helps to control network saturation
      incoming: 300
    lifetime:
      duration: 1800s  

in iPython (from VSCode-- environment is set right), once map completes

after <Future: finished, type: float, key: do-b4e893e37f05790c2d90f6f7a57bc865>,

calling x[0].result() yields expected result (1.0)
calling x[1].result() yields expected result (4.0)
calling custom range

for i in range(2,10):
 y.append(x[i].result())
print(y)

also prints expected answer.
CONCLUSION: code functions as desired

Now, when workers are launched froma different machine (no local workers present anymore)
Scheduler: ip1, vmware vm, windows os, 128 cores, 200GB RAM, dask 2021.9.1
Worker: ip2, vmware vm, windows os, dask 2021.9.1
Noteworthy mentions: 100Gbps network link between these two machines, firewall allows 8786, 8786 inbound on both

iPython output when calling the for loop; it goes on ad-inifinitium (until manual interruption):

distributed.client - DEBUG - Waiting on futures to clear before gather
distributed.comm.tcp - DEBUG - Setting TCP keepalive: idle=10, interval=2
distributed.comm.tcp - DEBUG - Setting TCP keepalive: idle=10, interval=2
distributed.comm.tcp - DEBUG - Setting TCP keepalive: idle=10, interval=2
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-9755ae4fe21461f71d76472d06ed69b7'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-4393edb8f90330ff85a14a347b57275f'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-b2354f9980c5f46e83153d4f2583f6f1'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-fa2cd7fe07c50c80dda8b5a308a7655c'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-80e7f4ee2b14bc0276a09b77db1e9d43'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-16311a9900ff7294aef0f0f3432b82fb'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-8ba4db219b30ddb8ee8ed086c752e72c'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-37eb0abfc3030871834fbb4ce6fdcf6e'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-5f894a60cd293af00b628a1c9d6972d0'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-e762686340db94c9b9617ba959694ecf'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-4e270441aebb8ed2a00f48ca3fc0391b'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-b1bd0b36b8e3b282d2caf42d373ccedd'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-a93dc9cc91d1a7e3fa7bc88673b7e33e'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-68324d1e2f16f3dd3ed2a4521280a1f1'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-94e3a7bdb4065cf34d47f6667ed5a695'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-0dbc31f32cb02f4925ca30ac92064eb1'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-70336d000742ece9b066b6c111f20fd3'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-ecd6c5c5caf0fe8d5499576f9f53764e'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-cb1d841d06994f4f5cba045e2c65e881'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-3015b21187427be5873d6477544063ae'}
distributed.client - DEBUG - Client receives message {'op': 'lost-data', 'key': 'do-f63da4cc3e1e26d64cc927f622d2631b'}
...

scheduler's sample output:

distributed.core - DEBUG - Message from 'tcp://<i3p>:53905': {'op': 'heartbeat_worker', 'address': 'tcp://<ip2>:53784', 'now': 1633623013.8946366, 'metrics': {'executing': 0, 'in_memory': 29, 'ready': 0, 'in_flight': 0, 'bandwidth': {'total': 10000000000, 'workers': {}, 'types': {}}, 'spilled_nbytes': 0, 'cpu': 0.0, 'memory': 104374272, 'time': 1633623013.8877068, 'read_bytes': 6227.017170580167, 'write_bytes': 19296.954301894224, 'read_bytes_disk': 0.0, 'write_bytes_disk': 10238.384044113825}, 'executing': {}, 'reply': True}
distributed.core - DEBUG - Calling into handler heartbeat_worker

bokeh status shows no errors.

lets pick a sample task from above logs to check its status @ worker:

distributed.worker - DEBUG - Send compute response to scheduler: do-b6553fa57215f934079c36275fe3985d, {'op': 'task-finished', 'status': 'OK', 'nbytes': 24, 'type': <class 'float'>, 'start': 1633623109.3889816, 'stop': 1633623109.3890078, 'thread': 14244, 'key': 'do-b6553fa57215f934079c36275f**e3985d**'}

once all tasks have state 'task-finished' at worker, debug messages only indicate normal heartbeat messages nothing special.

But, no result! Both scheduler + worker logs show activity, but thats it.

local output (iPython or VSCode run) shows:

    distributed.client - WARNING - Couldn't gather 1 keys, rescheduling {'do-473302f1b42c06c5a75fe179c7ffe473': ('tcp://10.61.68.33:54198',)}
distributed.client - DEBUG - Waiting on futures to clear before gather
distributed.client - DEBUG - Client receives message {'op': 'key-in-memory', 'key': 'do-473302f1b42c06c5a75fe179c7ffe473'}
distributed.client - WARNING - Couldn't gather 1 keys, rescheduling {'do-473302f1b42c06c5a75fe179c7ffe473': ('tcp://10.61.68.33:54264',)}
distributed.client - DEBUG - Waiting on futures to clear before gather
distributed.client - DEBUG - Client receives message {'op': 'key-in-memory', 'key': 'do-473302f1b42c06c5a75fe179c7ffe473'}
distributed.client - WARNING - Couldn't gather 1 keys, rescheduling {'do-473302f1b42c06c5a75fe179c7ffe473': ('tcp://10.61.68.33:54198',)}
distributed.client - DEBUG - Waiting on futures to clear before gather
distributed.client - DEBUG - Client receives message {'op': 'key-in-memory', 'key': 'do-473302f1b42c06c5a75fe179c7ffe473'}

and after few re-tries, it gives up, throws an exception, and program ends. This happens every single time i run it with a remote worker.

Behavior gets really weird if same code is re-factored using .persist() followed by client.gather(). In that case, workers get into, upon task completion, 'task forgotten' state and force-close the connection. I am happy to start a separate thread on it. But, i'm sure its something very basic that is tripping me here. Help would be much appreciated

0 Answers
Related