У меня есть задача, которая включает в себя обращение к веб-сервису несколько миллионов раз и сохранение ответа в файле (отдельный файл для каждого запроса).
У меня есть высокоуровневая работа. код, но немного запутывает пару вещей.
- В чем разница между двумя синтаксисами?
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