tf.distribute.experimental.coordinator.ClusterCoordinator
Объект для планирования и координации выполнения удалённых функций.
tf.distribute.experimental.coordinator.ClusterCoordinator(
strategy
)
Этот класс используется для создания отказоустойчивых ресурсов и отправки функций на удалённые серверы TensorFlow.
В настоящее время этот класс не поддерживается для использования в автономном режиме. Его следует использовать совместно со стратегией tf.distribute, разработанной для работы с ним. Класс ClusterCoordinator в настоящее время работает только tf.distribute.experimental.ParameterServerStrategy.
API schedule/join
Наиболее важными API, предоставляемыми этим классом, являются пара schedule/join. API schedule неблокирующий, так как он помещает tf.function в очередь и возвращает RemoteValue немедленно. Очередированные функции будут отправлены удалённым рабочим в фоновых потоках, и их RemoteValue будут заполнены асинхронно. Поскольку schedule не требует назначения рабочих, переданная tf.function может быть выполнена на любом доступном рабочем узле. Если рабочий узел становится недоступным до завершения, функция будет мигрирована на другой рабочий узел. Из-за этого и того, что выполнение функции не атомарно, функция может быть выполнена более одного раза.
Обработка сбоев задач
Этот класс при использовании с tf.distribute.experimental.ParameterServerStrategy имеет встроенную отказоустойчивость к сбоям рабочих узлов. То есть, когда некоторые рабочие узлы по какой-либо причине недоступны для координатора, процесс обучения продолжается с оставшимися рабочими узлами. При восстановлении отказавшего рабочего узла он будет добавлен для выполнения функции после повторной генерации наборов данных, созданных create_per_worker_dataset, на нём.
При отказе параметрического сервера координатор поднимает исключение tf.errors.UnavailableError, генерируемое schedule, join или done. В этом случае, помимо возврата отказавшего параметрического сервера, пользователи должны перезапустить координатор, чтобы он снова подключился к параметрическому серверу, пересоздал переменные и загрузил контрольные точки. Если координатор отказывает, пользователи также должны его вернуть. Программа автоматически подключится к параметрическим серверам и рабочим узлам и продолжит процесс с контрольной точки.
Поэтому крайне важно, чтобы в программе пользователя периодически сохранялся файл контрольной точки и восстанавливался в начале программы. Если tf.keras.optimizers.Optimizer сохраняется в контрольной точке, после восстановления из контрольной точки его свойство iterations приблизительно указывает количество выполненных шагов. Это можно использовать, чтобы определить, сколько эпох и шагов требуется перед завершением обучения.
См. строку документации tf.distribute.experimental.ParameterServerStrategy для примера использования этого API.
В настоящее время находится в разработке, и API, а также реализация могут быть изменены.
| Аргументы | |
|---|---|
strategy | объект поддерживаемой tf.distribute.Strategy. В настоящее время поддерживается только tf.distribute.experimental.ParameterServerStrategy. |
| Исключения | |
|---|---|
ValueError | если используемая стратегия не поддерживается. |
| Атрибуты | |
|---|---|
strategy | Возвращает Strategy ассоциированный с ClusterCoordinator. |
Методы
create_per_worker_dataset
create_per_worker_dataset(
dataset_fn
)
Создаёт набор данных на рабочих узлах, вызывая dataset_fn на устройствах рабочих узлов.
Это создаёт заданный набор данных, сгенерированный dataset_fn на рабочих узлах и возвращает объект, представляющий коллекцию этих отдельных наборов данных. Вызов iter на такой коллекции наборов данных возвращает tf.distribute.experimental.coordinator.PerWorkerValues, который представляет собой коллекцию итераторов, где итераторы размещены на соответствующих рабочих узлах.
Вызов next на PerWorkerValues итераторов не поддерживается. Итератор предназначен для передачи в качестве аргумента в tf.distribute.experimental.coordinator.ClusterCoordinator.schedule. Когда запланированная функция собирается быть выполнена рабочим узлом, функция получит индивидуальный итератор, соответствующий рабочему узлу. Метод next может быть вызван на итераторе внутри запланированной функции, когда итератор является входом функции.
В настоящее время метод schedule предполагает, что рабочие узлы одинаковы, и, следовательно, предполагает, что наборы данных на разных рабочих узлах одинаковы, за исключением того, что они могут быть перемешаны по-разному, если они содержат операцию dataset.shuffle и случайное семя не задано. Из-за этого мы также рекомендуем повторять наборы данных бесконечно и планировать конечное количество шагов вместо того, чтобы полагаться на OutOfRangeError из набора данных.
Пример:
strategy = tf.distribute.experimental.ParameterServerStrategy(
cluster_resolver=...)
coordinator = tf.distribute.experimental.coordinator.ClusterCoordinator(
strategy=strategy)
@tf.function
def worker_fn(iterator):
return next(iterator)
def per_worker_dataset_fn():
return strategy.distribute_datasets_from_function(
lambda x: tf.data.from_tensor_slices([3] * 3)
per_worker_dataset = coordinator.create_per_worker_dataset(
per_worker_dataset_fn)
per_worker_iter = iter(per_worker_dataset)
remote_value = coordinator.schedule(worker_fn, args=(per_worker_iter,))
assert remote_value.fetch() == 3
| Аргументы | |
|---|---|
dataset_fn | Функция набора данных, возвращающая набор данных. Это должно быть выполнено на рабочих узлах. |
| Возвращает | |
|---|---|
Объект, представляющий коллекцию этих отдельных наборов данных. iter ожидается, что будет вызван на этом объекте, который возвращает tf.distribute.experimental.coordinator.PerWorkerValues итераторов (которые находятся на рабочих узлах). |
done
done()
Возвращает, завершили ли все запланированные функции выполнение.
Если любая из ранее запланированных функций вызывает ошибку, done завершится, вызвав одну из этих ошибок.
Когда done возвращает True или вызывает ошибку, гарантируется, что нет функций, которые всё ещё выполняются.
| Возвращает | |
|---|---|
| Закончили ли все запланированные функции выполнение. |
| Исключения | |
|---|---|
Exception | одно из исключений, перехваченных координатором любой из ранее запланированных функций с момента последнего выбрасывания ошибки или с начала программы. |
fetch
fetch(
val
)
Блокирующий вызов для получения результатов из удалённых значений.
Это оболочка вокруг tf.distribute.experimental.coordinator.RemoteValue.fetch для структуры RemoteValue; она возвращает результаты выполнения RemoteValue. Если не готово, ожидает их, блокируя вызывающую функцию.
Пример:
strategy = ...
coordinator = tf.distribute.experimental.coordinator.ClusterCoordinator(
strategy)
def dataset_fn():
return tf.data.Dataset.from_tensor_slices([1, 1, 1])
with strategy.scope():
v = tf.Variable(initial_value=0)
@tf.function
def worker_fn(iterator):
def replica_fn(x):
v.assign_add(x)
return v.read_value()
return strategy.run(replica_fn, args=(next(iterator),))
distributed_dataset = coordinator.create_per_worker_dataset(dataset_fn)
distributed_iterator = iter(distributed_dataset)
result = coordinator.schedule(worker_fn, args=(distributed_iterator,))
assert coordinator.fetch(result) == 1
| Аргументы | |
|---|---|
val | Значение для получения результатов. Если это структура tf.distribute.experimental.coordinator.RemoteValue, fetch() вызывается на отдельных tf.distribute.experimental.coordinator.RemoteValue для получения результата. |
| Возвращает | |
|---|---|
Если val является tf.distribute.experimental.coordinator.RemoteValue или структурой tf.distribute.experimental.coordinator.RemoteValue, возвращает полученные tf.distribute.experimental.coordinator.RemoteValue значения немедленно, если они доступны, или блокирует вызов до тех пор, пока они не станут доступны, и возвращает полученные tf.distribute.experimental.coordinator.RemoteValue значения с той же структурой. Если val другие типы, возвращает их как есть. |
join
join()
Блокирует выполнение, пока все запланированные функции не завершат выполнение.
Если какая-либо из ранее запланированных функций вызывает ошибку, join завершится, вызвав одну из этих ошибок и очистит собранные до этого ошибки. Если это произойдёт, некоторые из ранее запланированных функций могут не быть выполнены. Пользователи могут вызвать fetch на возвращённом tf.distribute.experimental.coordinator.RemoteValue для проверки, были ли они выполнены, завершены или отменены. Если некоторые отменённые функции нужно запланировать повторно, пользователи должны вызвать schedule с функцией снова.
Когда join возвращает или вызывает ошибку, гарантируется, что нет функций, которые всё ещё выполняются.
| Исключения | |
|---|---|
Exception | одно из исключений, перехваченных координатором любой из ранее запланированных функций с момента последнего выбрасывания ошибки или с начала программы. |
schedule
schedule(
fn, args=None, kwargs=None
)
Планирует fn для отправки на рабочий узел для асинхронного выполнения.
Этот метод неблокирующий, так как он помещает fn в очередь, которая будет выполнена позже, и сразу возвращает объект tf.distribute.experimental.coordinator.RemoteValue. fetch можно вызвать на нём, чтобы дождаться завершения выполнения функции и получить её результат от удалённого работника. С другой стороны, вызовите tf.distribute.experimental.coordinator.ClusterCoordinator.join, чтобы дождаться завершения всех запланированных функций.
schedule гарантирует, что fn будет выполнена на работнике как минимум один раз; она может быть выполнена более одного раза, если соответствующий работник потерпит неудачу в середине её выполнения. Обратите внимание, что поскольку работник может потерпеть неудачу в любой момент во время выполнения функции, возможно, что функция выполнена частично, но tf.distribute.experimental.coordinator.ClusterCoordinator гарантирует, что в этих случаях функция в конечном итоге будет выполнена на любом доступном работнике.
Если любая ранее запланированная функция вызывает ошибку, schedule вызывает одну из этих ошибок и очищает собранные ошибки до этого момента. Здесь происходит ситуация, когда некоторые из ранее запланированных функций могут не быть выполнены. Пользователь может вызвать fetch на возвращённом объекте tf.distribute.experimental.coordinator.RemoteValue, чтобы проверить, были ли они выполнены, завершились с ошибкой или отменены, и повторно запланировать соответствующую функцию при необходимости.
Когда schedule возникает, гарантируется, что нет функции, которая всё ещё выполняется.
В настоящее время нет поддержки назначения работников для выполнения функций или приоритета работников.
args и kwargs являются аргументами, переданными в fn, когда fn выполняется на работнике. Они могут быть tf.distribute.experimental.coordinator.PerWorkerValues, и в этом случае аргумент будет заменён соответствующей компонентой на целевом работнике. Аргументы, которые не являются tf.distribute.experimental.coordinator.PerWorkerValues, будут переданы в fn как есть. В настоящее время tf.distribute.experimental.coordinator.RemoteValue не поддерживается в качестве входных args или kwargs.
| Args | |
|---|---|
fn | A tf.function; the function to be dispatched to a worker for execution asynchronously. |
args | Позиционные аргументы для fn. |
kwargs | Именные аргументы для fn. |
| Returns | |
|---|---|
Объект tf.distribute.experimental.coordinator.RemoteValue, представляющий результат запланированной функции. |
| Raises | |
|---|---|
Exception | одну из исключений, пойманных координатором от любой ранее запланированной функции, с момента последнего брошенного исключения или с начала программы. |
© 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.4/api_docs/python/tf/distribute/experimental/coordinator/ClusterCoordinator