concurrent.futures — Запуск параллельных задач
Добавлен в версии 3.2.
Исходный код: Lib/concurrent/futures/thread.py и Lib/concurrent/futures/process.py
Модуль concurrent.futures предоставляет высокоуровневый интерфейс для асинхронного выполнения вызываемых объектов.
Асинхронное выполнение можно осуществлять с помощью потоков, используя ThreadPoolExecutor, или отдельных процессов, используя ProcessPoolExecutor. Оба реализуют один и тот же интерфейс, который определен абстрактным классом Executor.
Доступность: не WASI.
Этот модуль не работает и недоступен в WebAssembly. См. 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), но:- списки собираются немедленно, а не лениво;
- fn выполняется асинхронно, и несколько вызовов fn могут быть выполнены одновременно.
Возвращаемый итератор поднимает исключение
TimeoutError, если__next__()вызывается, и результат недоступен через timeout секунд с момента первоначального вызоваExecutor.map(). timeout может быть целым или дробным числом. Если timeout не указан илиNone, времени ожидания нет.Если вызов fn вызывает исключение, то это исключение будет поднято при получении его значения из итератора.
При использовании
ProcessPoolExecutorэтот метод делит списки на несколько частей, которые отправляются в пул как отдельные задачи. Размер (приблизительный) этих частей можно указать, установив chunksize в положительное целое число. Для очень длинных списков использование большого значения для chunksize может значительно улучшить производительность по сравнению с величиной по умолчанию, равной 1. СThreadPoolExecutor, chunksize не имеет эффекта.Изменено в версии 3.5: Добавлен аргумент chunksize.
-
shutdown(wait=True, *, cancel_futures=False) -
Сигнализирует исполнителю о том, что он должен освободить любые используемые ресурсы, когда текущие ожидающие задачи завершат выполнение. Вызовы
Executor.submit()иExecutor.map(), сделанные после вызова shutdown, приведут к исключению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 вызывает исключение, все текущие ожидающие задания приведут к исключению
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 потоков работников.
Изменено в версии 3.13: Значение по умолчанию max_workers изменено на
min(32, (os.process_cpu_count() or 1) + 4).
Пример 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://nonexistent-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или не задано, оно будет установлено по умолчанию в значениеos.process_cpu_count(). Если 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()для получения дополнительной информации.Изменено в версии 3.13: По умолчанию для max_workers используется
os.process_cpu_count(), а неos.cpu_count().
Пример 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. Дубликаты задач в fs удаляются и будут возвращены только один раз. Возвращает именованную пару кортежей из множеств. Первое множество, названноеdone, содержит задачи, которые завершились (завершенные или отменённые задачи) до завершения ожидания. Второе множество, названноеnot_done, содержит задачи, которые не завершились (ожидающие или выполняющиеся задачи).timeout может использоваться для управления максимальным временем ожидания в секундах, прежде чем функция вернётся. timeout может быть целым или дробным числом. Если timeout не указан или
None, время ожидания не ограничено.return_when указывает, когда функция должна вернуть значение. Это должно быть одно из следующих констант:
Константа
Описание
-
concurrent.futures.FIRST_COMPLETED
Функция вернётся, когда любая задача завершится или будет отменена.
-
concurrent.futures.FIRST_EXCEPTION
Функция вернётся, когда любая задача завершится с исключением. Если ни одна задача не сгенерирует исключение, то это эквивалентно
ALL_COMPLETED.-
concurrent.futures.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–2024 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.13/library/concurrent.futures.html