TLDR;
- I want to read many files from disk in parallel
- My code loads these files very quickly, but merging them into the main process' memory is slow
- 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