tf.train.Coordinator
| Просмотреть исходный код на GitHub |
Координатор для потоков.
tf.train.Coordinator(
clean_stop_exception_types=None
)
Этот класс реализует простой механизм для координации завершения набора потоков.
Использование:
# Create a coordinator. coord = Coordinator() # Start a number of threads, passing the coordinator to each of them. ...start thread 1...(coord, ...) ...start thread N...(coord, ...) # Wait for all the threads to terminate. coord.join(threads)
Любой из потоков может вызвать coord.request_stop(), чтобы запросить остановку всех потоков. Для сотрудничества с запросами каждый поток должен регулярно проверять состояние coord.should_stop(). coord.should_stop() возвращает True, как только coord.request_stop() был вызван.
Типичный поток, работающий с координатором, будет выполнять что-то вроде:
while not coord.should_stop(): ...do some work...
Обработка исключений:
Поток может сообщить о возникшем исключении координатору в рамках вызова request_stop(). Исключение будет повторно поднято при вызове coord.join().
Код потока:
try:
while not coord.should_stop():
...do some work...
except Exception as e:
coord.request_stop(e)
Основной код:
try: ... coord = Coordinator() # Start a number of threads, passing the coordinator to each of them. ...start thread 1...(coord, ...) ...start thread N...(coord, ...) # Wait for all the threads to terminate. coord.join(threads) except Exception as e: ...exception that was passed to coord.request_stop()
Для упрощения реализации потоков координатор предоставляет обработчик контекста stop_on_exception(), который автоматически запрашивает остановку, если возникает исключение. Используя обработчик контекста, код потока выше можно переписать следующим образом:
with coord.stop_on_exception():
while not coord.should_stop():
...do some work...
Период ожидания остановки:
После того, как поток вызвал coord.request_stop(), у других потоков есть ограниченное время для остановки, это называется 'период ожидания остановки' и по умолчанию составляет 2 минуты. Если какой-либо из потоков остается активным после истечения срока ожидания, coord.join() вызывает RuntimeError, сообщая об отстающих потоках.
try: ... coord = Coordinator() # Start a number of threads, passing the coordinator to each of them. ...start thread 1...(coord, ...) ...start thread N...(coord, ...) # Wait for all the threads to terminate, give them 10s grace period coord.join(threads, stop_grace_period_secs=10) except RuntimeError: ...one of the threads took more than 10s to stop after request_stop() ...was called. except Exception: ...exception that was passed to coord.request_stop()
| Аргументы | |
|---|---|
clean_stop_exception_types | Необязательная кортеж типов исключений, которые должны вызывать чистую остановку координатора. Если исключение одного из этих типов сообщается координатору в request_stop(ex), координатор будет вести себя так, как если бы был вызван request_stop(None). По умолчанию используется (tf.errors.OutOfRangeError,), который используется очередями ввода для сигнализации об окончании ввода. При подаче обучающих данных из итератора Python часто добавляют StopIteration в этот список. |
| Атрибуты | |
|---|---|
joined | |
Методы
clear_stop
clear_stop()
Очищает флаг остановки.
После этого вызовы should_stop() будут возвращать False.
join
join(
threads=None, stop_grace_period_secs=120, ignore_live_threads=False
)
Ожидание завершения потоков.
Этот вызов блокируется, пока набор потоков не завершится. Набор потоков — это объединение потоков, переданных в аргументе threads, и список потоков, зарегистрированных в координаторе с помощью вызова Coordinator.register_thread().
После остановки потоков, если исключение exc_info было передано в request_stop, то это исключение будет повторно поднято.
Обработка периода ожидания: когда request_stop() вызывается, потокам даётся 'stop_grace_period_secs' секунд для завершения. Если какой-либо из них остается активным после истечения этого периода, возникает исключение RuntimeError. Обратите внимание, что если исключение exc_info было передано в request_stop(), то оно будет поднято вместо RuntimeError.
| Аргументы | |
|---|---|
threads | Список threading.Threads. Запущенные потоки для присоединения, помимо зарегистрированных потоков. |
stop_grace_period_secs | Количество секунд, предоставляемое потокам для остановки после вызова request_stop(). |
ignore_live_threads | Если False, вызовет ошибку, если какие-либо из потоков всё ещё активны после stop_grace_period_secs. |
| Возможные исключения | |
|---|---|
RuntimeError | Если какой-либо поток всё ещё активен после вызова request_stop() и истечения периода ожидания. |
raise_requested_exception
raise_requested_exception()
Если исключение было передано в request_stop, оно будет поднято.
register_thread
register_thread(
thread
)
Регистрация потока для присоединения.
| Аргументы | |
|---|---|
thread | Поток Python для присоединения. |
request_stop
request_stop(
ex=None
)
Запрос остановки потоков.
После этого вызовы should_stop() будут возвращать True.
Примечание: Если исключение передаётся, оно должно передаваться в контексте обработки исключения (т.е. try: ... except Exception as ex: ...), а не создаваться заново.
| Аргументы | |
|---|---|
ex | Необязательное исключение, или кортеж Python-исключений, как возвращается из sys.exc_info(). Если это первый вызов request_stop(), соответствующее исключение будет записано и повторно поднято из join(). |
should_stop
should_stop()
Проверка запроса на остановку.
| Возвращаемое значение | |
|---|---|
| True, если запрос на остановку получен. |
stop_on_exception
@contextlib.contextmanager stop_on_exception()
Блок кода для запроса остановки при возникновении исключения.
Код, использующий координатор, должен перехватывать исключения и передавать их в метод request_stop(), чтобы остановить другие потоки, управляемые координатором.
Этот обработчик контекста упрощает обработку исключений. Используйте его следующим образом:
with coord.stop_on_exception(): # Any exception raised in the body of the with # clause is reported to the coordinator before terminating # the execution of the body. ...body...
Это полностью эквивалентно несколько более длинному коду:
try: ...body... except: coord.request_stop(sys.exc_info())
Выходные данные:
ничего.
wait_for_stop
wait_for_stop(
timeout=None
)
Ожидание, пока координатор не получит команду на остановку.
| Аргументы | |
|---|---|
timeout | Число с плавающей точкой. Длительность ожидания в секундах, пока `should_stop()` не станет True. |
| Возвращаемое значение | |
|---|---|
| True, если координатор получил команду на остановку, False, если истекло время ожидания. |
© 2020 The TensorFlow Authors. All rights reserved.
Licensed under the Creative Commons Attribution License 3.0.
Code samples licensed under the Apache 2.0 License.
https://www.tensorflow.org/versions/r2.3/api_docs/python/tf/train/Coordinator