Многопроцессорность Python, как избежать создания кортежа с миллионом объектовPython

Программы на Python
Anonymous
Многопроцессорность Python, как избежать создания кортежа с миллионом объектов

Сообщение Anonymous »

Новичок в многопроцессорной обработке Python.
У меня есть задача, которая включает в себя обращение к веб-сервису несколько миллионов раз и сохранение ответа в файле (отдельный файл для каждого запроса).
У меня есть высокоуровневая работа. код, но немного запутывает пару вещей.
  • В чем разница между двумя синтаксисами?
    pool = Pool(processes=4)
    pool.starmap(task, listOfInputParametersTuple)
и
with Pool(processes=4) as pool:
pool.starmap(task, listOfInputParametersTuple)
  • Есть ли способ избежать чтения всего входного файла перед запуском многопроцессорного пула? По сути, прочитайте каждую строку и немедленно создайте задачу пула.
    with open (list_of_ids, 'r') as infile:
    for id in infile:
    listOfInputParametersTuple.append( tuple ((id, queue, requestBodyTemplate))
    pool.starmap(task, listOfInputParametersTuple)
  • В задаче я использую довольно большой requestBodyTemplate и передаю идентификатор для создания запроса JSON. Эта переменная requestBodyTemplate дублируется в каждом элементе входного кортежа. Есть ли способ передать его в функцию задачи за пределами входного кортежа.
  • Как гарантировать, что все порождены задачи завершаются до выхода из основной программы?
  • Есть какие-нибудь указания по тайм-ауту выполнения задачи, если освобождение пула занимает слишком много времени?
Не стесняйтесь поделиться любыми другими предложениями, которые могут у вас возникнуть.
Спасибо.
Цель:
  • Прочитать список из 1 миллиона идентификаторов из файла
    Цель:

    Прочитать из файла список из 1 миллиона идентификаторов
    li>
    Для каждого идентификатора создайте JSON.
  • Запросите веб-сервис с созданным выше JSON. (получение ответа занимает несколько секунд)
  • Сохранить вывод ответа в каталог.
    Сохранить статус в единый файл «статуса».

Задача
-
def task(i,queue):

# get the current process
process = current_process()

# generate some work
s = random.randint(1, 10)

# block to simulate work
print (f"TASK function : {process} - sleep for {s} sec")
sleep(s)

data = f"{process} - sleep {s} sec - {i} - {queue}"

print (f"TASK function: {data}")

# put it on the queue
queue.put(data)

Основной метод
def main ():

set_start_method('spawn')
pool = Pool(processes=4)

requestBodyTemplate = getRequestBodyTemplateJSON()

with open (list_of_ids, 'r') as infile:
for line in infile:
listOfInputParametersTuple.append( tuple ((line, queue , requestBodyTemplate))

# set the fork start method
# create the manager
with Manager() as manager:
# create the shared queue
queue = manager.Queue()

Process(target=listener, args=(queue,)).start()
print ("back in main after starting listener ")

# execute the tasks_i in parallel # use starmap to have multiple params
pool.starmap(task, listOfInputParametersTuple)
# pool.starmap(task, zip ( args_i, itertools.repeat(queue)))

pool.close()
pool.join()

# wait for all tasks to get over
sleep(10)

print ("\n Sending None message to queue ")
queue.put(None)


Подробнее здесь: https://stackoverflow.com/questions/788 ... on-objects

Вернуться в «Python»