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 на положительное целое число. Для очень длинных iterables использование большого значения для chunksize может значительно улучшить производительность по сравнению со значением по умолчанию 1. СThreadPoolExecutor, chunksize не имеет эффекта.Изменено в версии 3.5: Добавлен аргумент chunksize.
-
shutdown(wait=True) -
Уведомляет исполнителя о том, что он должен освободить все используемые ресурсы, когда текущие ожидающие задачи будут завершены. Вызовы
Executor.submit()иExecutor.map(), сделанные после вызова shutdown, вызовутRuntimeError.Если wait
True, этот метод не вернётся, пока все ожидающие задачи не будут завершены, и ресурсы, связанные с исполнителем, не будут освобождены. Если waitFalse, этот метод вернётся немедленно, а ресурсы, связанные с исполнителем, будут освобождены, когда все ожидающие задачи будут завершены. Независимо от значения wait, вся программа Python не завершится, пока все ожидающие задачи не будут завершены.Вы можете избежать явного вызова этого метода, если используете оператор
with, который выполнит shutdown для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')
-
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. Если 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. Возвращает именованную пару кортежей из множеств. Первое множество, названное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класс исключений, возбуждается, когда выполняет задачи Executor по какой-то причине неисправен и не может быть использован для отправки или выполнения новых задач.Новое в версии 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–2022 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.8/library/concurrent.futures.html