Dask Distributed - distributed read: scheduler assigns task to wrong worker

Viewed 308

I am currently working on a distributed read on a Dask Distributed cluster. Having all nodes fetch data from a common NFS is working, but my data is already scattered/stored on the local drive of each node. The scheduler runs on a separate node.

This is the code

from dask.distributed import Client
client = Client(scheduler_file='/path/to/scheduler.json')
import socket
nodes = client.run(socket.gethostname)
hosts =  nodes.values() # something like ['node01','node02',...,'node40']
import dask.dataframe as dd

columns = ['VendorID','tpep_pickup_datetime','tpep_dropoff_datetime','passenger_count', 'trip_distance', 'Pickup_longitude', 'Pickup_latitude', 'RatecodeID', 'store_and_fwd_flag', 'Dropoff_longitude', 'Dropoff_latitude',  'payment_type', 'fare_amount', 'extra', 'mta_tax', 'tip_amount', 'tolls_amount', 'improvement_surcharge', 'total_amount']
fdf = [ client.submit(dd.read_csv, '/data/nyc_taxi/yellow*.csv.part+worker[-2:]', workers=worker, names=columns ) for worker in hosts]

Any attempt to work on this data results in a FileNotFoundError: [Errno 2] No such file or directory: '/data/nyc_taxi/yellow*.csv.part30'

The scheduler log gives me:

distributed.scheduler - ERROR - error from worker tcp://192.168.1.141:43079: [Errno 2] No such file or directory: '/data/nyc_taxi/yellow_tripdata_2016-03.csv.part30'

which indicates that the scheduler assigns tasks to the wrong workers. Asking the client client.who_has(fdf[0]) gives me the correct address, though.

I had a look at this blog entry and this wiki entry, but it didnt offer me any useful clue. What am i missing here?

Optimally, i would love to do something like this:

df = dd.from_delayed(fdf) # is this supposed to work on futures?

I am working on dask 0.16.0/distributed 1.20.0, but that error occured on an older dask version as well (around 0.15.4).

0 Answers
Related