Spec-Zone.ru › TensorFlow 2.9

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)

Подробный учебник см. в Multi-worker training with Keras.

Сохранение

Вам необходимо сохранять и создавать контрольные точки на всех рабочих узлах, а не только на одном. Это связано с тем, что переменные, для которых synchronization=ON_READ, вызывают агрегацию при сохранении. Рекомендуется сохранять в разных директориях на каждом рабочем узле, чтобы избежать гонок. Каждый рабочий узел сохраняет то же самое. Примеры см. в учебнике Multi-worker training with 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.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). В тех случаях, когда набор данных бесконечный, это разбиение можно осуществить, создав копии наборов данных, отличающиеся только случайным начальным значением.

Функция-обработчик должна принимать экземпляр 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.

По умолчанию этот метод добавляет преобразование предварительной выборки в конце экземпляра пользовательского 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_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. Распределение значений в массиве на основе id реплики:
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(), or a variable created inscope`.
Возвращаемое значение
Кортеж значений, содержащихся в value, где i-й элемент соответствует i-й реплике. Если value представляет одно значение, это возвращает (value,).

gather

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

gather(
    value, axis
)

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

Учитывая tf.distribute.DistributedValues или объект наподобие tf.Tensor value, эта API собирает и конкатенирует value по репликам вдоль axis-й размерности. Результат копируется на "текущее" устройство, которое обычно является процессором рабочего узла, на котором выполняется программа. Для tf.distribute.TPUStrategy это первый хост TPU. Для многоклиентского tf.distribute.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(значение)).
Возвращаемое значение
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, так как они позволяют настроить место назначения результата. Они также вызываются в контексте межрепликационной связи.

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

Учитывая значение, возвращаемое 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, определяющее, как нужно комбинировать значения. Разрешает использовать строковое представление перечисления, например, "СУММА", "СРЕДНЕЕ".
value экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, для объединения в один тензор. Он также может быть обычным тензором при использовании с OneDeviceStrategy или по умолчанию стратегией.
axis определяет измерение, по которому следует выполнить уменьшение в тензоре каждой реплики. Обычно должно быть установлено на измерение пакета или None для выполнения сведения только по репликам (например, если тензор не имеет измерения пакета).
Возвращаемое значение
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 могут быть вложенной структурой тензоров, например, списком тензоров, в этом случае args и kwargs будут переданы в fn , вызываемом на каждой реплике. Или args или kwargs могут быть tf.distribute.DistributedValues, содержащие тензоры или составные тензоры, т.е. tf.compat.v1.TensorInfo.CompositeTensor, в этом случае каждый вызов fn получит компонент tf.distribute.DistributedValues, соответствующий его реплике. Обратите внимание, что произвольные значения Python, которые не являются указанных типов, не поддерживаются.

Важно: В зависимости от реализации 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. Его элемент может быть тензором, вложенной структурой тензоров или tf.distribute.DistributedValues.
kwargs Необязательные ключевые аргументы для fn. Его элемент может быть тензором, вложенной структурой тензоров или tf.distribute.DistributedValues.
options Необязательный экземпляр tf.distribute.RunOptions, определяющий параметры для выполнения fn.
Возвращаемое значение
Объединённое возвращаемое значение fn по всем репликам. Структура возвращаемого значения такая же, как и возвращаемое значение из fn. Каждый элемент структуры может быть tf.distribute.DistributedValues, объектами 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 или 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 иногда может потребоваться внутри области видимости, если оно создаёт переменные.
Возвращаемое значение
Менеджер контекста.

© 2022 The TensorFlow Authors. All rights reserved.
Licensed under the Creative Commons Attribution License 4.0.
Code samples licensed under the Apache 2.0 License.
https://www.tensorflow.org/versions/r2.9/api_docs/python/tf/distribute/MultiWorkerMirroredStrategy

Spec-Zone.ru

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