Spec-Zone.ru › TensorFlow

tf.distribute.MultiWorkerMirroredStrategy

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

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

tf.distribute.MultiWorkerMirroredStrategy(
    cluster_resolver=None, communication_options=None
)

Использование в ноутбуках

Использование в руководстве Использование в учебниках
  • Распределенное обучение с TensorFlow
  • Миграция распределенного обучения на нескольких рабочих узлах с CPU/GPU
  • Распределенное обучение с Keras на нескольких рабочих узлах
  • Пользовательский цикл обучения с Keras и MultiWorkerMirroredStrategy

Эта стратегия реализует синхронное распределенное обучение на нескольких рабочих узлах, каждый из которых может иметь несколько GPU. Подобно 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() определяет количество GPU, используемых стратегией. Если оно равно нулю, стратегия использует ЦП. Все рабочие узлы должны использовать одинаковое количество устройств, иначе поведение не определено.

Эта стратегия не предназначена для 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 не возвращает правильное количество ускорителей. Стратегия использует все доступные GPU, если 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 не разбивает на пакеты или не разделяет экземпляр набора данных, возвращаемый функцией ввода. 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, чтобы получить тип элементов, возвращаемых итератором. Это может использоваться для установки свойства 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 в стиле Python. x — это tf.distribute.DistributedValues, содержащий данные для всех реплик, и каждая реплика получает данные нового размера пакета. tf.distribute.Strategy.run позаботится о подаче правильных данных по каждой реплике в x в правильные replica_fn, выполняемые на каждой реплике.

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

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

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

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

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

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

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

experimental_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>)
        
  2. Распределение значений в массиве на основе 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)
        
  3. Указание значений с помощью 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)
        
  4. Размещение значений на устройствах и распределение:

    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, ранг(значение)).
Возвращает
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 GPU:

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Аргументы
reduce_op значение tf.distribute.ReduceOp, определяющее, как следует объединять значения. Разрешает использование строкового представления перечисления, такого как «СУММА», «СРЕДНЕЕ».
value экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, для объединения в один тензор. Он также может быть обычным тензором при использовании с OneDeviceStrategy или стратегией по умолчанию.
axis определяет размерность для сведения вдоль тензора каждой реплики. Обычно следует устанавливать в размерность пакета или None, чтобы производить сведение только по репликам (например, если тензор не имеет размерности пакета).
Возвращает
Tensor.

run

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

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

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

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

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

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

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

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

  1. Постоянный тензорный ввод.

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

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

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

scope

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

scope()

Менеджер контекста для установки стратегии по умолчанию и распределения переменных.

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

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

Что происходит при входе в Strategy.scope?

  • strategy устанавливается в глобальный контекст в качестве текущей стратегии. Внутри этого области tf.distribute.get_strategy() теперь вернёт эту стратегию. Вне этой области, он возвращает стратегию по умолчанию (без действий).
  • Вход в область также означает вход в "межрепличный контекст". См. 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/api_docs/python/tf/distribute/MultiWorkerMirroredStrategy

Spec-Zone.ru

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