get contiguous numbers for workers in a `multiprocessing.Pool`

Viewed 19

GPU models are often ran efficiently by parallelization over batches or features or other tensor dimensions, and theoretically don't need multiprocessing over one device. Nonetheless, from a (developer) productivity point of view it could become efficient to run the same model twice on one GPU device, particularly when memory assignment and GPU utilization allows for it and settings like the batch size are hard to set and require monkey patches. Moreover, forking cleans up GPU-memory, a small nightmare or even impossible to do (yes, I'm looking at you tensorflow) within one process between different stages and models.

However, when you have multiple devices, you'd need to assign the model and data to the right GPU device and keep the spread balanced. I'd like to use the worker's identity or number to assign a device. multiprocessing.current_process()._identity returns a number that seems contiguous, but doesn't have to be when multiple Pools are initializing their workers concurrently as per the example. Is there a way to get the worker number within a multiprocessing.Pool, such that resources can be balanced accordingly?

I admit concurrent Pools don't really occur in practice, but I'm still curious whether this is possible.

Concurrent device assignment over two pools

Figure: Run 1 assigns device 1 to both workers in pool A, and hence unbalances the device spread because of race conditions over _identity.

code in Colab

import time
import random
import pandas as pd
from tqdm.auto import tqdm
from multiprocessing import Pool, current_process


devices = [0, 1] # let's assume we have two devices (e.g. GPUs)
worker, pool_name, device = None, None, None

def initialize_model(pool_name):
  global devices
  worker = current_process()._identity[0]
  device = devices[worker % len(devices)]
  
  globals()['pool_name'] = pool_name
  globals()['device'] = device
  globals()['worker'] = worker


def apply(run):
  global pool_name, device, worker
  # time.sleep(random.random() * 0.04)
  return pool_name, device, worker, run


assignments = []
for run in tqdm(range(20)):
  with Pool(4, initializer=initialize_model, initargs=('A', )) as pool_a, Pool(4, initializer=initialize_model, initargs=('B', )) as pool_b:
    promises = [
      # apply in parallel
      [pool_a, pool_b][random.randint(0, 1)].apply_async(apply, (run, ))
      for _ in range(256)
    ]

    # wait for the results
    results = [promise.get() for promise in promises]
    assignment = pd.DataFrame(results, columns=['pool', 'device', 'worker', 'run']).groupby(['run', 'worker']).first().copy()
    assignments.append(assignment)

pd.options.display.max_columns = 50

device_assignment = pd.concat(assignments).groupby(['run', 'pool'])['device']

pd.DataFrame(dict(
    [(i, device_assignment.apply(lambda x: (x==i).sum())) for i in devices] +
    [('total', device_assignment.count())]
)).rename_axis('device',axis=1).T
0 Answers
Related