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(fn, *iterables, timeout=None, chunksize=1) -
Аналогично
map(fn, *iterables), за исключением:- iterables собираются немедленно, а не лениво;
- fn выполняется асинхронно и могут быть сделаны несколько вызовов fn одновременно.
Возвращаемый итератор вызывает
TimeoutError, если__next__()вызывается и результат недоступен через timeout секунд после первоначального вызоваExecutor.map(). timeout может быть целым числом или числом с плавающей запятой. Если timeout не указан илиNone, нет ограничения на время ожидания.Если вызов fn вызывает исключение, то это исключение будет вызвано при получении его значения из итератора.
При использовании
ProcessPoolExecutor, этот метод делит iterables на несколько кусков, которые он отправляет в пул в качестве отдельных задач. Приблизительный размер этих кусков можно указать, задав chunksize положительное целое число. Для очень длинных итерабельных объектов использование большого значения chunksize может значительно улучшить производительность по сравнению со стандартным значением 1. СThreadPoolExecutor, chunksize не оказывает никакого эффекта.Изменено в версии 3.5: Добавлен аргумент chunksize.
-
shutdown(wait=True, *, cancel_futures=False) -
Сигнализирует исполнителю, что он должен освободить все используемые ресурсы, когда текущие ожидающие задачи будут завершены. Вызовы
Executor.submit()иExecutor.map(), сделанные после завершения, вызовутRuntimeError.Если wait
True, этот метод не возвращается, пока все ожидающие задачи не будут завершены, и ресурсы, связанные с исполнителем, не будут освобождены. Если waitFalse, этот метод возвращается немедленно, и ресурсы, связанные с исполнителем, будут освобождены, когда все ожидающие задачи завершатся. Независимо от значения 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. Если 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.
Примечание
Метод запуска
multiprocessingпо умолчанию (см. Контексты и методы запуска) изменится с fork на Python 3.14. Код, которому необходимо использовать fork для своихProcessPoolExecutor, должен явно указать это, передав параметрmp_context=multiprocessing.get_context("fork").Изменено в версии 3.11: Добавлен аргумент max_tasks_per_child, позволяющий пользователям управлять временем жизни рабочих процессов в пуле.
Изменено в версии 3.12: В системах POSIX, если ваше приложение имеет несколько потоков и контекст
multiprocessingиспользует метод запуска"fork", функцияos.fork(), вызываемая внутри для запуска рабочих процессов, может генерироватьDeprecationWarning. Передайте mp_context, настроенный на использование другого метода запуска. Смотрите документацию кos.fork()для получения дополнительной информации.
Пример 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, в исключениеExceptionexception.Этот метод должен использоваться только реализациями
Executorи модульными тестами.Изменено в версии 3.8: Этот метод вызывает исключение
concurrent.futures.InvalidStateError, еслиFutureуже завершена.
-
Функции модуля
-
concurrent.futures.wait(fs, timeout=None, return_when=ALL_COMPLETED) -
Ожидание завершения экземпляров
Future(возможно, созданных различными экземплярамиExecutor) заданных в fs. Дублированные future в fs удаляются и будут возвращены только один раз. Возвращает именованную пару из двух множеств. Первое множество, названноеdone, содержит future, которые завершились (завершенные или отменённые future) до завершения ожидания. Второе множество, названноеnot_done, содержит future, которые не завершились (ожидающие или выполняющиеся future).timeout может использоваться для управления максимальным временем ожидания в секундах перед возвратом. timeout может быть целым или вещественным числом. Если timeout не указан или
None, время ожидания не ограничено.return_when указывает, когда эта функция должна возвратить значение. Должно быть одно из следующих констант:
Константа
Описание
-
concurrent.futures.FIRST_COMPLETED
Функция вернётся, когда любое future завершится или будет отменено.
-
concurrent.futures.FIRST_EXCEPTION
Функция вернётся, когда любое future завершится сгенерированием исключения. Если ни одно future не генерирует исключение, то это эквивалентно
ALL_COMPLETED.-
concurrent.futures.ALL_COMPLETED
Функция вернётся, когда все future завершатся или будут отменены.
-
-
concurrent.futures.as_completed(fs, timeout=None) -
Возвращает итератор по экземплярам
Future(возможно, созданных различными экземплярамиExecutor) заданных в fs, который возвращает future по мере их завершения (завершенные или отменённые future). Любые дублированные future в fs будут возвращены один раз. Любые future, завершившиеся до вызоваas_completed(), будут возвращены первыми. Возвращаемый итератор генерируетTimeoutError, если__next__()вызывается, а результат недоступен после timeout секунд с момента первоначального вызоваas_completed(). timeout может быть целым или вещественным числом. Если timeout не указан илиNone, время ожидания не ограничено.
См. также
- PEP 3148 – futures - выполнение вычислений асинхронно
-
Предложение, в котором описана эта возможность для включения в стандартную библиотеку Python.
Классы исключений
-
exception concurrent.futures.CancelledError -
Сгенерировано, когда future отменено.
-
exception concurrent.futures.TimeoutError -
Устаревшее алиас
TimeoutError, генерируется, когда операция future превышает заданное время ожидания.Изменено в версии 3.11: Этот класс стал алиасом
TimeoutError.
-
exception concurrent.futures.BrokenExecutor -
Производный от
RuntimeError, этот класс исключений генерируется, когда executor по какой-то причине неисправен и не может быть использован для отправки или выполнения новых задач.Добавлен в версии 3.7.
-
exception concurrent.futures.InvalidStateError -
Сгенерировано, когда операция выполняется на future, которая не разрешена в текущем состоянии.
Добавлен в версии 3.8.
-
exception concurrent.futures.thread.BrokenThreadPool -
Производный от
BrokenExecutor, этот класс исключений генерируется, когда один из работниковThreadPoolExecutorне смог инициализироваться.Добавлен в версии 3.7.
-
exception concurrent.futures.process.BrokenProcessPool -
Производный от
BrokenExecutor(ранееRuntimeError), этот класс исключений генерируется, когда один из работниковProcessPoolExecutorзавершился некорректно (например, если был убит извне).Добавлен в версии 3.3.
© 2001–2024 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.12/library/concurrent.futures.html