Таймеры истечения срока действия
Создано: 04 мая 2021 г. | Последнее обновление: 06 мая 2026 г.
Таймеры истечения срока действия настраиваются в том же процессе, что и агент, и используются в вашем скрипте для обработки зависших рабочих процессов. Когда вы переходите к блоку кода, который может зависнуть, можно активировать таймер истечения срока действия. Он указывает серверу таймеров завершить процесс, если тот не освободит таймер до установленного срока его действия.
Использование:
import torchelastic.timer as timer
import torchelastic.agent.server as agent
def main():
start_method = "spawn"
message_queue = mp.get_context(start_method).Queue()
server = timer.LocalTimerServer(message, max_interval=0.01)
server.start() # non-blocking
spec = WorkerSpec(
fn=trainer_func,
args=(message_queue,),
...<OTHER_PARAMS...>)
agent = agent.LocalElasticAgent(spec, start_method)
agent.run()
def trainer_func(message_queue):
timer.configure(timer.LocalTimerClient(message_queue))
with timer.expires(after=60): # 60 second expiry
# do some work
В приведённом выше примере, если trainer_func выполняется более 60 секунд, рабочий процесс завершается, и агент повторно запускает группу рабочих процессов.
Методы клиента
-
torch.distributed.elastic.timer.configure(timer_client)[исходный код] -
Настраивает клиент таймера. Необходимо вызвать перед использованием
expires.
-
torch.distributed.elastic.timer.expires(after, scope=None, client=None)[исходный код] -
Активирует таймер обратного отсчёта, срок действия которого истекает через
afterсекунд, если заключённый в него блок кода не завершится за это время. После истечения срока действия таймера этот рабочий процесс может быть принудительно завершён. Точное значение термина «завершён» зависит от реализации клиента. В большинстве случаев завершение означает остановку процесса рабочего процесса. Обратите внимание: рабочий процесс НЕ гарантированно будет завершён точно в моментtime.now() + after. Он лишь становится «подлежащим» завершению, аTimerServer, с которым взаимодействует клиент, в конечном итоге принимает решение о том, когда и как завершать рабочие процессы с истёкшими таймерами.Использование:
torch.distributed.elastic.timer.configure(LocalTimerClient()) with expires(after=10): torch.distributed.all_reduce(...)
Реализации сервера и клиента
Ниже представлены пары сервера и клиента таймера, предоставляемые torchelastic.
Примечание
Сервер таймера и его клиенты всегда должны реализовываться и использоваться парами, поскольку между сервером и клиентом действует протокол обмена сообщениями.
Ниже представлена пара сервера и клиента таймера, реализованная на основе multiprocess.Queue.
-
class torch.distributed.elastic.timer.LocalTimerServer(mp_queue, max_interval=60, daemon=True)[исходный код] -
Сервер, работающий с
LocalTimerClient. Предполагается, что клиенты являются дочерними процессами родительского процесса, в котором работает этот сервер. Ожидается, что каждый узел задания запустит собственный локальный сервер таймеров. Каждый экземпляр сервера управляет таймерами локальных рабочих процессов (запущенных в процессах на том же узле).
-
class torch.distributed.elastic.timer.LocalTimerClient(mp_queue)[исходный код] -
Клиентская часть
LocalTimerServer. Этот клиент предназначен для использования на том же узле, где работаетLocalTimerServer, и использует pid для уникальной идентификации рабочего процесса. Это особенно полезно в ситуациях, когда на узле с несколькими GPU-устройствами запускается отдельный дочерний процесс (тренировщик) для каждого GPU.
Ниже представлена ещё одна пара сервера и клиента таймера, реализованная на основе именованного канала.
-
class torch.distributed.elastic.timer.FileTimerServer(file_path, run_id, max_interval=10, daemon=True, log_event=None)[исходный код] -
Сервер, работающий с
FileTimerClient. Предполагается, что клиенты работают на том же узле, что и процесс, в котором работает этот сервер. Ожидается, что каждый узел задания запустит собственный локальный сервер таймеров. Каждый экземпляр сервера управляет таймерами локальных рабочих процессов (запущенных в процессах на том же узле).- Параметры:
-
- file_path (str) – str, путь к специальному FIFO-файлу, который будет создан.
- max_interval (float) – float, максимальный интервал в секундах для каждого цикла наблюдения.
- daemon (bool) – bool, указывает, следует ли запускать поток наблюдения в режиме демона. Поток-демон не будет препятствовать остановке процесса.
- log_event (Callable[[str, FileTimerRequest | None], None] | None) – Callable[[Dict[str, str]], None], необязательная функция обратного вызова для записи событий в формате JSON.
-
class torch.distributed.elastic.timer.FileTimerClient(file_path, signal=Signals.SIGKILL)[исходный код] -
Клиентская часть
FileTimerServer. Этот клиент предназначен для использования на том же узле, где работаетFileTimerServer, и использует pid для уникальной идентификации рабочего процесса. Клиент использует named_pipe для отправки запросов таймера вFileTimerServer. Этот клиент является производителем, аFileTimerServer— потребителем. Несколько клиентов могут работать с одним и тем жеFileTimerServer.- Параметры:
-
-
file_path (str) – str, путь к специальному FIFO-файлу.
FileTimerServerдолжен создать его вызовом os.mkfifo(). - signal – сигнал, который следует использовать для завершения процесса. Отрицательное значение сигнала или ноль не приведёт к завершению процесса.
-
file_path (str) – str, путь к специальному FIFO-файлу.
Создание собственного сервера и клиента таймера
Чтобы создать собственные сервер и клиент таймера, унаследуйте сервер от torch.distributed.elastic.timer.TimerServer, а клиент — от torch.distributed.elastic.timer.TimerClient. Объект TimerRequest используется для передачи сообщений между сервером и клиентом.
-
class torch.distributed.elastic.timer.TimerRequest(worker_id, scope_id, expiration_time)[исходный код] -
Объект данных, представляющий активацию и освобождение таймера обратного отсчёта, который используется между
TimerClientиTimerServer. Отрицательное значениеexpiration_timeследует интерпретировать как запрос на освобождение.Примечание
Тип
worker_idзависит от реализации. Это значение, которое реализации TimerServer и TimerClient используют для уникальной идентификации рабочего процесса.
-
class torch.distributed.elastic.timer.TimerServer(request_queue, max_interval, daemon=True)[исходный код] -
Компонент, который отслеживает активные таймеры и своевременно завершает их действие. Этот сервер отвечает за завершение рабочих процессов с истёкшими таймерами.
-
abstract clear_timers(worker_ids)[исходный код] -
Удаляет все таймеры для указанного
worker_ids.
-
abstract get_expired_timers(deadline)[исходный код] -
Возвращает все истёкшие таймеры для каждого worker_id. Таймер считается истёкшим, если его expiration_time меньше или равно указанному сроку.
- Тип возвращаемого значения:
-
dict[str, list[TimerRequest]]
-
abstract register_timers(timer_requests)[исходный код] -
Обрабатывает входящие запросы таймера и регистрирует их на сервере. Запрос таймера может быть запросом на активацию или освобождение таймера. Запросы таймера с отрицательным значением expiration_time следует интерпретировать как запросы на освобождение таймера.
-
-
class torch.distributed.elastic.timer.TimerClient[исходный код] -
Клиентская библиотека для активации и освобождения таймеров обратного отсчёта посредством взаимодействия с TimerServer.
-
abstract acquire(scope_id, expiration_time)[исходный код] -
Активирует таймер для рабочего процесса, которому принадлежит этот объект клиента, с учётом scope_id и expiration_time. Обычно регистрирует таймер на TimerServer.
-
abstract release(scope_id)[исходный код] -
Освобождает таймер для
scope_idрабочего процесса, представленного этим клиентом. После вызова этого метода таймер обратного отсчёта для данной области действия перестаёт действовать.
-
Запись отладочной информации в журнал
-
torch.distributed.elastic.timer.debug_info_logging.log_debug_info_for_expired_timers(run_id, expired_timers)[исходный код]
© 2026, PyTorch Contributors
PyTorch has a BSD-style license, as found in the LICENSE file.
https://docs.pytorch.org/docs/2.14/elastic/timer.html