Spec-Zone.ru › Python 3.11

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

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

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

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

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

Доступность: нет в Emscripten, нет в WASI.

Этот модуль не работает или недоступен на платформах WebAssembly wasm32-emscripten и wasm32-wasi. См. Платформы WebAssembly для получения дополнительной информации.

Объекты 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 могут выполняться одновременно.

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

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

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

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

shutdown(wait=True, *, cancel_futures=False)

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

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

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

Если оба cancel_futures и wait True, все задачи, которые исполнитель начал выполнять, завершатся до того, как этот метод вернётся. Остальные задачи отменяются.

Вы можете избежать явного вызова этого метода, если используете оператор 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')

Изменено в версии 3.9: Добавлен cancel_futures.

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 потоков для асинхронного выполнения вызовов.

Все потоки, запущенные в ThreadPoolExecutor , будут объединены до того, как интерпретатор сможет выйти. Обратите внимание, что обработчик выхода, который делает это, выполняется *до* любых обработчиков выхода, добавленных с помощью atexit. Это означает, что исключения в основном потоке необходимо перехватывать и обрабатывать, чтобы сигнализировать потокам о прекращении работы корректно. По этой причине рекомендуется использовать ThreadPoolExecutor не для долгосрочных задач.

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

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

Новое в версии 3.6: Добавлен аргумент thread_name_prefix, позволяющий пользователям контролировать имена threading.Thread для рабочих потоков, создаваемых пулом, для более удобной отладки.

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

Изменено в версии 3.8: Значение по умолчанию для max_workers изменено на min(32, os.cpu_count() + 4). Это значение по умолчанию сохраняет как минимум 5 рабочих процессов для задач ввода-вывода. Оно использует не более 32 ядер процессора для задач, связанных с ЦП, которые освобождают GIL. И оно избегает неявного использования очень больших ресурсов на машинах с множеством ядер.

ThreadPoolExecutor теперь повторно использует свободные рабочие потоки, прежде чем запускать max_workers рабочих потоков.

Пример 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://nonexistant-subdomain.python.org/']

# 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=(), max_tasks_per_child=None)

Подкласс 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, а также любая попытка отправить новые задачи в пул.

max_tasks_per_child — необязательный аргумент, определяющий максимальное количество задач, которые может выполнить один процесс, прежде чем он завершит работу и будет заменён новым рабочим процессом. По умолчанию max_tasks_per_child равно None, что означает, что рабочие процессы будут существовать до тех пор, пока существует пул. При указании максимального значения метод запуска multiprocessing «spawn» будет использоваться по умолчанию при отсутствии параметра mp_context. Эта функция несовместима с методом запуска «fork».

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

Изменено в версии 3.7: Добавлен аргумент mp_context, позволяющий пользователям управлять методом запуска для рабочих процессов, созданных пулом.

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

Изменено в версии 3.11: Добавлен аргумент max_tasks_per_child, позволяющий пользователям управлять временем жизни рабочих процессов в пуле.

Пример ProcessPoolExecutor

import concurrent.futures
import math

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

def is_prime(n):
    if n < 2:
        return False
    if n == 2:
        return True
    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 секунд, будет возбуждено исключение TimeoutError. timeout может быть целым или дробным числом. Если timeout не указан или None, ограничений по времени ожидания нет.

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

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

exception(timeout=None)

Возвращает исключение, возбуждённое вызовом. Если вызов ещё не завершён, этот метод будет ожидать до timeout секунд. Если вызов не завершится в течение timeout секунд, будет возбуждено исключение 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 и тестами на единицу.

Изменено в версии 3.8: Этот метод возбуждает concurrent.futures.InvalidStateError, если Future уже завершена.

set_exception(exception)

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

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

Изменено в версии 3.8: Этот метод возбуждает concurrent.futures.InvalidStateError, если Future уже завершена.

END_OF_DOCUMENT_MARKER

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

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

Ожидание завершения экземпляров Future (возможно, созданных различными экземплярами Executor), переданных в fs. Дублируемые значения в 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(), будут возвращены в первую очередь. Возвращаемый итератор генерирует исключение TimeoutError, если к нему обращается метод __next__(), и результат недоступен через timeout секунд с момента первоначального вызова as_completed(). timeout может быть целым или вещественным числом. Если timeout не указан или None, времени ожидания нет предела.

См. также

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

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

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

exception concurrent.futures.CancelledError

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

exception concurrent.futures.TimeoutError

Устаревшее псевдоним TimeoutError, вызывается, когда операция задачи превышает заданный интервал ожидания.

Изменено в версии 3.11: Этот класс был преобразован в псевдоним TimeoutError.

exception concurrent.futures.BrokenExecutor

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

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

exception concurrent.futures.InvalidStateError

Вызывается, когда операция выполняется над задачей, которая недопустима в текущем состоянии.

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

exception concurrent.futures.thread.BrokenThreadPool

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

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

exception concurrent.futures.process.BrokenProcessPool

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

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

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

Spec-Zone.ru

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