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 в положительное целое число. Для очень длинных итерируемых объектов использование большого значения для 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 вызовет исключение, все текущие ожидающие задания вызовут
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=()) -
Подкласс
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:
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 секунд, будет поднято исключение
concurrent.futures.TimeoutError. timeout может быть целым или дробным числом. Если timeout не указан илиNone, ограничений по времени ожидания нет.Если будущее (future) отменено до завершения, будет поднято исключение
CancelledError.Если вызов поднял исключение, этот метод поднимет то же самое исключение.
-
exception(timeout=None) -
Возвращает исключение, поднятое вызовом. Если вызов ещё не завершён, метод будет ждать до timeout секунд. Если вызов не завершится в течение timeout секунд, будет поднято исключение
concurrent.futures.TimeoutError. timeout может быть целым или дробным числом. Если timeout не указан илиNone, ограничений по времени ожидания нет.Если будущее (future) отменено до завершения, будет поднято исключение
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 указывает, когда эта функция должна вернуть результат. Он должен быть одним из следующих констант:
Константа
Описание
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.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.10/library/concurrent.futures.html