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