Spec-Zone.ru › TensorFlow 2.4

tf.distribute.experimental.MultiWorkerMirroredStrategy

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

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

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

tf.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.
Аргументы
communication необязательная tf.distribute.experimental.CommunicationImplementation. Это подсказка о предпочтительной реализации коллективного обмена данными. Возможные значения включают AUTO, RING, и NCCL.
cluster_resolver необязательная tf.distribute.cluster_resolver.ClusterResolver. Если None, используется tf.distribute.cluster_resolver.TFConfigClusterResolver.
Атрибуты
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 не разбивает и не распределяет экземпляр tf.data.Dataset, возвращаемый функцией-входом. dataset_fn будет вызываться на процессоре CPU каждого узла, и каждый создаёт набор данных, где каждая реплика на этом узле будет извлекать один пакет входов (т.е. если у узла две реплики, с набора данных 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, чтобы посмотреть пример.

Ключевая точка: Возвращаемый dataset_fn tf.data.Dataset должен иметь размер пакета на реплику, в отличие от 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 в стиле 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.

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

Args
dataset tf.data.Dataset, который будет разделен по всем репликам в соответствии с вышеуказанными правилами.
options tf.distribute.InputOptions для управления параметрами распределения этого набора данных.
Returns
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 или другие методы, принимающие распределённые значения, когда не используются наборы данных.

Args
value_fn Функция для выполнения для генерации значений. Она вызывается для каждой реплики с tf.distribute.ValueContext в качестве единственного аргумента. Она должна возвращать тензор или тип, который может быть преобразован в тензор.
Returns
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, каждый рабочий узел будет своим клиентом, и эта функция вернёт только значения, вычисленные на этом рабочем узле.
Args
value Значение, возвращаемое experimental_run(), run(), extended.call_for_each_replica(), или переменная, созданная в scope
Returns
Кортеж значений, содержащихся в 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)>
Args
value экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который нужно объединить в один тензор. Он также может быть обычным тензором, если используется с tf.distribute.OneDeviceStrategy или по умолчанию стратегией. Тензоры, составляющие DistributedValues, могут быть только плотными тензорами с ненулевым рангом, а НЕ tf.IndexedSlices.
axis 0-мерный тензор типа int32. Размерность, по которой собирать. Должно находиться в диапазоне [0, rank(value)).
Returns
A 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: Эти API представляют собой более продвинутые версии Strategy.reduce, так как они позволяют настраивать место назначения результата. Они также вызываются в контексте между репликами.

Что должно быть значением axis?

Учитывая значение на каждую реплику, возвращаемое 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.

Args
reduce_op значение tf.distribute.ReduceOp, указывающее, как должны комбинироваться значения. Разрешает использовать строковое представление перечисления, например, "SUM", "MEAN".
value экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который нужно объединить в один тензор. Он также может быть обычным тензором, если используется с OneDeviceStrategy или стратегией по умолчанию.
axis указывает размерность для сведения вдоль тензора каждой реплики. Обычно следует устанавливать в размерность пакета или None для сведения только по репликам (например, если у тензора нет размерности пакета).
Returns
A 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>
}
Args
fn Функция для выполнения на каждой реплике.
args Необязательные позиционные аргументы для fn. Его элемент может быть значением Python, тензором или tf.distribute.DistributedValues.
kwargs Необязательные ключевые аргументы для fn. Его элемент может быть значением Python, тензором или tf.distribute.DistributedValues.
options Необязательный экземпляр tf.distribute.RunOptions, определяющий параметры выполнения fn.
Returns
Объединённое возвращаемое значение 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/experimental/MultiWorkerMirroredStrategy

Spec-Zone.ru

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