Spec-Zone.ru › TensorFlow 2.9

tf.distribute.experimental.coordinator.ClusterCoordinator

Объект для планирования и координации выполнения удалённых функций.

Просмотр псевдонимов

Основные псевдонимы

tf.distribute.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 находится в разработке, и API, а также реализация могут быть изменены.

Args
strategy Поддерживаемый объект tf.distribute.Strategy. В настоящее время поддерживается только tf.distribute.experimental.ParameterServerStrategy.
Raises
ValueError Если используемая стратегия не поддерживается.
Attributes
strategy Возвращает Strategy ассоциированный с ClusterCoordinator .

Методы

create_per_worker_dataset

Просмотреть исходный код

create_per_worker_dataset(
    dataset_fn
)

Создание наборов данных на рабочих узлах путём вызова dataset_fn на устройствах рабочих узлов.

Это создаёт заданный набор данных, сгенерированный dataset_fn на рабочих узлах, и возвращает объект, представляющий коллекцию этих отдельных наборов данных. Вызов iter на такой коллекции наборов данных возвращает tf.distribute.experimental.coordinator.PerWorkerValues, который представляет собой коллекцию итераторов, где итераторы размещены на соответствующих рабочих узлах.

Вызов next на коллекции итераторов не поддерживается. Итератор должен быть передан в качестве аргумента в tf.distribute.experimental.coordinator.ClusterCoordinator.schedule. Когда планируемая функция собирается быть выполнена рабочим узлом, функция получит индивидуальный итератор, соответствующий рабочему узлу. Метод next может быть вызван на итераторе внутри планируемой функции, когда итератор является входом функции.

В настоящее время метод schedule предполагает, что все рабочие узлы одинаковы, и поэтому предполагает, что наборы данных на разных рабочих узлах одинаковы, за исключением случаев, когда они могут быть перетасованы по-разному, если они содержат операцию dataset.shuffle и не задан seed для генерации случайных чисел. Из-за этого мы также рекомендуем повторять наборы данных неопределённое количество раз и планировать конечное число шагов вместо того, чтобы полагаться на 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
Args
dataset_fn Функция набора данных, возвращающая набор данных. Она должна выполняться на рабочих узлах.
Returns
Объект, представляющий коллекцию этих отдельных наборов данных. iter ожидается, что вызов будет выполнен на этом объекте, который возвращает tf.distribute.experimental.coordinator.PerWorkerValues итераторов (которые находятся на рабочих узлах).

done

Просмотреть исходный код

done()

Возвращает, завершились ли все запланированные функции.

Если любая из ранее запланированных функций вызывает ошибку, done завершится с ошибкой, вызвав одну из этих ошибок.

Когда done возвращает True или вызывает ошибку, это гарантирует, что нет функций, которые всё ещё выполняются.

Returns
Завершились ли все запланированные функции.
Raises
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
Args
val Значение для получения результатов. Если это структура tf.distribute.experimental.coordinator.RemoteValue, fetch() будет вызвано для каждого отдельного tf.distribute.experimental.coordinator.RemoteValue для получения результата.
Returns
Если 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 возвращает значение или вызывает ошибку, это гарантирует, что нет функций, которые всё ещё выполняются.

Raises
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 Функция tf.function; функция, которая будет асинхронно отправлена на работника для выполнения. Обычная функция Python не поддерживается для планирования.
args Позиционные аргументы для fn.
kwargs Именованные аргументы для fn.
Returns
Объект tf.distribute.experimental.coordinator.RemoteValue, который представляет результат запланированной функции.
Raises
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/versions/r2.9/api_docs/python/tf/distribute/experimental/coordinator/ClusterCoordinator

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API