Loading Large Datasets in Parallel With Python Multiprocessing

Viewed 130

TLDR;

  1. I want to read many files from disk in parallel
  2. My code loads these files very quickly, but merging them into the main process' memory is slow
  3. Am I doing this correctly? Is there a better way to share memory across processes?

Background: I'm currently working on a project with a private hospital dataset (only 3 GB!) and would like to be able to load it faster (their disks are slow). I work on a 36 CPU server, so I should be able to speed it up using multiple processes. However, I am having trouble seeing an improvement.

The setup: I have split my dataset over 12 binary PyTorch files ('.pt') which I load with torch.load(filename). These files each contain a list of images inside a dict E.G. the 0th list: {0:ListOfImages}. I would like to load these files simultaneously with different processes, then merge them into a list in their original order. I am currently able to do that, but cannot seem to merge the data into the main process faster than it would take to load sequentially.

Sequential baseline load: 74 seconds

What I tried:

  • multiprocessing.Pool().map() (150 Seconds)
  • multiple iterations of process.start() + process.join() (300 Seconds +)
  • multithreading: I load the dataset in exactly the same amount of time as sequentially (74 Seconds)

Multiprocessing code thus far:

import multiprocessing as mp
import time
CORES = 12

def loadWorker(sharedDict,filename):
  t1 = time.time()
  temp = torch.load(filename)
  print("Dict Loaded in :{} seconds".format(time.time()-t1))
  t1=time.time()
  sharedDict.update(temp)
  print("Dict updated in :{} seconds".format(time.time()-t1))
  

filenames = ["namingconvention{}.pt".format(i) for i in range(CORES)]

jobs=[]
manager = mp.Manager()
sharedDict = manager.dict()
t1 = time.time()
for i,filename in enumerate(filenames):
  print("starting Process: {}".format(i))
  p = mp.Process(target=loadWorker, args=(sharedDict,filename))
  jobs.append(p)
  p.start()

for i,proc in enumerate(jobs):
  print("join Process {}".format(i))
  process.join()

print("dict created of length: {} in {} seconds".format(
  len(sharedDict),time.time()-t1)
)
  
images = []
for i in range(CORES):
  print("Appending Images: {}".format(i))
  images = images + list(images)

print("Images appended in {} seconds".format(time.time()-t1))

Am I doing something wrong? Is there a better way to share memory between these processes?

OUTPUT:

The output of the above: Image of the stdout

The output of the above code with CORES=1 and logger.setLevel(mp.SUBDEBUG) Image of debug stdout

0 Answers
Related