Spec-Zone.ru › TensorFlow 2.9

tf.distribute.experimental.MultiWorkerMirroredStrategy

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

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

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

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

Эта стратегия реализует синхронное распределённое обучение на нескольких рабочих узлах, каждый из которых может иметь несколько графических процессоров. Подобно tf.distribute.MirroredStrategy, она дублирует все переменные и вычисления на каждом локальном устройстве. Разница заключается в том, что она использует распределённую коллективную реализацию (например, all-reduce), чтобы несколько рабочих узлов могли работать вместе.

Необходимо запустить вашу программу на каждом рабочем узле и правильно настроить cluster_resolver. Например, если вы используете tf.distribute.cluster_resolver.TFConfigClusterResolver, каждый рабочий узел должен иметь соответствующие task_type и task_id, установленные в переменной среды TF_CONFIG. Пример TF_CONFIG для рабочего узла-0 двух узловой кластера:

TF_CONFIG = '{"cluster": {"worker": ["localhost:12345", "localhost:23456"]}, "task": {"type": "worker", "index": 0} }'

Ваша программа работает на каждом рабочем узле без изменений. Обратите внимание, что для коллективных операций требуется участие каждого рабочего узла. Все tf.distribute и не tf.distribute API могут использовать коллективные операции внутри, например, сохранение и восстановление контрольных точек, так как чтение tf.Variable с tf.VariableSynchronization.ON_READ выполняет all-reduce значения. Поэтому рекомендуется запускать на каждом рабочем узле точно ту же программу. Распределение задач на основе task_type или task_id рабочего узла чревато ошибками.

cluster_resolver.num_accelerators() определяет количество графических процессоров, используемых стратегией. Если оно равно нулю, стратегия использует ЦП. Все рабочие узлы должны использовать одинаковое количество устройств, в противном случае поведение не определено.

Эта стратегия не предназначена для TPU. Используйте tf.distribute.TPUStrategy вместо неё.

После настройки TF_CONFIG использование этой стратегии аналогично использованию tf.distribute.MirroredStrategy и tf.distribute.TPUStrategy.

strategy = tf.distribute.MultiWorkerMirroredStrategy()

with strategy.scope():
  model = tf.keras.Sequential([
    tf.keras.layers.Dense(2, input_shape=(5,)),
  ])
  optimizer = tf.keras.optimizers.SGD(learning_rate=0.1)

def dataset_fn(ctx):
  x = np.random.random((2, 5)).astype(np.float32)
  y = np.random.randint(2, size=(2, 1))
  dataset = tf.data.Dataset.from_tensor_slices((x, y))
  return dataset.repeat().batch(1, drop_remainder=True)
dist_dataset = strategy.distribute_datasets_from_function(dataset_fn)

model.compile()
model.fit(dist_dataset)

Вы также можете написать свою собственную цикл обучения:

@tf.function
def train_step(iterator):

  def step_fn(inputs):
    features, labels = inputs
    with tf.GradientTape() as tape:
      logits = model(features, training=True)
      loss = tf.keras.losses.sparse_categorical_crossentropy(
          labels, logits)

    grads = tape.gradient(loss, model.trainable_variables)
    optimizer.apply_gradients(zip(grads, model.trainable_variables))

  strategy.run(step_fn, args=(next(iterator),))

for _ in range(NUM_STEP):
  train_step(iterator)

См. Обучение на нескольких рабочих узлах с помощью Keras для подробного руководства.

Сохранение

Вам необходимо сохранять и создавать контрольные точки на всех рабочих узлах, а не только на одном. Это связано с тем, что переменные, у которых synchronization=ON_READ, вызывают агрегацию во время сохранения. Рекомендуется сохранять в разные пути на каждом рабочем узле, чтобы избежать гонок. Каждый рабочий узел сохраняет одно и то же. Смотрите учебное пособие Обучение на нескольких рабочих узлах с помощью Keras для примеров.

Известные проблемы

  • tf.distribute.cluster_resolver.TFConfigClusterResolver не возвращает правильное количество ускорителей. Стратегия использует все доступные графические процессоры, если cluster_resolver является tf.distribute.cluster_resolver.TFConfigClusterResolver или None.
  • В режиме eager стратегия должна быть создана перед вызовом любого другого API TensorFlow.
Аргументы
communication необязательная tf.distribute.experimental.CommunicationImplementation. Это подсказка о предпочтительной реализации коллективной коммуникации. Возможные значения включают AUTO, RING, и NCCL.
cluster_resolver необязательный tf.distribute.cluster_resolver.ClusterResolver. Если None, используется tf.distribute.cluster_resolver.TFConfigClusterResolver.
Атрибуты
cluster_resolver Возвращает решатель кластера, связанный с этой стратегией.

Как стратегия для нескольких рабочих узлов, tf.distribute.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 будет вызываться на процессоре ЦП каждого рабочего узла, и каждый генерирует набор данных, в котором каждая реплика на этом рабочем узле будет извлекать по одному пакету входных данных (т. е. если у рабочего узла две реплики, две порции будут извлечены из Dataset на каждом шаге).

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

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

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

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

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

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

experimental_distribute_dataset

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

experimental_distribute_dataset(
    dataset, options=None
)

Создаёт tf.distribute.DistributedDataset из tf.data.Dataset.

Возвращаемый tf.distribute.DistributedDataset можно перебирать так же, как и обычные наборы данных. ПРИМЕЧАНИЕ: пользователь не может добавить больше преобразований к tf.distribute.DistributedDataset. Вы можете только создать итератор или изучить tf.TypeSpec генерируемых им данных. См. документацию API tf.distribute.DistributedDataset для получения дополнительной информации.

Следующий пример:

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

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

В приведенном выше фрагменте кода dataset пакетируется global_batch_size, а вызов experimental_distribute_dataset на нём повторно пакетирует dataset до нового размера пакета, равного глобальному размеру пакета, делённому на количество реплик в синхронизации. Мы перебираем его с помощью питоновского цикла for. 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. Распределение значений в массиве на основе идентификатора реплики:
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-й размерности. Результат копируется на "текущее" устройство, которое, как правило, является процессором 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 тензор int32 размерности 0. Размеры, по которым выполняется сбор. Должно быть в диапазоне [0, rank(значение)).
Возвращаемое значение
Tensor, который является конкатенацией value по всем репликам вдоль axis измерения.

reduce

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

reduce(
    reduce_op, value, axis
)

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

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

per_replica_result = strategy.run(step_fn)
total = strategy.reduce("SUM", per_replica_result, axis=None)
total
<tf.Tensor: shape=(), dtype=int32, numpy=1>

Чтобы увидеть, как это будет выглядеть с несколькими репликами, рассмотрите тот же пример с MirroredStrategy с 2 графическими процессорами:

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Аргументы
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 (режим 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>
}
  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, или 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/experimental/MultiWorkerMirroredStrategy

Spec-Zone.ru

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