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), за исключением:- итерируемые объекты собираются немедленно, а не лениво;
- 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 потоков, для асинхронного выполнения вызовов.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://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 вызывает исключение, все текущие ожидающие задачи вызовут исключение
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, время ожидания не ограничено.Если будущее было отменено до завершения, будет поднято исключение
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и модульными тестами.Изменено в версии 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 указывает, когда эта функция должна возвратить результат. Он должен быть одним из следующих констант:
Константа
Описание
FIRST_COMPLETEDФункция вернётся, когда завершится или отменится любая future.
FIRST_EXCEPTIONФункция вернётся, когда любая future завершится с исключением. Если ни одна future не вызовет исключение, это эквивалентно
ALL_COMPLETED.ALL_COMPLETEDФункция вернётся, когда завершатся или отменятся все future.
-
concurrent.futures.as_completed(fs, timeout=None) -
Возвращает итератор по экземплярам
Future(возможно, созданных различными экземплярамиExecutor) в списке fs, который возвращает future по мере их завершения (завершённые или отменённые future). Любые дублированные future в fs будут возвращены один раз. Любые future, которые завершились до вызоваas_completed(), будут возвращены первыми. Возвращённый итератор поднимает исключениеconcurrent.futures.TimeoutError, если к нему обращаются с помощью метода__next__()и результат не доступен через timeout секунд от первоначального вызоваas_completed(). timeout может быть целым или вещественным числом. Если timeout не указан илиNone, время ожидания не ограничено.
См. также
- PEP 3148 – futures - асинхронное выполнение вычислений
-
Предложение, описывающее эту функцию для включения в стандартную библиотеку Python.
Классы исключений
-
exception concurrent.futures.CancelledError -
Вызывается, когда future отменена.
-
exception concurrent.futures.TimeoutError -
Вызывается, когда операция future превышает заданное время ожидания.
-
exception concurrent.futures.BrokenExecutor -
Производный от
RuntimeError, этот класс исключений генерируется, когда исполняющий механизм по какой-то причине сломан и не может использоваться для отправки или выполнения новых задач.Введено в версии 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–2022 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.9/library/concurrent.futures.html