У меня много данных, которые я сначала загружаю в оперативную память с помощью SharedMemory, а затем читаю со многими дочерними процессами с помощью multiprocessing.Pool.map< /code>.
Код
Это упрощенная версия (не совсем пример жалобы), которую я использую:
Код: Выделить всё
def SharedObject: # wrapper of SharedMemory passed to subprocesses
def __init__(self, blablabla):
self.shmem = SharedMemory(create=True, size=numbers) # reference to shmem will live as long as the instance of the class will (shmem not getting GC)
temp_arr = np.ndarray(self.shape, dtype=self.dtype, buffer=self.shmem.buf)
temp_arr[:] = ...lot of data... # this array will get destroyed after __init__ finishes
def __getitem__(self, indices) -> np.ndarray: # used by subprocesses
selected = np.ndarray(self.shape, dtype=self.dtype, buffer=self.shmem.buf)
return selected.__getitem__(indices)
# This is in main process
shobj = SharedObject()
with multiprocessing.Pool() as pool:
result= list(pool.map(f, shobj)) # f does shobj.__getitem__
У меня на компьютере 64 ГБ оперативной памяти. Если я запускаю описанный выше алгоритм с небольшим количеством данных, все работает гладко, но если я загружаю много данных (~ 40 ГБ), я получаю следующую ошибку:
Код: Выделить всё
n __enter__
return self._semlock.__enter__()
^^^^^^^^^^^^^^^^^^^^^^^^^
File "C:\Users\mnandcharvim\AppData\Local\Programs\Python\Python312\Lib\multiprocessing\connection.py", line 321, in _recv_bytes
waitres = _winapi.WaitForMultipleObjects(
Данные доступны только для чтения, поэтому для меня было бы лучше, если бы я мог загрузить их в файл, доступный только для чтения. часть памяти, чтобы не было блокировок. Этот вопрос указывает на то, что SharedMemory не блокируется, но на данный момент, судя по полученной ошибке, я не уверен).
2-я версия
Я также попытался привести код в соответствие с примером официальной документации, на которую я дал ссылку:
Код: Выделить всё
shmems = [] # module variable (same module of SharedObject) as suggested in some answers
def SharedObject:
def __init__(self, blablabla):
shmem = SharedMemory(name=self.name, create=True, size=numbers) # this reference destroyed after __init__
shmems.append(shmem) # but this one outlives it
temp_arr = np.ndarray(self.shape, dtype=self.dtype, buffer=self.shmem.buf)
temp_arr[:] = ...lot of data...
def __getitem__(self, indices) -> np.ndarray: # used by subprocesses
shmem = SharedMemory(name=self.name) # added line
selected = np.ndarray(self.shape, dtype=self.dtype, buffer=shmem.buf)
return selected.__getitem__(indices)
# This is in main process
shobj = SharedObject()
with multiprocessing.Pool() as pool:
result= list(pool.map(f, shobj)) # f does shobj.__getitem__
Эта вторая версия выдает мне ошибку: shmem = SharedMemory(name=self.name) # добавлена строка, говорящая что мне не хватает системных ресурсов для выполнения mmap (странно, ведь данные уже были отображены в оперативной памяти: ресурсов хватило на якобы первую и единоразовую загрузку).О чем следует помнить
- Обработка данных предполагает только их чтение
- I не могу использовать потоки (поскольку при чтении контента я выполняю вычисления, не освобождающие GIL).
- Я просто хочу знать четкий, простой и функциональный способ, для которого я должен был бы использовать SharedMemory в сочетание с массивами numpy и многопроцессорностью.
Подробнее здесь: https://stackoverflow.com/questions/790 ... ng-probelm