Spec-Zone.ru › Python 3.7

concurrent.futures — Запуск параллельных задач

Новое в версии 3.2.

Исходный код: Lib/concurrent/futures/thread.py и Lib/concurrent/futures/process.py

Модуль concurrent.futures предоставляет высокоуровневый интерфейс для асинхронного выполнения вызываемых объектов.

Асинхронное выполнение может быть осуществлено с использованием потоков, с помощью ThreadPoolExecutor, или отдельных процессов, с помощью ProcessPoolExecutor. Оба реализуют одинаковый интерфейс, который определяется абстрактным классом Executor.

Объекты Executor

class concurrent.futures.Executor

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

submit(fn, *args, **kwargs)

Распределяет вызываемый объект, fn, для выполнения как fn(*args **kwargs) и возвращает объект Future, представляющий выполнение вызываемого объекта.

with ThreadPoolExecutor(max_workers=1) as executor:
    future = executor.submit(pow, 323, 1235)
    print(future.result())
map(func, *iterables, timeout=None, chunksize=1)

Аналогично map(func, *iterables), за исключением:

  • iterables собираются немедленно, а не лениво;
  • func выполняется асинхронно и несколько вызовов func могут быть выполнены одновременно.

Возвращаемый итератор генерирует исключение concurrent.futures.TimeoutError, если __next__() вызывается, и результат недоступен после timeout секунд с момента первоначального вызова Executor.map(). timeout может быть целым или вещественным числом. Если timeout не указан или None, время ожидания не ограничено.

Если вызов func генерирует исключение, то это исключение будет сгенерировано при получении его значения из итератора.

При использовании ProcessPoolExecutor, этот метод разбивает iterables на несколько частей, которые он отправляет в пул как отдельные задачи. Примерный размер этих частей можно указать, установив chunksize в положительное целое число. Для очень длинных iterables использование большого значения для chunksize может значительно улучшить производительность по сравнению с значением по умолчанию 1. При использовании ThreadPoolExecutor, chunksize не имеет эффекта.

Изменено в версии 3.5: Добавлен аргумент chunksize.

shutdown(wait=True)

Указывает исполнителю, что он должен освободить все используемые ресурсы, когда текущие ожидающие задачи завершат выполнение. Вызовы Executor.submit() и Executor.map(), сделанные после завершения работы, вызовут RuntimeError.

Если wait True , этот метод не вернётся, пока все ожидающие задачи не завершат выполнение, и ресурсы, связанные с исполнителем, не будут освобождены. Если wait False , этот метод вернётся немедленно, а ресурсы, связанные с исполнителем, будут освобождены, когда все ожидающие задачи завершат выполнение. Независимо от значения wait, вся программа Python не завершится, пока все ожидающие задачи не завершат выполнение.

Вы можете избежать необходимости явного вызова этого метода, если вы используете оператор with, который завершит работу Executor (ожидая, как если бы Executor.shutdown() был вызван со значением wait, установленным в True):

import shutil
with ThreadPoolExecutor(max_workers=4) as e:
    e.submit(shutil.copy, 'src1.txt', 'dest1.txt')
    e.submit(shutil.copy, 'src2.txt', 'dest2.txt')
    e.submit(shutil.copy, 'src3.txt', 'dest3.txt')
    e.submit(shutil.copy, 'src4.txt', 'dest4.txt')

ThreadPoolExecutor

ThreadPoolExecutor — это подкласс Executor, который использует пул потоков для асинхронного выполнения вызовов.

Тупики могут возникнуть, когда вызываемый объект, связанный с Future, ждёт результатов другого Future. Например:

import time
def wait_on_b():
    time.sleep(5)
    print(b.result())  # b will never complete because it is waiting on a.
    return 5

def wait_on_a():
    time.sleep(5)
    print(a.result())  # a will never complete because it is waiting on b.
    return 6


executor = ThreadPoolExecutor(max_workers=2)
a = executor.submit(wait_on_b)
b = executor.submit(wait_on_a)

И:

def wait_on_future():
    f = executor.submit(pow, 5, 2)
    # This will never complete because there is only one worker thread and
    # it is executing this function.
    print(f.result())

executor = ThreadPoolExecutor(max_workers=1)
executor.submit(wait_on_future)
class concurrent.futures.ThreadPoolExecutor(max_workers=None, thread_name_prefix='', initializer=None, initargs=())

Подкласс Executor, использующий пул из не более чем max_workers потоков для асинхронного выполнения вызовов.

initializer — необязательный вызываемый объект, который вызывается в начале каждого рабочего потока; initargs — кортеж аргументов, передаваемых инициализатору. Если initializer сгенерирует исключение, все текущие ожидающие задания сгенерируют BrokenThreadPool, а также любая попытка отправить больше заданий в пул.

Изменено в версии 3.5: Если max_workers None или не задано, оно по умолчанию равно количеству процессоров на машине, умноженному на 5, предполагая, что ThreadPoolExecutor часто используется для перекрытия ввода-вывода вместо работы с процессором, и количество рабочих потоков должно быть больше, чем количество рабочих потоков для ProcessPoolExecutor.

Новое в версии 3.6: Аргумент thread_name_prefix был добавлен для управления именами threading.Thread для потоков пула для более удобной отладки.

Изменено в версии 3.7: Добавлены аргументы initializer и initargs.

Пример ThreadPoolExecutor

import concurrent.futures
import urllib.request

URLS = ['http://www.foxnews.com/',
        'http://www.cnn.com/',
        'http://europe.wsj.com/',
        'http://www.bbc.co.uk/',
        'http://some-made-up-domain.com/']

# Retrieve a single page and report the URL and contents
def load_url(url, timeout):
    with urllib.request.urlopen(url, timeout=timeout) as conn:
        return conn.read()

# We can use a with statement to ensure threads are cleaned up promptly
with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor:
    # Start the load operations and mark each future with its URL
    future_to_url = {executor.submit(load_url, url, 60): url for url in URLS}
    for future in concurrent.futures.as_completed(future_to_url):
        url = future_to_url[future]
        try:
            data = future.result()
        except Exception as exc:
            print('%r generated an exception: %s' % (url, exc))
        else:
            print('%r page is %d bytes' % (url, len(data)))

ProcessPoolExecutor

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

Модуль __main__ должен быть импортируемым рабочими подпроцессами. Это означает, что ProcessPoolExecutor не будет работать в интерактивном интерпретаторе.

Вызов методов Executor или Future из вызываемого объекта, отправленного в ProcessPoolExecutor, приведёт к тупику.

class concurrent.futures.ProcessPoolExecutor(max_workers=None, mp_context=None, initializer=None, initargs=())

Подкласс Executor, который выполняет вызовы асинхронно с помощью пула из не более чем max_workers процессов. Если max_workers равно None или не указано, оно будет по умолчанию равно количеству процессоров на компьютере. Если max_workers меньше или равно 0, будет вызвано исключение ValueError. В Windows max_workers должно быть равно или меньше 61. В противном случае будет вызвано исключение ValueError. Если max_workers равно None, по умолчанию будет выбрано значение не более 61, даже если доступно больше процессоров. mp_context может быть контекстом multiprocessing или None. Он будет использоваться для запуска рабочих процессов. Если mp_context равно None или не указано, используется контекст multiprocessing по умолчанию.

initializer — это необязаемая функция, которая вызывается в начале каждого рабочего процесса; initargs — это кортеж аргументов, передаваемых в функцию initializer. Если initializer вызывает исключение, все текущие ожидающие задачи вызовут исключение BrokenProcessPool, а также любая попытка отправить больше задач в пул.

Изменено в версии 3.3: Если один из рабочих процессов завершается неожиданно, теперь возникает ошибка BrokenProcessPool. Раньше поведение было неопределенным, а операции с исполнителем или его будущими задачами часто замирали или попадали в тупик.

Изменено в версии 3.7: Аргумент mp_context был добавлен для предоставления пользователям возможности управлять методом start_method для рабочих процессов, созданных пулом.

Добавлены аргументы initializer и initargs.

ProcessPoolExecutor Пример

import concurrent.futures
import math

PRIMES = [
    112272535095293,
    112582705942171,
    112272535095293,
    115280095190773,
    115797848077099,
    1099726899285419]

def is_prime(n):
    if n % 2 == 0:
        return False

    sqrt_n = int(math.floor(math.sqrt(n)))
    for i in range(3, sqrt_n + 1, 2):
        if n % i == 0:
            return False
    return True

def main():
    with concurrent.futures.ProcessPoolExecutor() as executor:
        for number, prime in zip(PRIMES, executor.map(is_prime, PRIMES)):
            print('%d is prime: %s' % (number, prime))

if __name__ == '__main__':
    main()

Объекты Future

Класс Future инкапсулирует асинхронное выполнение вызываемого объекта. Экземпляры Future создаются методом Executor.submit().

class concurrent.futures.Future

Инкапсулирует асинхронное выполнение вызываемого объекта. Экземпляры Future создаются методом Executor.submit() и не должны создаваться напрямую, за исключением тестирования.

cancel()

Попытка отменить вызов. Если вызов выполняется или уже завершен и не может быть отменён, метод вернёт False, иначе вызов будет отменён и метод вернёт True.

cancelled()

Возвращает True если вызов был успешно отменён.

running()

Возвращает True если вызов выполняется в данный момент и не может быть отменён.

done()

Возвращает True если вызов был успешно отменён или завершён.

result(timeout=None)

Возвращает значение, возвращённое вызовом. Если вызов ещё не завершен, этот метод будет ждать до timeout секунд. Если вызов не завершится в течение timeout секунд, будет возбуждено исключение concurrent.futures.TimeoutError. timeout может быть целым или дробным числом. Если timeout не указан или None, ограничений по времени ожидания нет.

Если задача отменена до завершения, будет возбуждено исключение CancelledError.

Если при вызове возникло исключение, этот метод возбудит то же исключение.

exception(timeout=None)

Возвращает исключение, возбужденное вызовом. Если вызов ещё не завершен, этот метод будет ждать до timeout секунд. Если вызов не завершится в течение timeout секунд, будет возбуждено исключение concurrent.futures.TimeoutError. timeout может быть целым или дробным числом. Если timeout не указан или None, ограничений по времени ожидания нет.

Если задача отменена до завершения, будет возбуждено исключение CancelledError.

Если вызов завершился без возбуждения исключения, возвращается None.

add_done_callback(fn)

Присоединяет вызываемый объект fn к задаче. fn будет вызван с задачей в качестве единственного аргумента, когда задача отменена или завершена.

Вызываемые объекты вызываются в том порядке, в котором они были добавлены, и всегда вызываются в потоке, принадлежащем процессу, который их добавил. Если вызываемый объект возбуждает исключение, являющееся подклассом Exception, оно будет зарегистрировано и проигнорировано. Если вызываемый объект возбуждает исключение, являющееся подклассом BaseException, поведение не определено.

Если задача уже завершена или отменена, fn будет вызван немедленно.

Следующие методы Future предназначены для использования в тестах на единицу и реализациях Executor.

set_running_or_notify_cancel()

Этот метод должен вызываться только реализациями Executor перед выполнением работы, связанной с Future, и тестами на единицу.

Если метод возвращает False, то Future была отменена, т.е. был вызван Future.cancel() и возвращено значение True. Любые потоки, ожидающие завершения Future (например, через as_completed() или wait()), будут разбужены.

Если метод возвращает True, то Future не была отменена и переведена в состояние выполнения, т.е. вызовы Future.running() вернут True.

Этот метод может быть вызван только один раз и не может быть вызван после того, как были вызваны Future.set_result() или Future.set_exception().

set_result(result)

Устанавливает результат работы, связанной с Future, в result.

Этот метод должен использоваться только реализациями Executor и тестами на единицу.

set_exception(exception)

Устанавливает результат работы, связанной с Future, в исключение Exception exception.

Этот метод должен использоваться только реализациями Executor и тестами на единицу.

Функции модуля

concurrent.futures.wait(fs, timeout=None, return_when=ALL_COMPLETED)

Ожидание завершения экземпляров Future (возможно, созданных различными экземплярами Executor) в списке fs. Возвращает именованную пару из двух множеств. Первое множество, названное done, содержит фьючеры, которые завершились (завершенные или отменённые фьючеры) до завершения ожидания. Второе множество, названное not_done, содержит фьючеры, которые не завершились (ожидающие или работающие фьючеры).

timeout может использоваться для управления максимальным временем ожидания в секундах, прежде чем функция вернётся. timeout может быть целым или дробным числом. Если timeout не указан или None, то время ожидания не ограничено.

return_when указывает, когда функция должна вернуть значение. Оно должно быть одним из следующих констант:

Константа

Описание

FIRST_COMPLETED

Функция вернётся, когда любой фьючер завершится или будет отменён.

FIRST_EXCEPTION

Функция вернётся, когда любой фьючер завершится, вызвав исключение. Если ни один фьючер не вызывает исключения, то это эквивалентно ALL_COMPLETED.

ALL_COMPLETED

Функция вернётся, когда все фьючеры завершатся или будут отменены.

concurrent.futures.as_completed(fs, timeout=None)

Возвращает итератор по экземплярам Future (возможно, созданным различными экземплярами Executor в списке fs), который возвращает фьючеры по мере их завершения (завершенные или отменённые фьючеры). Любые дублирующиеся фьючеры в fs будут возвращены один раз. Любые фьючеры, которые завершились до вызова as_completed(), будут возвращены в первую очередь. Возвращённый итератор вызывает concurrent.futures.TimeoutError, если к нему обращается метод __next__(), а результат недоступен через timeout секунд с момента первоначального вызова as_completed(). timeout может быть целым или дробным числом. Если timeout не указан или None, то время ожидания не ограничено.

См. также

PEP 3148 – futures - асинхронное выполнение вычислений

Предложение, в котором описывается эта функция для включения в стандартную библиотеку Python.

Классы исключений

exception concurrent.futures.CancelledError

Вызывается, когда фьючер отменяется.

exception concurrent.futures.TimeoutError

Вызывается, когда операция фьючера превышает заданное время ожидания.

exception concurrent.futures.BrokenExecutor

Производный от RuntimeError класс исключений. Вызывается, когда исполняющий процесс по каким-либо причинам неисправен и не может быть использован для отправки или выполнения новых задач.

Добавлена в версии 3.7.

exception concurrent.futures.thread.BrokenThreadPool

Производный от BrokenExecutor класс исключений. Вызывается, когда один из рабочих процессов пула потоков ThreadPoolExecutor не смог завершить инициализацию.

Добавлена в версии 3.7.

exception concurrent.futures.process.BrokenProcessPool

Производный от BrokenExecutor (ранее RuntimeError) класс исключений. Вызывается, когда один из рабочих процессов пула процессов ProcessPoolExecutor завершился некорректно (например, был прерван извне).

Добавлена в версии 3.3.

© 2001–2020 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.7/library/concurrent.futures.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API