Spec-Zone.ru › TensorFlow

tf.compat.v1.distribute.experimental.MultiWorkerMirroredStrategy

Стратегия распределения для синхронного обучения на нескольких рабочих узлах.

Наследуется от: Strategy

tf.compat.v1.distribute.experimental.MultiWorkerMirroredStrategy(
    communication=tf.distribute.experimental.CollectiveCommunication.AUTO,
    cluster_resolver=None
)

Эта стратегия реализует синхронное распределённое обучение на нескольких рабочих узлах, каждый из которых может иметь несколько графических процессоров. Подобно tf.distribute.MirroredStrategy, она дублирует все переменные и вычисления на каждом локальном устройстве. Разница в том, что она использует распределённую коллективную реализацию (например, all-reduce), чтобы несколько рабочих узлов могли работать вместе.

Вам необходимо запустить свою программу на каждом рабочем узле и правильно настроить cluster_resolver. Например, если вы используете tf.distribute.cluster_resolver.TFConfigClusterResolver, каждый рабочий узел должен иметь соответствующие task_type и task_id, установленные в переменной среды TF_CONFIG. Пример TF_CONFIG на рабочем узле-0 двух-узловой кластера:

TF_CONFIG = '{"cluster": {"worker": ["localhost:12345", "localhost:23456"]}, "task": {"type": "worker", "index": 0} }'

Ваша программа работает на каждом рабочем узле как есть. Обратите внимание, что для коллективных операций необходимо участие каждого рабочего узла. Все tf.distribute и не tf.distribute API могут использовать коллективные операции внутри, например, для создания контрольных точек и сохранения, так как чтение tf.Variable с tf.VariableSynchronization.ON_READ выполняет all-reduce значения. Поэтому рекомендуется запускать точно одну и ту же программу на каждом рабочем узле. Распределение задач на основе task_type или task_id рабочего узла чревато ошибками.

cluster_resolver.num_accelerators() определяет количество графических процессоров, используемых стратегией. Если значение равно нулю, стратегия использует центральный процессор. Все рабочие узлы должны использовать одинаковое количество устройств, в противном случае поведение не определено.

Эта стратегия не предназначена для TPU. Используйте tf.distribute.TPUStrategy вместо неё.

После настройки TF_CONFIG использование этой стратегии аналогично использованию tf.distribute.MirroredStrategy и tf.distribute.TPUStrategy.

strategy = tf.distribute.MultiWorkerMirroredStrategy()

with strategy.scope():
  model = tf.keras.Sequential([
    tf.keras.layers.Dense(2, input_shape=(5,)),
  ])
  optimizer = tf.keras.optimizers.SGD(learning_rate=0.1)

def dataset_fn(ctx):
  x = np.random.random((2, 5)).astype(np.float32)
  y = np.random.randint(2, size=(2, 1))
  dataset = tf.data.Dataset.from_tensor_slices((x, y))
  return dataset.repeat().batch(1, drop_remainder=True)
dist_dataset = strategy.distribute_datasets_from_function(dataset_fn)

model.compile()
model.fit(dist_dataset)

Вы также можете написать свой собственный цикл обучения:

@tf.function
def train_step(iterator):

  def step_fn(inputs):
    features, labels = inputs
    with tf.GradientTape() as tape:
      logits = model(features, training=True)
      loss = tf.keras.losses.sparse_categorical_crossentropy(
          labels, logits)

    grads = tape.gradient(loss, model.trainable_variables)
    optimizer.apply_gradients(zip(grads, model.trainable_variables))

  strategy.run(step_fn, args=(next(iterator),))

for _ in range(NUM_STEP):
  train_step(iterator)

Подробный учебник см. в разделе Обучение с использованием нескольких рабочих узлов с Keras.

Сохранение

Вам необходимо сохранять и создавать контрольные точки на всех рабочих узлах, а не только на одном. Это происходит потому, что переменные, для которых synchronization=ON_READ, вызывают агрегацию при сохранении. Рекомендуется сохранять в разные пути на каждом рабочем узле, чтобы избежать гонок. Каждый рабочий узел сохраняет то же самое. Примеры см. в учебнике Обучение с использованием нескольких рабочих узлов с Keras.

Известные проблемы

  • tf.distribute.cluster_resolver.TFConfigClusterResolver не возвращает правильное количество ускорителей. Стратегия использует все доступные графические процессоры, если cluster_resolver имеет значение tf.distribute.cluster_resolver.TFConfigClusterResolver или None.
  • В режиме eager стратегия должна быть создана до вызова любого другого API Tensorflow.
Атрибуты
cluster_resolver Возвращает резольвер кластера, связанный с этой стратегией.

В общем случае, при использовании стратегии распределения для нескольких рабочих узлов, например, tf.distribute.experimental.MultiWorkerMirroredStrategy или tf.distribute.TPUStrategy(), существует связанный с стратегией tf.distribute.cluster_resolver.ClusterResolver, и такая экземпляр возвращается этим свойством.

Стратегии, которые намереваются иметь связанный tf.distribute.cluster_resolver.ClusterResolver, должны установить соответствующий атрибут или переопределить это свойство; в противном случае, по умолчанию возвращается None. Эти стратегии также должны предоставить информацию о том, что возвращается этим свойством.

Стратегии для одного рабочего узла обычно не имеют tf.distribute.cluster_resolver.ClusterResolver, и в этих случаях это свойство возвращает None.

При необходимости получить информацию, такую как спецификация кластера, тип задачи или идентификатор задачи, полезно использовать tf.distribute.cluster_resolver.ClusterResolver. Например,

os.environ['TF_CONFIG'] = json.dumps({
  'cluster': {
      'worker': ["localhost:12345", "localhost:23456"],
      'ps': ["localhost:34567"]
  },
  'task': {'type': 'worker', 'index': 0}
})

# This implicitly uses TF_CONFIG for the cluster and current task info.
strategy = tf.distribute.experimental.MultiWorkerMirroredStrategy()

...

if strategy.cluster_resolver.task_type == 'worker':
  # Perform something that's only applicable on workers. Since we set this
  # as a worker above, this block will run on this particular instance.
elif strategy.cluster_resolver.task_type == 'ps':
  # Perform something that's only applicable on parameter servers. Since we
  # set this as a worker above, this block will not run on this particular
  # instance.

Дополнительную информацию можно найти в документации API tf.distribute.cluster_resolver.ClusterResolver.

extended tf.distribute.StrategyExtended с дополнительными методами.
num_replicas_in_sync Возвращает число реплик, по которым агрегируются градиенты.

Методы

distribute_datasets_from_function

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

distribute_datasets_from_function(
    dataset_fn, options=None
)

Распределяет экземпляры tf.data.Dataset, созданные вызовами к dataset_fn.

Передаваемая пользователем переменная dataset_fn — это функция ввода, которая принимает аргумент tf.distribute.InputContext и возвращает экземпляр tf.data.Dataset. Ожидается, что возвращаемый из dataset_fn набор данных уже разделен по размеру пакета на реплику (т. е. общий размер пакета, деленный на количество реплик в синхронизации) и фрагментирован. tf.distribute.Strategy.distribute_datasets_from_function не выполняет разделение и группировку экземпляра набора данных, возвращаемого функцией ввода. dataset_fn будет вызвана на устройстве процессора каждого из рабочих узлов, и каждый из них сгенерирует набор данных, где каждая реплика на этом рабочем узле будет извлекать по одному пакету ввода (т. е. если рабочий узел имеет две реплики, то с Dataset на каждом шаге будет извлечено два пакета).

Этот метод может использоваться для нескольких целей. Во-первых, он позволяет указать собственную логику разделения и группировки. (В отличие от tf.distribute.experimental_distribute_dataset, которая выполняет группировку и разделение за вас.) Например, там, где experimental_distribute_dataset не может разделить входные файлы, этот метод может использоваться для ручного разделения набора данных (избегая медленного поведения по умолчанию в experimental_distribute_dataset). В тех случаях, когда набор данных является бесконечным, это разделение может выполняться путём создания реплик набора данных, отличающихся только начальным значением генератора случайных чисел.

dataset_fn должна принимать экземпляр tf.distribute.InputContext, где можно получить информацию о группировке и дублировании входных данных.

Вы можете использовать свойство element_spec возвращённого tf.distribute.DistributedDataset API для получения tf.TypeSpec элементов, возвращаемых итератором. Это можно использовать для установки свойства input_signature функции tf.function. См. tf.distribute.DistributedDataset.element_spec для примера.

Важно: tf.data.Dataset, возвращаемый dataset_fn, должен иметь размер пакета на реплику, в отличие от experimental_distribute_dataset, который использует глобальный размер пакета. Это можно вычислить с помощью input_context.get_per_replica_batch_size.
Примечание: Если вы используете TPUStrategy, порядок обработки данных рабочими узлами при использовании tf.distribute.Strategy.experimental_distribute_dataset или tf.distribute.Strategy.distribute_datasets_from_function не гарантируется. Это обычно требуется, если вы используете tf.distribute для масштабирования прогнозирования. Однако вы можете вставить индекс для каждого элемента в пакете и упорядочить результаты соответственно. См. этот пример для примера того, как упорядочить результаты.
Примечание: Состояние преобразований набора данных в настоящее время не поддерживается с tf.distribute.experimental_distribute_dataset или tf.distribute.distribute_datasets_from_function. Любые операции состояния, которые может иметь набор данных, в настоящее время игнорируются. Например, если ваш набор данных содержит преобразование map_fn, которое использует tf.random.uniform для поворота изображения, то у вас есть граф набора данных, зависящий от состояния (т. е. от генератора случайных чисел) на локальном компьютере, где выполняется процесс Python.

Для ознакомления с дополнительным использованием и свойствами этого метода, обратитесь к тьюториалу по распределенному вводу). Если вас интересует обработка последней частичной пачки, прочитайте эту секцию.

Аргументы
dataset_fn Функция, принимающая экземпляр tf.distribute.InputContext и возвращающая tf.data.Dataset.
options tf.distribute.InputOptions, используемый для управления параметрами распределения этого набора данных.
Возвращаемое значение
tf.distribute.DistributedDataset.

experimental_distribute_dataset

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

experimental_distribute_dataset(
    dataset, options=None
)

Создаёт tf.distribute.DistributedDataset из tf.data.Dataset.

Возвращённый tf.distribute.DistributedDataset может быть итерирован подобно обычным наборам данных. ПРИМЕЧАНИЕ: Пользователь не может добавить больше преобразований к tf.distribute.DistributedDataset. Вы можете только создать итератор или проверить tf.TypeSpec данных, генерируемых им. См. документацию API tf.distribute.DistributedDataset для получения дополнительной информации.

Вот пример:

global_batch_size = 2
# Passing the devices is optional.
strategy = tf.distribute.MirroredStrategy(devices=["GPU:0", "GPU:1"])
# Create a dataset
dataset = tf.data.Dataset.range(4).batch(global_batch_size)
# Distribute that dataset
dist_dataset = strategy.experimental_distribute_dataset(dataset)
@tf.function
def replica_fn(input):
  return input*2
result = []
# Iterate over the `tf.distribute.DistributedDataset`
for x in dist_dataset:
  # process dataset elements
  result.append(strategy.run(replica_fn, args=(x,)))
print(result)
[PerReplica:{
  0: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([0])>,
  1: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([2])>
}, PerReplica:{
  0: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([4])>,
  1: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([6])>
}]

Три ключевых действия, происходящих под капотом этого метода, — это группировка, фрагментация и предварительная загрузка.

В приведенном выше фрагменте кода dataset группируется по global_batch_size, а вызов experimental_distribute_dataset на нём перегруппировывает dataset в новый размер пачки, равный глобальному размеру пачки, делённому на количество реплик в синхронизации. Мы перебираем его с помощью цикла Python. x — это tf.distribute.DistributedValues, содержащий данные для всех реплик, и каждая реплика получает данные нового размера пачки. tf.distribute.Strategy.run позаботится о предоставлении правильных данных на реплику в x в соответствующей replica_fn, выполняемой на каждой реплике.

Фрагментация включает в себя автоматическое фрагментирование по нескольким рабочим узлам и внутри каждого рабочего узла. Во-первых, при распределённом обучении с несколькими рабочими узлами (т. е. при использовании tf.distribute.experimental.MultiWorkerMirroredStrategy или tf.distribute.TPUStrategy) автоматическое фрагментирование набора данных по нескольким рабочим узлам означает, что каждому рабочему узлу назначается подмножество всего набора данных (если установлен правильный tf.data.experimental.AutoShardPolicy). Это делается для того, чтобы на каждом шаге глобальный размер пачки неперекрывающихся элементов набора данных обрабатывался каждым рабочим узлом. Автоматическое фрагментирование имеет несколько различных параметров, которые можно указать с помощью tf.data.experimental.DistributeOptions.

Примечание: для автоматического фрагментирования по нескольким рабочим узлам, по умолчанию используется режим tf.data.experimental.AutoShardPolicy.AUTO. Этот режим попытается фрагментировать набор данных по файлам, если набор данных создаётся из наборов данных читателей (например, tf.data.TFRecordDataset, tf.data.TextLineDataset и т. д.) или фрагментировать набор данных по данным, где каждый из рабочих узлов прочитает весь набор данных и обработает только назначенный ему фрагмент. Однако, если у вас меньше одного входного файла на рабочий узел, мы рекомендуем отключить автоматическое фрагментирование набора данных по рабочим узлам, установив tf.data.experimental.DistributeOptions.auto_shard_policy в tf.data.experimental.AutoShardPolicy.OFF.

По умолчанию этот метод добавляет преобразование предварительной загрузки в конце предоставленного пользователем экземпляра tf.data.Dataset. Аргументом преобразования предварительной загрузки, который является buffer_size, является количество реплик в синхронизации.

Если логика разделения пачек и фрагментации набора данных выше нежелательна, пожалуйста, используйте tf.distribute.Strategy.distribute_datasets_from_function вместо этого, который не выполняет автоматического группирования или фрагментирования.

Примечание: Если вы используете TPUStrategy, порядок обработки данных рабочими узлами при использовании tf.distribute.Strategy.experimental_distribute_dataset или tf.distribute.Strategy.distribute_datasets_from_function не гарантируется. Это обычно требуется, если вы используете tf.distribute для масштабирования прогнозирования. Однако вы можете вставить индекс для каждого элемента в пачке и упорядочить результаты соответствующим образом. Обратитесь к этому фрагменту для примера того, как упорядочить результаты.
Примечание: Трансформации наборов данных с состоянием в настоящее время не поддерживаются с tf.distribute.experimental_distribute_dataset или tf.distribute.distribute_datasets_from_function. Любые операторы с состоянием, которые может иметь набор данных, в настоящее время игнорируются. Например, если ваш набор данных имеет map_fn, который использует tf.random.uniform для поворота изображения, у вас есть граф набора данных, зависящий от состояния (т. е. начального значения генератора случайных чисел) на локальной машине, где выполняется процесс Python.

Для ознакомления с дополнительным использованием и свойствами этого метода, обратитесь к тьюториалу по распределенному вводу. Если вас интересует обработка последней частичной пачки, прочитайте эту секцию.

Аргументы
dataset tf.data.Dataset, который будет фрагментирован по всем репликам в соответствии с описанными выше правилами.
options tf.distribute.InputOptions для управления параметрами распределения набора данных.
Возвращаемое значение
tf.distribute.DistributedDataset.

experimental_local_results

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

experimental_local_results(
    value
)

Возвращает список всех локальных значений на реплику, содержащихся в value.

Примечание: Это возвращает только значения на рабочем узле, инициированном этим клиентом. При использовании tf.distribute.Strategy, такого как tf.distribute.experimental.MultiWorkerMirroredStrategy, каждый рабочий узел будет своим клиентом, и эта функция вернёт только значения, вычисленные на этом рабочем узле.
Аргументы
value Значение, возвращённое experimental_run(), run(), or a variable created inscope`.
Возвращаемое значение
Кортеж значений, содержащихся в value, где i-й элемент соответствует i-й реплике. Если value представляет одно значение, возвращается (value,).

experimental_make_numpy_dataset

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

experimental_make_numpy_dataset(
    numpy_input, session=None
)

Создаёт tf.data.Dataset для ввода, предоставленного через массив NumPy.

Это предотвращает добавление numpy_input как большой константы в граф и копирует данные на машину или машины, которые будут обрабатывать ввод.

Обратите внимание, что, вероятно, вам потребуется использовать tf.distribute.Strategy.experimental_distribute_dataset с возвращённым набором данных для его дальнейшего распределения с помощью стратегии.

Пример:

numpy_input = np.ones([10], dtype=np.float32)
dataset = strategy.experimental_make_numpy_dataset(numpy_input)
dist_dataset = strategy.experimental_distribute_dataset(dataset)
Аргументы
numpy_input Вложенный набор массивов NumPy, которые будут преобразованы в набор данных. Обратите внимание, что списки массивов NumPy складываются, так как это обычное поведение tf.data.Dataset.
session (Только для выполнения графов TensorFlow v1.x) Сессия, используемая для инициализации.
Возвращаемое значение
tf.data.Dataset, представляющий numpy_input.

experimental_run

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

experimental_run(
    fn, input_iterator=None
)

Выполняет операции в fn на каждой реплике, используя входные данные из input_iterator. (устарело)

Устарело: ЭТА ФУНКЦИЯ УСТАРЕЛА. Она будет удалена в будущей версии. Инструкции по обновлению: Этот метод недоступен в TF 2.x. Пожалуйста, переключитесь на использование run вместо него.
Устарело: Этот метод недоступен в TF 2.x. Пожалуйста, переключитесь на использование run вместо него.

При включенном режиме выполнения Eager, выполняет операции, указанные в fn, на каждой реплике. В противном случае строит граф для выполнения операций на каждой реплике.

Каждая реплика получит один уникальный входной параметр из входных данных, предоставленных одним вызовом get_next для итератора входных данных.

fn может вызвать tf.distribute.get_replica_context() для доступа к элементам, таким как replica_id_in_sync_group.

Важно: В зависимости от используемой реализации tf.distribute.Strategy и включения режима eager execution, fn может быть вызван один или несколько раз (по одному разу для каждой реплики).
Аргументы
fn Функция для выполнения. Входные данные функции должны соответствовать выходам input_iterator.get_next(). Выход должен быть tf.nest из Tensor.
input_iterator (Необязательно) итератор входных данных, из которого берутся входные значения.
Возвращаемое значение
Объединённое значение возврата fn по всем репликам. Структура возвращаемого значения такая же, как у возвращаемого значения от fn. Каждый элемент структуры может быть PerReplica (если значения не синхронизированы), Mirrored (если значения синхронизированы), или Tensor (если выполняется на одной реплике).

make_dataset_iterator

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

make_dataset_iterator(
    dataset
)

Создаёт итератор для входных данных, предоставленных через dataset.

Устарело: Этот метод недоступен в TF 2.x.

Данные из заданного набора данных будут распределены равномерно по всем вычислительным репликам. Предполагается, что входной набор данных сгруппирован по глобальному размеру пакета. С этим предположением мы сделаем всё возможное, чтобы разделить каждый пакет по всем репликам (одному или нескольким работникам). Если эта попытка провалится, будет выброшено исключение, и пользователь должен вместо этого использовать make_input_fn_iterator, которое предоставляет больше контроля пользователю и не пытается разделить пакет по репликам.

Пользователь также может использовать make_input_fn_iterator, если хочет настроить, какой вход подаётся на какую реплику/работник и т. д.

Аргументы
dataset tf.data.Dataset, который будет распределён равномерно по всем репликам.
Возвращаемое значение
Итератор tf.distribute.InputIterator, который возвращает входные данные для каждого шага вычисления. Пользователь должен вызвать initialize на возвращённом итераторе.

make_input_fn_iterator

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

make_input_fn_iterator(
    input_fn,
    replication_mode=tf.distribute.InputReplicationMode.PER_WORKER
)

Возвращает итератор, разделённый по репликам, созданный из функции входных данных.

Устарело: Этот метод недоступен в TF 2.x.

Функция input_fn должна принимать объект tf.distribute.InputContext, где можно получить информацию о пакета и фрагментации входных данных:

def input_fn(input_context):
  batch_size = input_context.get_per_replica_batch_size(global_batch_size)
  d = tf.data.Dataset.from_tensors([[1.]]).repeat().batch(batch_size)
  return d.shard(input_context.num_input_pipelines,
                 input_context.input_pipeline_id)
with strategy.scope():
  iterator = strategy.make_input_fn_iterator(input_fn)
  replica_results = strategy.experimental_run(replica_fn, iterator)

Возвращаемый tf.data.Dataset объектом input_fn должен иметь размер пакета на реплику, который может быть вычислен с помощью input_context.get_per_replica_batch_size.

Аргументы
input_fn Функция, принимающая объект tf.distribute.InputContext и возвращающая tf.data.Dataset.
replication_mode значение перечисления tf.distribute.InputReplicationMode. В настоящее время поддерживается только PER_WORKER, что означает, что будет один вызов input_fn на каждый работник. Реплики будут извлекать данные из локального tf.data.Dataset на своих работниках.
Возвращаемое значение
Объект итератора, который сначала необходимо вызвать методом .initialize(). Затем он может быть передан в strategy.experimental_run() или вы можете использовать iterator.get_next() для получения следующего значения, которое нужно передать в strategy.extended.call_for_each_replica().

reduce

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

reduce(
    reduce_op, value, axis=None
)

Сводит value по репликам и возвращает результат на текущем устройстве.

strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
def step_fn():
  i = tf.distribute.get_replica_context().replica_id_in_sync_group
  return tf.identity(i)

per_replica_result = strategy.run(step_fn)
total = strategy.reduce("SUM", per_replica_result, axis=None)
total
<tf.Tensor: shape=(), dtype=int32, numpy=1>

Чтобы увидеть, как это будет выглядеть с несколькими репликами, рассмотрим тот же пример с MirroredStrategy с 2-мя графическими процессорами:

strategy = tf.distribute.MirroredStrategy(devices=["GPU:0", "GPU:1"])
def step_fn():
  i = tf.distribute.get_replica_context().replica_id_in_sync_group
  return tf.identity(i)

per_replica_result = strategy.run(step_fn)
# Check devices on which per replica result is:
strategy.experimental_local_results(per_replica_result)[0].device
# /job:localhost/replica:0/task:0/device:GPU:0
strategy.experimental_local_results(per_replica_result)[1].device
# /job:localhost/replica:0/task:0/device:GPU:1

total = strategy.reduce("SUM", per_replica_result, axis=None)
# Check device on which reduced result is:
total.device
# /job:localhost/replica:0/task:0/device:CPU:0

Этот API обычно используется для агрегирования результатов, возвращаемых различными репликами, например, для отчётности. Например, потеря, вычисленная из различных реплик, может быть усреднена с помощью этого API перед печатью.

Примечание: Результат копируется на "текущее" устройство — обычно это процессор на работнике, на котором выполняется программа. Для TPUStrategy это первый хост TPU. Для многоклиентского MultiWorkerMirroredStrategy это процессор каждого работника.

Существует несколько различных API tf.distribute для сведения значений по репликам:

  • tf.distribute.ReplicaContext.all_reduce: Отличается от Strategy.reduce тем, что предназначен для контекста реплики и не копирует результаты на устройство хоста. all_reduce обычно используется для редукций внутри шага обучения, таких как градиенты.
  • tf.distribute.StrategyExtended.reduce_to и tf.distribute.StrategyExtended.batch_reduce_to: Эти API являются более продвинутыми версиями Strategy.reduce, поскольку позволяют настраивать место назначения результата. Они также вызываются в межрепличном контексте.

Какой должен быть ось?

Учитывая значение на реплику, возвращаемое run, например, потерю на пример, пакет будет разделён по всем репликам. Эта функция позволяет агрегировать по репликам и, при необходимости, по элементам пакета, указав параметр оси соответствующим образом.

Например, если у вас есть глобальный размер пакета 8 и 2 реплики, значения для примеров [0, 1, 2, 3] будут на реплике 0, а [4, 5, 6, 7] — на реплике 1. С axis=None, reduce будет агрегировать только по репликам, возвращая [0+4, 1+5, 2+6, 3+7]. Это полезно, когда каждая реплика вычисляет скаляр или какое-либо другое значение, у которого нет «размерности пакета» (например, градиент или потеря).

strategy.reduce("sum", per_replica_result, axis=None)

Иногда вам нужно будет агрегировать как по глобальному пакету, так и по всем репликам. Этого можно добиться, указав размер пакета как axis, обычно axis=0. В этом случае он вернёт скалярное значение 0+1+2+3+4+5+6+7.

strategy.reduce("sum", per_replica_result, axis=0)

Если есть последний частичный пакет, вам необходимо указать ось, чтобы размер получающегося массива был согласован между репликами. Так, если последний пакет имеет размер 6 и разделен на [0, 1, 2, 3] и [4, 5], у вас возникнет несоответствие размеров, если вы не укажете axis=0. Если вы укажете tf.distribute.ReduceOp.MEAN, с использованием axis=0 будет использоваться правильный знаменатель 6. В отличие от вычисления reduce_mean для получения скалярного значения на каждой реплике и этой функции для усреднения этих средних значений, что будет учитывать некоторые значения 1/8 и другие значения 1/4.

Аргументы
reduce_op значение tf.distribute.ReduceOp, определяющее, как должны быть объединены значения. Разрешает использовать строковое представление перечисления, например, "SUM", "MEAN".
value экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который должен быть объединён в один тензор. Также может быть обычным тензором, если используется с OneDeviceStrategy или стратегией по умолчанию.
axis определяет размерность, по которой должно происходить сведение внутри тензора каждой реплики. Обычно устанавливается на размерность пакета или None, чтобы выполнить сведение только по репликам (например, если тензор не имеет размерности пакета).
Возвращаемое значение
Tensor.

run

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

run(
    fn, args=(), kwargs=None, options=None
)

Вызывает fn на каждой реплике с указанными аргументами.

Этот метод является основным способом распределения вычислений с помощью объекта tf.distribute. Он вызывает fn на каждой реплике. Если args или kwargs имеют tf.distribute.DistributedValues, такие, как полученные из tf.distribute.DistributedDataset из tf.distribute.Strategy.experimental_distribute_dataset или tf.distribute.Strategy.distribute_datasets_from_function, когда fn выполняется на конкретной реплике, он будет выполнен с компонентом tf.distribute.DistributedValues, соответствующим этой реплике.

fn вызывается в контексте реплики. fn может вызвать tf.distribute.get_replica_context() для доступа к членам, таким как all_reduce. Подробнее о контексте реплики см. в документации модуля tf.distribute.

Все аргументы в args или kwargs могут быть вложенной структурой тензоров, например, списком тензоров, в этом случае args и kwargs будут переданы в вызываемую на каждой реплике fn. Или args или kwargs могут быть tf.distribute.DistributedValues, содержащими тензоры или составные тензоры, т.е. tf.compat.v1.TensorInfo.CompositeTensor, в этом случае каждый вызов fn получит компонент tf.distribute.DistributedValues, соответствующий его реплике. Обратите внимание, что произвольные значения Python, которые не являются вышеперечисленными типами, не поддерживаются.

Важно: В зависимости от реализации tf.distribute.Strategy и от того, включено ли выполнение в режиме eager, fn может быть вызвано один или несколько раз. Если fn анотировано tf.function или tf.distribute.Strategy.run вызывается внутри tf.function (выполнение в режиме eager отключено внутри tf.function по умолчанию), fn вызывается один раз на реплику для генерации графа TensorFlow, который затем будет повторно использован для выполнения с новыми входными данными. В противном случае, если выполнение в режиме eager включено, fn будет вызываться один раз на реплику на каждом шаге, как обычный код Python.

Пример использования:

  1. Входной тензор-константа.

    strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
    tensor_input = tf.constant(3.0)
    @tf.function
    def replica_fn(input):
      return input*2.0
    result = strategy.run(replica_fn, args=(tensor_input,))
    result
        PerReplica:{
          0: <tf.Tensor: shape=(), dtype=float32, numpy=6.0>,
          1: <tf.Tensor: shape=(), dtype=float32, numpy=6.0>
        }
        
  2. Входные данные DistributedValues.

    strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
    @tf.function
    def run():
      def value_fn(value_context):
        return value_context.num_replicas_in_sync
      distributed_values = (
        strategy.experimental_distribute_values_from_function(
          value_fn))
      def replica_fn2(input):
        return input*2
      return strategy.run(replica_fn2, args=(distributed_values,))
    result = run()
    result
        <tf.Tensor: shape=(), dtype=int32, numpy=4>
        
  3. Использование tf.distribute.ReplicaContext для allreduce значений.

    strategy = tf.distribute.MirroredStrategy(["gpu:0", "gpu:1"])
    @tf.function
    def run():
       def value_fn(value_context):
         return tf.constant(value_context.replica_id_in_sync_group)
       distributed_values = (
           strategy.experimental_distribute_values_from_function(
               value_fn))
       def replica_fn(input):
         return tf.distribute.get_replica_context().all_reduce(
             "sum", input)
       return strategy.run(replica_fn, args=(distributed_values,))
    result = run()
    result
        PerReplica:{
          0: <tf.Tensor: shape=(), dtype=int32, numpy=1>,
          1: <tf.Tensor: shape=(), dtype=int32, numpy=1>
        }
        
Аргументы
fn Функция, которую нужно запустить на каждой реплике.
args Необязательные позиционные аргументы для fn. Его элемент может быть тензором, вложенной структурой тензоров или tf.distribute.DistributedValues.
kwargs Необязательные именованные аргументы для fn. Его элемент может быть тензором, вложенной структурой тензоров или tf.distribute.DistributedValues.
options Необязательный экземпляр tf.distribute.RunOptions, указывающий опции для запуска fn.
Возвращает
Объединенное значение возврата fn по всем репликам. Структура возвращаемого значения такая же, как и значение возврата fn. Каждый элемент структуры может быть tf.distribute.DistributedValues, объектами Tensor или Tensor (например, при выполнении на одной реплике).

scope

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

scope()

Управляющая конструкция для установки текущей стратегии и распределения переменных.

Этот метод возвращает управляющую конструкцию, и используется следующим образом:

strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
# Variable created inside scope:
with strategy.scope():
  mirrored_variable = tf.Variable(1.)
mirrored_variable
MirroredVariable:{
  0: <tf.Variable 'Variable:0' shape=() dtype=float32, numpy=1.0>,
  1: <tf.Variable 'Variable/replica_1:0' shape=() dtype=float32, numpy=1.0>
}
# Variable created outside scope:
regular_variable = tf.Variable(1.)
regular_variable
<tf.Variable 'Variable:0' shape=() dtype=float32, numpy=1.0>

Что происходит при входе в область действия Strategy.scope?

  • strategy устанавливается в глобальном контексте как «текущая» стратегия. Внутри этой области tf.distribute.get_strategy() теперь вернёт эту стратегию. За пределами этой области он возвращает стратегию по умолчанию no-op.
  • Вход в область также приводит к входу в «межрепличный контекст». Подробнее о межрепличных и репличных контекстах см. в tf.distribute.StrategyExtended.
  • Создание переменных внутри scope перехватывается стратегией. Каждая стратегия определяет, как она хочет повлиять на создание переменной. Синхронные стратегии, такие как MirroredStrategy, TPUStrategy и MultiWorkerMiroredStrategy создают переменные, дублированные на каждой реплике, в то время как ParameterServerStrategy создаёт переменные на параметровых серверах. Это выполняется с помощью настраиваемого tf.variable_creator_scope.
  • В некоторых стратегиях также может быть введён контекст по умолчанию для устройства: в MultiWorkerMiroredStrategy на каждом работнике вводится контекст по умолчанию для устройства "/CPU:0".
Примечание: Вход в область не автоматически распределяет вычисления, за исключением случаев высокоуровневых обучающих фреймворков, таких как keras model.fit. Если вы не используете model.fit, вам необходимо использовать API strategy.run, чтобы явно распределить эти вычисления. Пример см. в учебнике по созданию пользовательской петли обучения https://www.tensorflow.org/tutorials/distribute/custom_training.

Что должно быть внутри, а что снаружи области?

Есть ряд требований к тому, что должно произойти внутри области. Однако в тех местах, где у нас есть информация о используемой стратегии, мы часто входим в область для пользователя, чтобы он не делал это явно (т.е. вызов внутри или вне области приемлем).

  • Любой код, создающий переменные, которые должны быть распределёнными переменными, должен вызываться в области strategy.scope. Это можно сделать, либо вызвав функцию создания переменной непосредственно в контексте области, либо полагаясь на другой API, например, strategy.run или keras.Model.fit, для автоматического входа в него. Любые переменные, созданные вне области, не будут распределены и могут повлиять на производительность. Некоторые общие объекты, создающие переменные в TF, это модели, оптимизаторы, метрики. Такие объекты всегда должны инициализироваться в области, а любые функции, которые могут лениво создавать переменные (например, Model.call(), отслеживание tf.function и т.д.), аналогичным образом должны вызываться в области. Ещё один источник создания переменных — восстановление контрольной точки; когда переменные создаются лениво. Обратите внимание, что любая переменная, созданная внутри стратегии, сохраняет информацию о стратегии. Поэтому чтение и запись этих переменных вне strategy.scope также могут работать беспрепятственно, без необходимости входа пользователя в область.
  • Некоторые API стратегии (например, strategy.run и strategy.reduce), которые требуют наличия области стратегии, автоматически входят в неё, что означает, что при использовании этих API вам не нужно явно входить в область.
  • Когда tf.keras.Model создаётся внутри strategy.scope, объект модели сохраняет информацию об области. Когда затем вызываются методы высокоуровневого фреймворка обучения, такие как model.compile, model.fit и т.д., захваченная область будет автоматически введена, и используемая стратегия будет использована для распределения обучения и т.д. Подробный пример см. в учебнике по распределённому Keras. ПРЕДУПРЕЖДЕНИЕ: Простое вызов model(..) не обеспечивает автоматического входа в захваченную область — только высокоуровневые API обучения поддерживают такое поведение: model.compile, model.fit, model.evaluate, model.predict и model.save могут вызываться внутри или вне области.
  • Следующее может быть как внутри, так и вне области:
    • Создание наборов входных данных
    • Определение tf.function, представляющих ваш обучающий шаг
    • Сохранение API, такие как tf.saved_model.save. Загрузка создаёт переменные, поэтому это должно находиться внутри области, если вы хотите обучить модель в распределённом режиме.
    • Сохранение контрольных точек. Как упоминалось выше, checkpoint.restore иногда может потребоваться находиться внутри области, если он создаёт переменные.
Возвращает
Управляющая конструкция.

update_config_proto

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

update_config_proto(
    config_proto
)

Возвращает копию config_proto, изменённую для использования с этой стратегией.

Устарело: Этот метод недоступен в TF 2.x.

Обновлённая конфигурация содержит что-то необходимое для работы стратегии, например, настройки для запуска коллективных операций или фильтры устройств для повышения производительности распределённого обучения.

Аргументы
config_proto объект tf.ConfigProto.
Возвращаемое значение
Обновленная копия config_proto.

© 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/compat/v1/distribute/experimental/MultiWorkerMirroredStrategy

Spec-Zone.ru

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