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
)
Создаёт набор данных на каждом рабочем узле.
Создаёт наборы данных на рабочих узлах из входных данных, которые могут быть либо tf.data.Dataset, либо tf.distribute.DistributedDataset, либо функцией, возвращающей набор данных, и возвращает объект, представляющий коллекцию этих отдельных наборов данных. Вызов 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.Dataset.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.RemoteValues, возвращает полученные значения 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.
| Аргументы | |
|---|---|
fn | tf.function; функция, которая должна быть передана на рабочий узел для асинхронного выполнения. Регулярные python-функции не поддерживаются для планирования. |
args | Позиционные аргументы для fn. |
kwargs | Именованные аргументы для fn. |
| Возвращаемое значение | |
|---|---|
Объект tf.distribute.experimental.coordinator.RemoteValue, представляющий результат запланированной функции. |
| Возможные исключения | |
|---|---|
Exception | одно из исключений, перехваченных координатором любой из ранее запланированных функций с момента последнего выброса ошибки или с начала программы. |
© 2022 The TensorFlow Authors. All rights reserved.
Licensed under the Creative Commons Attribution License 4.0.
Code samples licensed under the Apache 2.0 License.
https://www.tensorflow.org/api_docs/python/tf/distribute/experimental/coordinator/ClusterCoordinator