Spec-Zone.ru › TensorFlow 2.4

tf.distribute.MultiWorkerMirroredStrategy

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

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

tf.distribute.MultiWorkerMirroredStrategy(
    cluster_resolver=None, communication_options=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.cluster_resolver.ClusterResolver. Если None, используется tf.distribute.cluster_resolver.TFConfigClusterResolver.
communication_options Необязательные tf.distribute.experimental.CommunicationOptions. Это настраивает опции по умолчанию для межъузеловых коммуникаций. Их можно переопределить, указав опции в API для коммуникации, например, tf.distribute.ReplicaContext.all_reduce. Подробности см. в tf.distribute.experimental.CommunicationOptions.
Атрибуты
cluster_resolver Возвращает решатель кластера, связанный с этой стратегией.

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

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 возвращаемого этим API tf.distribute.DistributedDataset, чтобы получить 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 до нового размера пакета, равного глобальному размеру пакета, деленному на количество реплик в синхронизации. Мы проходим по нему с помощью цикла for по-питонически. 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.

По умолчанию этот метод добавляет преобразование prefetch в конце предоставленного пользователем экземпляра tf.data.Dataset. Аргумент для преобразования prefetch, который является 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_distribute_values_from_function

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

experimental_distribute_values_from_function(
    value_fn
)

Генерирует tf.distribute.DistributedValues из value_fn.

Эта функция предназначена для генерации tf.distribute.DistributedValues для передачи в run, reduce или другие методы, принимающие распределенные значения, при отсутствии наборов данных.

Аргументы
value_fn Функция, которая выполняется для генерации значений. Она вызывается для каждой реплики с tf.distribute.ValueContext в качестве единственного аргумента. Она должна возвращать тензор или тип, который может быть преобразован в тензор.
Возвращает
tf.distribute.DistributedValues, содержащий значение для каждой реплики.

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

  1. Возврат постоянного значения для каждой реплики:
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
def value_fn(ctx):
  return tf.constant(1.)
distributed_values = (
     strategy.experimental_distribute_values_from_function(
       value_fn))
local_result = strategy.experimental_local_results(distributed_values)
local_result
(<tf.Tensor: shape=(), dtype=float32, numpy=1.0>,
 <tf.Tensor: shape=(), dtype=float32, numpy=1.0>)
  1. Распределение значений в массиве на основе идентификатора реплики:
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
array_value = np.array([3., 2., 1.])
def value_fn(ctx):
  return array_value[ctx.replica_id_in_sync_group]
distributed_values = (
     strategy.experimental_distribute_values_from_function(
       value_fn))
local_result = strategy.experimental_local_results(distributed_values)
local_result
(3.0, 2.0)
  1. Указание значений с использованием num_replicas_in_sync:
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
def value_fn(ctx):
  return ctx.num_replicas_in_sync
distributed_values = (
     strategy.experimental_distribute_values_from_function(
       value_fn))
local_result = strategy.experimental_local_results(distributed_values)
local_result
(2, 2)
  1. Размещение значений на устройствах и распределение:
strategy = tf.distribute.TPUStrategy()
worker_devices = strategy.extended.worker_devices
multiple_values = []
for i in range(strategy.num_replicas_in_sync):
  with tf.device(worker_devices[i]):
    multiple_values.append(tf.constant(1.0))

def value_fn(ctx):
  return multiple_values[ctx.replica_id_in_sync_group]

distributed_values = strategy.
  experimental_distribute_values_from_function(
  value_fn)

experimental_local_results

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

experimental_local_results(
    value
)

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

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

gather

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

gather(
    value, axis
)

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

Учитывая tf.distribute.DistributedValues или tf.Tensor-подобный объект value, этот API собирает и конкатенирует value по репликам вдоль axis-й размерности. Результат копируется на "текущее" устройство

  • что обычно является процессором рабочего процесса, на котором выполняется программа. Для tf.distribute.TPUStrategy это первый хост TPU. Для многоклиентской MultiWorkerMirroredStrategy, это процессор каждого рабочего процесса.

Этот API может быть вызван только в контексте между репликами. Для аналога в контексте реплики см. tf.distribute.ReplicaContext.all_gather.

Примечание: Для всех стратегий, кроме tf.distribute.TPUStrategy, вход value на разных репликах должен иметь одинаковый ранг, а их формы должны быть одинаковыми во всех измерениях, кроме axis-й размерности. Другими словами, их формы не могут отличаться в измерении d, где d не равно аргументу axis. Например, учитывая tf.distribute.DistributedValues с компонентами тензоров формы (1, 2, 3) и (1, 3, 3) на двух репликах, вы можете вызвать gather(..., axis=1, ...), но не gather(..., axis=0, ...) или gather(..., axis=2, ...). Однако для tf.distribute.TPUStrategy.gather все тензоры должны иметь ровно одинаковый ранг и одинаковую форму.
Примечание: Учитывая tf.distribute.DistributedValues value, компоненты тензоров должны иметь ненулевой ранг. В противном случае рассмотрите использование tf.expand_dims перед их сбором.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
# A DistributedValues with component tensor of shape (2, 1) on each replica
distributed_values = strategy.experimental_distribute_values_from_function(lambda _: tf.identity(tf.constant([[1], [2]])))
@tf.function
def run():
  return strategy.gather(distributed_values, axis=0)
run()
<tf.Tensor: shape=(4, 1), dtype=int32, numpy=
array([[1],
       [2],
       [1],
       [2]], dtype=int32)>

Рассмотрим следующий пример для большего количества комбинаций:

strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1", "GPU:2", "GPU:3"])
single_tensor = tf.reshape(tf.range(6), shape=(1,2,3))
distributed_values = strategy.experimental_distribute_values_from_function(lambda _: tf.identity(single_tensor))
@tf.function
def run(axis):
  return strategy.gather(distributed_values, axis=axis)
axis=0
run(axis)
<tf.Tensor: shape=(4, 2, 3), dtype=int32, numpy=
array([[[0, 1, 2],
        [3, 4, 5]],
       [[0, 1, 2],
        [3, 4, 5]],
       [[0, 1, 2],
        [3, 4, 5]],
       [[0, 1, 2],
        [3, 4, 5]]], dtype=int32)>
axis=1
run(axis)
<tf.Tensor: shape=(1, 8, 3), dtype=int32, numpy=
array([[[0, 1, 2],
        [3, 4, 5],
        [0, 1, 2],
        [3, 4, 5],
        [0, 1, 2],
        [3, 4, 5],
        [0, 1, 2],
        [3, 4, 5]]], dtype=int32)>
axis=2
run(axis)
<tf.Tensor: shape=(1, 2, 12), dtype=int32, numpy=
array([[[0, 1, 2, 0, 1, 2, 0, 1, 2, 0, 1, 2],
        [3, 4, 5, 3, 4, 5, 3, 4, 5, 3, 4, 5]]], dtype=int32)>
Аргументы
value экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который необходимо объединить в один тензор. Также может быть обычным тензором, когда используется с tf.distribute.OneDeviceStrategy или по умолчанию. Тензоры, составляющие DistributedValues, могут быть только плотными тензорами с ненулевым рангом, а НЕ tf.IndexedSlices.
axis 0-мерный тензор типа int32. Размеры, по которым выполнять сбор. Должно находиться в диапазоне [0, rank(value)).
Возвращаемое значение
Объект Tensor, представляющий собой конкатенацию value по всем репликам вдоль axis измерения.

reduce

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

reduce(
    reduce_op, value, axis
)

Свести 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: Более продвинутые версии Strategy.reduce, позволяющие настраивать место назначения результата. Также вызываются в контексте между репликами.

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

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

Например, если у вас есть глобальный размер пакета 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 должны быть либо Python-значениями вложенной структуры тензоров, например, список тензоров, в этом случае args и kwargs будут переданы в fn, вызванный на каждой реплике. Или args или kwargs могут быть tf.distribute.DistributedValues, содержащими тензоры или составные тензоры, например, tf.compat.v1.TensorInfo.CompositeTensor, в этом случае каждый fn вызов получит компонент tf.distribute.DistributedValues, соответствующий его реплике.

Ключевой момент: В зависимости от реализации tf.distribute.Strategy и того, включена ли жадная реализация, fn может быть вызван один или несколько раз. Если fn анотирован с tf.function или tf.distribute.Strategy.run вызывается внутри tf.function (жадная реализация выключена внутри tf.function по умолчанию), fn вызывается один раз на реплику для генерации графа Tensorflow, который затем будет повторно использован для выполнения с новыми входными данными. В противном случае, если жадная реализация включена, 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>
}
  1. Входной 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>
  1. Используйте 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. Его элементы могут быть Python-значениями, тензорами или tf.distribute.DistributedValues.
kwargs необязательные именованные аргументы для fn. Его элементы могут быть Python-значениями, тензорами или 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() теперь вернёт эту стратегию. За пределами этого контекста возвращается стратегия по умолчанию, ничего не выполняющая.
  • Вхождение в контекст также вводит «межрепличный контекст». См. tf.distribute.StrategyExtended для объяснения межрепличных и репличных контекстов.
  • Создание переменных внутри scope перехватывается стратегией. Каждая стратегия определяет, как она хочет повлиять на создание переменной. Стратегии синхронизации, такие как MirroredStrategy, TPUStrategy и MultiWorkerMiroredStrategy, создают переменные, реплицированные на каждой реплике, тогда как ParameterServerStrategy создаёт переменные на серверах параметров. Это делается с помощью пользовательского tf.variable_creator_scope.
  • В некоторых стратегиях также может вводиться контекст устройства по умолчанию: в MultiWorkerMiroredStrategy, на каждом рабочем узле вводится контекст устройства по умолчанию «/CPU:0».
Примечание: Вхождение в контекст не автоматически распределяет вычисление, за исключением случаев с высокоуровневыми фреймворками обучения, такими как Keras model.fit. Если вы не используете model.fit, вам необходимо использовать API strategy.run для явного распределения этого вычисления. См. пример в учебнике по пользовательским циклам обучения.

Что должно находиться в контексте, а что — вне его?

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

  • Любое создание переменных, которые должны быть распределёнными переменными, должно происходить в strategy.scope. Это можно сделать, либо напрямую поместив его в контекст, либо используя другой API, например, strategy.run или model.fit, чтобы он входил в него за вас. Любая переменная, созданная вне контекста, не будет распределена и может повлиять на производительность. Типичные вещи, создающие переменные в TF: модели, оптимизаторы, метрики. Они всегда должны создаваться внутри контекста. Другим источником создания переменных может быть восстановление контрольной точки — когда переменные создаются лениво. Обратите внимание, что любая переменная, созданная внутри стратегии, сохраняет информацию о стратегии. Поэтому чтение и запись этих переменных за пределами 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 иногда может потребоваться находиться в контексте, если оно создаёт переменные.
Возвращает
Менеджер контекста.

© 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/MultiWorkerMirroredStrategy

Spec-Zone.ru

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