Постобработка многопроцессорного пула Python ⇐ Python
-
Anonymous
Постобработка многопроцессорного пула Python
У меня есть набор задач, которые нужно выполнить (миллионы), и мне нужно POST результат каждого «выполнения». POST предлагает мне возможность группировать результаты.
У меня есть следующий скелет:
def init_pool_processes(): глобальные собранные_данные собранные_данные = [] класс Создать: защита __init__(сам): self.dummy = 0 def __call__(я, сообщение): печать (обработка {сообщения}) # сообщение обработано. Предположим, что результатом является сообщение собранные_данные.append (сообщение) если len(собранные_данные) == 3: # POST собрал данные пакетом из (скажем) 3 и # ясно, так как нам это больше не нужно собранные_данные.очистить() процесс def (я, сообщение): пытаться: пул = Пул (инициализатор = init_pool_processes) пул.карта(я, сообщение) окончательно: пул.закрытие() пул.join() функция функции(): message = ['Это', 'есть', 'конец', 'конец', 'мой', 'единственный', 'друг', 'ты', 'конец', '.'] объект = Создать() obj.process (сообщение) print («Обработка завершена!!!») функция() Теперь очевидно, что некоторые исполнения в конце могут не попасть в вызов POST. Поэтому мне нужен способ вызова функции (один раз) для каждого созданного работника.
Я попробовал добавить следующее после вызова карты:
pool.map(post_process, list(range(pool._processes)), chunksize=1) но это не гарантирует, что функция post_process будет вызвана для всех воркеров.
Есть какие-нибудь советы о том, как этого добиться? Я хочу, если возможно, избежать глобальных решений очередей/таймеров
У меня есть набор задач, которые нужно выполнить (миллионы), и мне нужно POST результат каждого «выполнения». POST предлагает мне возможность группировать результаты.
У меня есть следующий скелет:
def init_pool_processes(): глобальные собранные_данные собранные_данные = [] класс Создать: защита __init__(сам): self.dummy = 0 def __call__(я, сообщение): печать (обработка {сообщения}) # сообщение обработано. Предположим, что результатом является сообщение собранные_данные.append (сообщение) если len(собранные_данные) == 3: # POST собрал данные пакетом из (скажем) 3 и # ясно, так как нам это больше не нужно собранные_данные.очистить() процесс def (я, сообщение): пытаться: пул = Пул (инициализатор = init_pool_processes) пул.карта(я, сообщение) окончательно: пул.закрытие() пул.join() функция функции(): message = ['Это', 'есть', 'конец', 'конец', 'мой', 'единственный', 'друг', 'ты', 'конец', '.'] объект = Создать() obj.process (сообщение) print («Обработка завершена!!!») функция() Теперь очевидно, что некоторые исполнения в конце могут не попасть в вызов POST. Поэтому мне нужен способ вызова функции (один раз) для каждого созданного работника.
Я попробовал добавить следующее после вызова карты:
pool.map(post_process, list(range(pool._processes)), chunksize=1) но это не гарантирует, что функция post_process будет вызвана для всех воркеров.
Есть какие-нибудь советы о том, как этого добиться? Я хочу, если возможно, избежать глобальных решений очередей/таймеров