Understanding the deserialization time in dask distributed

Hi, I am trying to benchmark the performance of an HPC Cluster using Dask’s SLURMCluster. Thanks to the help I previously received from you, I managed to make it scale well on my cluster.

There is one issue that remained unsolved. I realize that if I initiliaze a scheduler and I ask it to perform multiple times the same task for me, then the first run is going to be several seconds slower.

I am attaching the results from the 3 runs on the same cluster:
Run 1 → Dask Performance Report
Run 2 → Dask Performance Report
Run 3 → Dask Performance Report

I read in the report that the first run took several seconds in “deserialize time”, and I wonder from where this extra time is added. If I read the Task Stream, I can confirm that only in the first run, there is a “deserialize-dask_mapper”.

What I have tried was to import ROOT as an external module, by adding relative paths as:


I again produced 3 reports for each run:
Run 1 → Dask Performance Report
Run 2 → Dask Performance Report
Run 3 → Dask Performance Report

For runs 2,3 reports are as before. But I see that in report 1, there is no time labeled as “deserialize time”, nor there is a “deserialize-dask_mapper” in the task stream. And again the first run took few seconds longer than the other runs.

Question 1: I would like to ask you what might be the cause of the slower first runs.
Question 2: What is the transfer time in the summary of the report? I see that it is minimal for the first run, which seems unreasonable to me.

Deseriailzation time is the time spent to recreate the function that you want to run on a remote machine. A common cause of initial long deserialization time followed by short (or non-existent) deserialization time is that the library you’re using takes a while to import the first time it is used. Subsequent imports are fast because Python keeps libraries in memory.

Import times can be particularly challenging on HPC systems when code is mounted on a network file system (NFS) because many workers all try to import the same code at once and the NFS thrashes a bit. Typically people just live with this. However, if you want to avoid it you could do a few things:

  1. Install libraries onto local disk (if it exists)
  2. Import the libraries when you start up the Dask workers (things will still be slow, but you’ll move around the slowness if that’s helpful). One way to do this would be to use the --preload my_library flag when starting a dask-worker

Hopefully that helps?

1 Like