Не работает multiprocessing pool python

Python multiprocessing pool.map doesn’t work parallel

I wrote a simple parallel python program

What I expect to see in the output:

What I actually see:

When I’m looking in the output It’s look like pool.map running a function and waits till it’s done and then run another, but when I calculate the duration of whole program its about 2 seconds and it’s impossible unless the test_function is running parallel

This code is working well in MacOS and Linux but It’s not showing the expected output on windows 10. python version is 3.6.4

2 Answers 2

The multiprocessing.Pool() documentation ( since ever, Py27 incl. ) is clear in intentionally blocking in processing the queue-of-calls as created by the iterator-generated set of the just -4- calls, produced sequentially from the above posted example.

The multiprocessing -module documentation says this about its Pool.map() method:

A parallel equivalent of the map() built-in function (it supports only one iterable argument though). It blocks until the result is ready.

This should be the observed behaviour, whereas different instantiation methods would accrue different add-on ( process copying-related ) overhead costs.

Anyway, the mp.cpu_count() need not be the number of CPU-cores any such dispatched .Pool() -instance workers’ tasks will get on to get executed, because of the O/S ( user/process-related restriction policies ) settings of affinity:

Your code will have to «obey» the sub-set of those CPU-cores, that are permitted to be harnessed by any such multiprocessing -requested sub-process,
the number of which is not higher than: len( os.sched_getaffinity( 0 ) )

The Best Next Step : re-evaluate your whole code-execution eco-system

might look about something like this on linux-class O/S:

Check the hidden detail — what your O/S uses for invoking the test_function() — the mapstar() ( not being a sure choice universally ) was the local SMP-linux-class O/S’s choice for its default sub-process instantiation method, performed via ‘ fork ‘.

Источник

Python multiprocessing example not working

I am trying to learn how to use multiprocessing but I can’t get it to work. Here is the code right out of the documentation

it should output

but instead i get

no errors or other messages, it just sits there, It is running in IDLE from a saved .py file on a Windows 7 machine with the 32-bit version of Python 2.7

6 Answers 6

My guess is that you are using IDLE to try to run this script. Unfortunately, this example will not run correctly in IDLE. Note the comment at the beginning of the docs:

Note Functionality within this package requires that the main module be importable by the children. This is covered in Programming guidelines however it is worth pointing out here. This means that some examples, such as the multiprocessing.Pool examples will not work in the interactive interpreter.

The __main__ module is not importable by children in IDLE, even if you run the script as a file with IDLE (which is commonly done with F5).

The problem is not IDLE. The problem is trying to print to sys.stdout in a process that has no sys.stdout. That is why Spyder has the same problem. Any GUI program on Windows is likely to have the same problem.

On Windows, at least, GUI programs are usually run in a process without stdin, stdout, or stderr streams. Windows expects GUI programs to interact with users through widgets that paint pixels on the screen (the G in Graphical) and receive key and mouse events from Windows event system. That is what the IDLE GUI does, using the tkinter wrapper of the tcl tk GUI framework.

Читайте также:  Справка windows не работает

When IDLE runs user code in a subprocess, idlelib.run runs first, and it replaces None for the standard streams with objects that interact with IDLE itself through a socket. Then it exec()s user code. When the user code runs multiprocessing, multiprocessing starts further processes that have no std streams, but never get them.

The solution is to start IDLE in a console: python -m idlelib.idle (the .idle is not needed on 3.x). Processes started in a console get std streams connect to the console. So do further subprocesses. The real stdout (as opposed to the sys.stdout) of all the processes is the console. If one runs the third example in the doc,

then the ‘main line’ block goes to the IDLE shell and the ‘function f’ block goes to the console.

This result shows that Justin Barber’s claim that the user file run by IDLE cannot be imported into processes started by multiprocessing is not correct.

EDIT: Python saves the original stdout of a process in sys.__stdout__ . Here is the result in IDLE’s shell when IDLE is started normally on Windows, as a pure GUI process.

Here is the result when IDLE is started from CommandPrompt.

The standard file numbers for stdin, stdout, and stderr are 0, 1, 2. Run a file with

in IDLE started in the console and the output is the same.

Источник

Python/Multiprocessing: процессы не запускаются

У меня есть функция, которая читает двоичный файл и преобразует каждый байт в соответствующую последовательность символов. Например, 0x05 становится «AACC», 0x2A становится «AGGG» и т. Д. Функция, которая читает файл и преобразует байты, в настоящее время является линейной, и поскольку файлы для преобразования находятся в диапазоне от 25 КБ до 2 МБ, это может занять довольно много времени. какое-то время.

Поэтому я пытаюсь использовать многопроцессорность, чтобы разделить задачу и, надеюсь, повысить скорость. Тем не менее, я просто не могу заставить его работать. Ниже приведена линейная функция, которая работает, хотя и медленно;

Примечание. Функция «decSequenceToRNA» принимает данные из буфера и преобразует каждый байт в требуемую строку. После выполнения функция возвращает кортеж, который содержит номер блока и строку, например (1, ‘ACCGTAGATTA. ‘), и в конце у меня есть массив этих доступных кортежей.

Я пытался преобразовать функцию, чтобы использовать многопроцессорность Python;

Однако ни один из процессов, кажется, даже не запускается, так как при запуске этой функции возвращается пустой массив. Любое сообщение, напечатанное на консоли в ‘decSequenceToRNA‘, не отображается;

В отличие от этого вопроса, я использую Linux shiva 3.14-kali1-amd64 #1 SMP Debian 3.14.5-1kali1 (2014-06-07) x86_64 GNU/Linux и использую PyCrust для тестирования функций в версии Python: 2.7.3, Я использую следующие пакеты:

Я хотел бы помочь выяснить, почему мой код не работает, если я что-то упускаю в другом месте, чтобы заставить Процесс работать. Также открыты предложения по улучшению кода. Ниже ‘decSequenceToRNA‘ для справки:

2 ответа

decSequenceToRNA выполняется в своем собственном процессе, что означает, что он получает свою собственную, отдельную копию каждой структуры данных в основном процессе. Это означает, что когда вы добавляете _rnaSequences в decSequenceToRNA не влияет на rnaSequences в родительском процессе. Это объясняет, почему возвращается пустой список.

У вас есть два варианта решения этой проблемы. Во-первых, это создать list которые могут быть разделены между процессами, используя multiprocessing.Manager , Например:

Применение этого к вашему коду просто потребует замены

В качестве альтернативы, вы можете (и, вероятно, должны) использовать multiprocessing.Pool вместо создания индивидуального Process за каждый кусок. Я не уверен, насколько большой hFile или насколько большие куски вы читаете, но если их больше multiprocessing.cpu_count() куски, вы будете снижать производительность, порождая процессы для каждого чанка. Используя Pool , вы можете поддерживать постоянный счетчик процессов и легко создавать свои rnaSequence список:

Обратите внимание, что мы больше не передаем rnaSequences список ребенку. Вместо этого мы просто возвращаем результат, который мы добавили бы к родителю (что мы не можем сделать с Process ), и построить список там.

Читайте также:  Не работает фитнес браслет как вернуть

Источник

Класс Pool() модуля multiprocessing в Python.

Создание, запуск и получение результатов от пула процессов.

Синтаксис:

Параметры:

  • processes — количество используемых рабочих процессов,
  • initializer — вызываемый объект (функция),
  • initargs — аргументами для initializer
  • maxtasksperchild — количество задач рабочего процесса до обновления,
  • context — контекста для запуска рабочих процессов.

Возвращаемое значение:

Описание:

Класс Pool() модуля multiprocessing создает объект, управляющий пулом рабочих процессов, в который могут быть отправлены задания. Пул рабочих процессов поддерживает асинхронное выполнение задач с тайм-аутами и обратными вызовами и имеет параллельную реализацию.

Аргумент processes — это количество используемых рабочих процессов. Если аргумент processes не указан, то используется число, возвращаемое функцией os.cpu_count() .

Если аргумент initializer не равен None , то при запуске каждый рабочий процесс будет выполнять вызываемый объект initializer с аргументами *initargs в виде: initializer(*initargs)

Аргумент maxtasksperchild — это количество задач, которые рабочий процесс может выполнить до того, как он выйдет и будет заменен новым рабочим процессом, чтобы освободить неиспользуемые ресурсы. По умолчанию для maxtasksperchild установлено значение None , это означает, что рабочие процессы будут жить столько же, сколько и пул.

Аргумент context можно использовать для указания контекста, используемого для запуска рабочих процессов. Обычно пул создается с помощью создания экземпляра класса multiprocessing.Pool() или объекта контекста context.Pool() , созданного функцией multiprocessing.get_context() . В обоих случаях контекст устанавливается соответствующим образом.

Обратите внимание, что методы объекта пула должны вызываться только процессом, создавшим пул. Проще говоря методы должны вызываться из кода, где был создан и запущен пул рабочих процессов.

Предупреждение. Объекты multiprocessing.Pool имеют внутренние ресурсы, которыми необходимо правильно управлять, используя пул с менеджером контекста или вручную вызывая методы пула Pool.close() и Pool.terminate() . Невыполнение этого требования может привести к зависанию процесса при завершении.

Обратите внимание, что неправильно полагаться на сборщик мусора для уничтожения пула, поскольку CPython не гарантирует, что будет вызван финализатор пула.

Примечание. Рабочие процессы в пуле обычно живут в течение всего срока рабочей очереди пула. Часто используемый в других системах (например, Apache, mod_wsgi и т. д.) шаблон для освобождения ресурсов, удерживаемых рабочими процессами, заключается в том, чтобы позволить им выполнить только заданный объем работы перед его выходом, очисткой и созданием нового процесса. Для пула процессов, эту возможность пользователю предоставляет аргумент maxtasksperchild .

Методы объекта Pool .

  • Pool.apply() вызывает функцию с аргументами,
  • Pool.apply_async() асинхронный вариант метода Pool.apply() ,
  • Pool.map() многопроцессорный эквивалент встроенной функции map() ,
  • Pool.map_async() асинхронный вариант метода Pool.map() ,
  • Pool.imap() более ленивая версия метода Pool.map() ,
  • Pool.imap_unordered() то же самое, что и Pool.imap() , только результаты идут по готовности,
  • Pool.starmap() аналогичен методу Pool.map() , только другая передача аргументов,
  • Pool.starmap_async() комбинация методов Pool.starmap() и Pool.map_async() ,
  • Pool.close() предотвращает отправку задач в пул,
  • Pool.terminate() останавливает рабочие процессы,
  • Pool.join() ждет, пока рабочие процессы закончатся,
  • Объект AsyncResult результат вызовов методов Pool.apply_async() и Pool.map_async()
    • AsyncResult.get() возвращает результат, как только он придет,
    • AsyncResult.wait() ждет, пока будет доступен результат,
    • AsyncResult.ready() проверяет, завершился ли вызов,
    • AsyncResult.successful() проверяет, был ли завершен вызов без исключения.
  • Пример создания и запуска и использования пула рабочих процессов,
  • Тестирование основной функциональности объекта Pool .

Pool.apply(func[, args[, kwds]]) :

Метод Pool.apply() вызывает функцию func с аргументами args и ключевыми аргументами kwds . Блокирует, пока не будет готов результат.

С учетом этой блокировки, для параллельной работы лучше подходит метод Pool.apply_async() . Кроме того, аргумент метода func выполняется только в одном рабочем процессе пула.

Pool.apply_async(func[, args[, kwds[, callback[, error_callback]]]]) :

Метод Pool.apply_async() асинхронный вариант метода Pool.apply() , который возвращает объект AsyncResult .

Если указан обратный вызов callback , то это должен быть вызываемый объект, принимающий единственный аргумент. Когда результат становится готовым, к нему применяется этот обратный вызов. Если вызов callback не удался, то в этом случае вместо него применяется error_callback .

Если указан error_callback , то это должен быть вызываемый объект, который принимает единственный аргумент. Если целевая функция обратного вызова терпит неудачу, то вызывается error_callback с экземпляром исключения в качестве аргумента.

Обратные вызовы должны завершаться немедленно, в противном случае асинхронный поток, обрабатывающий результаты, будет заблокирован.

Pool.map(func, iterable[, chunksize]) :

Метод Pool.map() представляет собой параллельный эквивалент встроенной функции map() .

Читайте также:  Indesit wiun 102 не работает

Аргумент метода iterable представляет из себя итерацию, элементы которой не распаковываются при передаче в функцию func . То есть функция func(x) , может иметь только один аргумент x . Другими словами, если iterable=[(1, 2), (10, 20), . ] , то запуск задач будет происходить следующим образом: [func((1, 2)), func((10, 20)), . ] , где x будет равен кортежу, например (1, 2) .

Если необходимо использовать функцию с несколькими аргументами, например func(x, y) и распаковывать итерации при запуске, например: [func(1, 2), func(10, 20), . ] , то посмотрите в сторону метода Pool.starmap() .

Метод блокирует выполнение программы до получения результатов работы всеми запущенными процессами.

Метод Pool.map() разбивает итерируемый объект на несколько частей, которые отправляет в пул процессов как отдельные задачи. Приблизительный размер этих фрагментов можно указать, задав для аргумента метода chunksize положительное целое число.

Обратите внимание, что Pool.map() может привести к высокому использованию памяти для очень длинных итераций. Для большей эффективности рассмотрите возможность использования метода Pool.imap() или Pool.imap_unordered() с явной опцией chunksize .

Pool.map_async(func, iterable[, chunksize[, callback[, error_callback]]]) :

Метод Pool.map_async() вариант метода Pool.map() , который возвращает объект AsyncResult .

Если указан обратный вызов callback , то это должен быть вызываемый объект, принимающий единственный аргумент. Когда результат становится готовым, к нему применяется этот обратный вызов. Если вызов callback не удался, то в этом случае вместо него применяется error_callback .

Если указан error_callback , то это должен быть вызываемый объект, который принимает единственный аргумент. Если целевая функция обратного вызова терпит неудачу, то вызывается error_callback с экземпляром исключения в качестве аргумента.

Обратные вызовы должны завершаться немедленно, в противном случае асинхронный поток, обрабатывающий результаты, будет заблокирован.

Pool.imap(func, iterable[, chunksize]) :

Метод Pool.imap() представляет из себя более ленивую версию метода Pool.map() .

Аргумент метода chunksize совпадает с аргументом, используемым методом Pool.map() . Для очень длинных итераций использование большого значения chunksize может значительно ускорить выполнение задания, чем использование значения по умолчанию 1.

Также, если chunksize=1 , то метод итератора it.next() , возвращаемый методом Pool.imap() , будет иметь необязательный параметр тайм-аута: .next(timeout) . Этот метод итератора вызовет TimeoutError , если результат не может быть возвращен в течение timeout секунд.

Pool.imap_unordered(func, iterable[, chunksize]) :

Метод Pool.imap_unordered() то же самое, что и Pool.imap() , за исключением того, что порядок результатов возвращаемого итератора следует считать произвольным.

Только когда есть только один рабочий процесс, порядок гарантированно будет «правильным».

Pool.starmap(func, iterable[, chunksize]) :

Метод Pool.starmap() аналогичен методу Pool.map() , за исключением того, что элементы iterable , как ожидается, будут то же итерируемыми, которые будут распаковываться в качестве аргументов.

Например, если iterable=[(1, 2), (3, 4)] , то результат передачи аргументов в вызываемый объект func будет таким: [func(1, 2), func(3, 4)] .

Pool.starmap_async(func, iterable[, chunksize[, callback[, error_callback]]]) :

Метод Pool.starmap_async() представляет собой комбинацию методов Pool.starmap() и Pool.map_async() .

Метод выполняет итерацию по iterable , распаковывает их и вызывает функцию func с распакованными аргументами из iterable . Возвращает объект результата.

Pool.close() :

Метод Pool.close() предотвращает отправку задач в пул. Как только все задачи будут выполнены, рабочие процессы завершатся.

Pool.terminate() :

Метод Pool.terminate() останавливает рабочие процессы немедленно, не давая завершить невыполненную работу. Метод будет вызван немедленно, когда объект пула будет обработан сборщиком мусора.

Pool.join() :

Метод Pool.join() ждет, пока запущенные рабочие процессы закончатся. Перед использованием Pool.join() необходимо вызвать Pool.close() или Pool.terminate() .

Объект AsyncResult .

Объект AsyncResult представляет собой результат, возвращаемый методами Pool.apply_async() и Pool.map_async() .

Объект AsyncResult определяет следующие методы:

AsyncResult.get([timeout]) :

Метод AsyncResult.get() возвращает результат, как только он придет.

Если тайм-аут timeout не равен None и результат не приходит в течение этого тайм-аута, то возникает исключение TimeoutError . Если удаленный вызов вызвал исключение, это исключение будет вызвано повторно.

AsyncResult.wait([timeout]) :

Метод AsyncResult.wait() ждет, пока будет доступен результат или пока не пройдет время таймаута timeout .

AsyncResult.ready() :

Метод AsyncResult.ready() проверяет, завершился ли вызов.

AsyncResult.successful() :

Метод AsyncResult.successful() проверяет, был ли завершен вызов без возникновения исключения. Поднимет исключение ValueError , если результат не готов.

Изменено в Python 3.7: Если результат не готов, вместо AssertionError возникает ValueError .

Пример создания, запуска и получение результатов от пула процессов:

В следующем примере показано использование пула:

Тестирование основной функциональности объекта Pool .

Вывод теста основной функциональности и методов объекта Pool очень длинный и приводиться не будет. Кому интересно, скопируйте код и запустите сами.

Источник

Оцените статью