Spec-Zone.ru › TensorFlow

tf.distribute.experimental.MultiWorkerMirroredStrategy

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

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

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

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

Используется в учебниках
  • Обучение на нескольких рабочих узлах с помощью Estimator

Эта стратегия реализует синхронное распределение обучения на нескольких рабочих узлах, каждый из которых может иметь несколько 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.
Аргументы
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.MultiWorkerMirroredStrategy предоставляет связанный tf.distribute.cluster_resolver.ClusterResolver. Если пользователь предоставляет его в __init__, возвращается этот экземпляр; если пользователь не предоставляет, предоставляется по умолчанию TFConfigClusterResolver.

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

Методы

distribute_datasets_from_function

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

distribute_datasets_from_function(
    dataset_fn, options=None
)

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

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

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

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

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

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

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

END_OF_DOCUMENT_MARKER ```
Args
dataset_fn Функция, принимающая экземпляр tf.distribute.InputContext и возвращающая tf.data.Dataset.
options tf.distribute.InputOptions, используемая для управления параметрами распределения набора данных.
Returns
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 генерируемых им данных. Подробнее см. в документации tf.distribute.DistributedDataset.

Вот пример:

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

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

В приведённом фрагменте кода dataset группируется по global_batch_size, а вызов experimental_distribute_dataset на нём перегруппировывает dataset в новый размер пакета, равный глобальному размеру пакета, делённому на число реплик в синхронизации. Мы итерируем через него с помощью цикла Python 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.

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

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>)
        
  2. Распределение значений в массиве на основе идентификатора реплики:

    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), каждый рабочий узел будет своим клиентом, и эта функция будет возвращать только значения, вычисленные на этом рабочем узле.
Args
value Значение, возвращённое experimental_run(), run(), or a variable created inscope`.
Returns
Кортеж значений, содержащихся в value, где i-й элемент соответствует i-й реплике. Если value представляет одно значение, возвращается (value,).

gather

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

gather(
    value, axis
)

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

При заданном объекте tf.distribute.DistributedValues или tf.Tensor подобного типа value, этот API собирает и конкатенирует value по всем репликам вдоль axis-й размерности. Результат копируется на «текущее» устройство, которое обычно является процессором (CPU) рабочего узла, на котором выполняется программа. Для tf.distribute.TPUStrategy это первый хост TPU. Для многоклиентского tf.distribute.MultiWorkerMirroredStrategy это процессор (CPU) каждого рабочего узла.

Этот 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 перед выводом.

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

Существует ряд различных 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.

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

run

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

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

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

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

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

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

Важно: В зависимости от реализации tf.distribute.Strategy и от того, включена ли жадная (eager) обработка, fn может быть вызвана один или несколько раз. Если fn аннотирована с tf.function или tf.distribute.Strategy.run вызывается внутри tf.function (жадная обработка по умолчанию отключена внутри 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>
        }
        
  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/experimental/MultiWorkerMirroredStrategy

Spec-Zone.ru

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