Spec-Zone.ru › TensorFlow 2.9

tf.distribute.MirroredStrategy

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

Синхронное обучение на нескольких репликах на одной машине.

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

tf.distribute.MirroredStrategy(
    devices=None, cross_device_ops=None
)

Эта стратегия обычно используется для обучения на одной машине с несколькими графическими процессорами. Для TPUs используйте tf.distribute.TPUStrategy. Чтобы использовать MirroredStrategy с несколькими рабочими узлами, обратитесь к tf.distribute.experimental.MultiWorkerMirroredStrategy.

Например, переменная, созданная в рамках MirroredStrategy, является MirroredVariable. Если в конструкторе стратегии не указаны устройства, она будет использовать все доступные графические процессоры. Если графических процессоров не найдено, она будет использовать доступные центральные процессоры. Обратите внимание, что TensorFlow обрабатывает все центральные процессоры на машине как одно устройство и использует потоки внутри для параллелизма.

strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
with strategy.scope():
  x = tf.Variable(1.)
x
MirroredVariable:{
  0: <tf.Variable ... shape=() dtype=float32, numpy=1.0>,
  1: <tf.Variable ... shape=() dtype=float32, numpy=1.0>
}

При использовании стратегий распределения все создание переменных должно выполняться в области действия стратегии. Это позволит продублировать переменные на всех репликах и поддерживать их синхронизацию с помощью алгоритма all-reduce.

Переменные, созданные внутри MirroredStrategy, который обернут в tf.function, по-прежнему MirroredVariables.

x = []
@tf.function  # Wrap the function with tf.function.
def create_variable():
  if not x:
    x.append(tf.Variable(1.))
  return x[0]
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
with strategy.scope():
  _ = create_variable()
  print(x[0])
MirroredVariable:{
  0: <tf.Variable ... shape=() dtype=float32, numpy=1.0>,
  1: <tf.Variable ... shape=() dtype=float32, numpy=1.0>
}

experimental_distribute_dataset можно использовать для распределения набора данных по репликам при написании собственного цикла обучения. Если вы используете .fit и .compile методы, доступные в tf.keras, то tf.keras будет обрабатывать распределение за вас.

Например:

my_strategy = tf.distribute.MirroredStrategy()
with my_strategy.scope():
  @tf.function
  def distribute_train_epoch(dataset):
    def replica_fn(input):
      # process input and return result
      return result

    total_result = 0
    for x in dataset:
      per_replica_result = my_strategy.run(replica_fn, args=(x,))
      total_result += my_strategy.reduce(tf.distribute.ReduceOp.SUM,
                                         per_replica_result, axis=None)
    return total_result

  dist_dataset = my_strategy.experimental_distribute_dataset(dataset)
  for _ in range(EPOCHS):
    train_result = distribute_train_epoch(dist_dataset)
Аргументы
devices список строк устройств, таких как ['/gpu:0', '/gpu:1']. Если None, используются все доступные графические процессоры. Если графические процессоры не найдены, используется центральный процессор.
cross_device_ops необязательно, потомок CrossDeviceOps. Если это не задано, по умолчанию будет использоваться NcclAllReduce(). Это нужно настраивать, если NCCL недоступен или если доступна специальная реализация, которая использует конкретное оборудование.
Атрибуты
cluster_resolver Возвращает решатель кластера, связанный с этой стратегией.

В общем случае, при использовании многоузловой tf.distribute стратегии, такой как tf.distribute.experimental.MultiWorkerMirroredStrategy или tf.distribute.TPUStrategy(), существует tf.distribute.cluster_resolver.ClusterResolver, связанный с используемой стратегией, и такой экземпляр возвращается этим свойством.

Стратегии, которые намерены иметь связанный tf.distribute.cluster_resolver.ClusterResolver, должны установить соответствующий атрибут или переопределить это свойство; в противном случае по умолчанию возвращается None. Эти стратегии также должны предоставить информацию о том, что возвращается этим свойством.

Стратегии с одним рабочим узлом обычно не имеют tf.distribute.cluster_resolver.ClusterResolver, и в этих случаях это свойство вернет None.

tf.distribute.cluster_resolver.ClusterResolver может быть полезен, когда пользователю нужно получить доступ к информации, такой как описание кластера, тип задачи или идентификатор задачи. Например,

os.environ['TF_CONFIG'] = json.dumps({
  'cluster': {
      'worker': ["localhost:12345", "localhost:23456"],
      'ps': ["localhost:34567"]
  },
  'task': {'type': 'worker', 'index': 0}
})

# This implicitly uses TF_CONFIG for the cluster and current task info.
strategy = tf.distribute.experimental.MultiWorkerMirroredStrategy()

...

if strategy.cluster_resolver.task_type == 'worker':
  # Perform something that's only applicable on workers. Since we set this
  # as a worker above, this block will run on this particular instance.
elif strategy.cluster_resolver.task_type == 'ps':
  # Perform something that's only applicable on parameter servers. Since we
  # set this as a worker above, this block will not run on this particular
  # instance.

Для получения дополнительной информации, пожалуйста, обратитесь к документации API tf.distribute.cluster_resolver.ClusterResolver.

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 для просмотра примера.

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

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

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

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

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

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

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

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

experimental_distribute_values_from_function

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

experimental_distribute_values_from_function(
    value_fn
)

Генерирует tf.distribute.DistributedValues из value_fn.

Эта функция предназначена для генерации tf.distribute.DistributedValues для передачи в run, reduce, или другие методы, принимающие распределённые значения, когда не используются наборы данных.

Аргументы
value_fn Функция для запуска генерации значений. Она вызывается для каждой реплики с tf.distribute.ValueContext в качестве единственного аргумента. Она должна возвращать тензор или тип, который можно преобразовать в тензор.
Возвращаемые значения
tf.distribute.DistributedValues, содержащий значение для каждой реплики.

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

  1. Возвратить постоянное значение для каждой реплики:
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
def value_fn(ctx):
  return tf.constant(1.)
distributed_values = (
     strategy.experimental_distribute_values_from_function(
       value_fn))
local_result = strategy.experimental_local_results(distributed_values)
local_result
(<tf.Tensor: shape=(), dtype=float32, numpy=1.0>,
 <tf.Tensor: shape=(), dtype=float32, numpy=1.0>)
  1. Распределить значения в массиве на основе id реплики:
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
array_value = np.array([3., 2., 1.])
def value_fn(ctx):
  return array_value[ctx.replica_id_in_sync_group]
distributed_values = (
     strategy.experimental_distribute_values_from_function(
       value_fn))
local_result = strategy.experimental_local_results(distributed_values)
local_result
(3.0, 2.0)
  1. Указать значения с помощью num_replicas_in_sync:
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
def value_fn(ctx):
  return ctx.num_replicas_in_sync
distributed_values = (
     strategy.experimental_distribute_values_from_function(
       value_fn))
local_result = strategy.experimental_local_results(distributed_values)
local_result
(2, 2)
  1. Разместить значения на устройствах и распределить:
strategy = tf.distribute.TPUStrategy()
worker_devices = strategy.extended.worker_devices
multiple_values = []
for i in range(strategy.num_replicas_in_sync):
  with tf.device(worker_devices[i]):
    multiple_values.append(tf.constant(1.0))

def value_fn(ctx):
  return multiple_values[ctx.replica_id_in_sync_group]

distributed_values = strategy.
  experimental_distribute_values_from_function(
  value_fn)

experimental_local_results

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

experimental_local_results(
    value
)

Возвращает список всех локальных значений каждой реплики, содержащихся в value.

Примечание: Это возвращает только значения на рабочем процессе, инициированном этим клиентом. При использовании tf.distribute.Strategy, такого как tf.distribute.experimental.MultiWorkerMirroredStrategy, каждый рабочий процесс будет своим клиентом, и эта функция вернёт только значения, вычисленные на этом рабочем процессе.
Аргументы
value Значение, возвращённое experimental_run(), run(), or a variable created inscope`.
Возвращаемые значения
Кортеж значений, содержащихся в value , где i-й элемент соответствует i-й реплике. Если value представляет единственное значение, возвращается (value,).

gather

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

gather(
    value, axis
)

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

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

Этот API можно вызывать только в контексте между репликами. Для аналога в контексте реплики см. tf.distribute.ReplicaContext.all_gather.

Примечание: Для всех стратегий, кроме tf.distribute.TPUStrategy, входной value на разных репликах должен иметь одинаковый ранг, а их формы должны быть одинаковыми во всех размерностях, кроме axis-й размерности. Другими словами, их формы не могут отличаться в размерности d , где d не равно аргументу axis. Например, при tf.distribute.DistributedValues с тензорами компонентов формы (1, 2, 3) и (1, 3, 3) на двух репликах вы можете вызвать gather(..., axis=1, ...), но не gather(..., axis=0, ...) или gather(..., axis=2, ...). Однако для tf.distribute.TPUStrategy.gather все тензоры должны иметь точно такой же ранг и такую же форму.
Примечание: Учитывая tf.distribute.DistributedValues value, его компоненты тензоров должны иметь ненулевой ранг. В противном случае рассмотрите использование tf.expand_dims перед их сборкой.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
# A DistributedValues with component tensor of shape (2, 1) on each replica
distributed_values = strategy.experimental_distribute_values_from_function(lambda _: tf.identity(tf.constant([[1], [2]])))
@tf.function
def run():
  return strategy.gather(distributed_values, axis=0)
run()
<tf.Tensor: shape=(4, 1), dtype=int32, numpy=
array([[1],
       [2],
       [1],
       [2]], dtype=int32)>

Рассмотрим следующий пример для других комбинаций:

strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1", "GPU:2", "GPU:3"])
single_tensor = tf.reshape(tf.range(6), shape=(1,2,3))
distributed_values = strategy.experimental_distribute_values_from_function(lambda _: tf.identity(single_tensor))
@tf.function
def run(axis):
  return strategy.gather(distributed_values, axis=axis)
axis=0
run(axis)
<tf.Tensor: shape=(4, 2, 3), dtype=int32, numpy=
array([[[0, 1, 2],
        [3, 4, 5]],
       [[0, 1, 2],
        [3, 4, 5]],
       [[0, 1, 2],
        [3, 4, 5]],
       [[0, 1, 2],
        [3, 4, 5]]], dtype=int32)>
axis=1
run(axis)
<tf.Tensor: shape=(1, 8, 3), dtype=int32, numpy=
array([[[0, 1, 2],
        [3, 4, 5],
        [0, 1, 2],
        [3, 4, 5],
        [0, 1, 2],
        [3, 4, 5],
        [0, 1, 2],
        [3, 4, 5]]], dtype=int32)>
axis=2
run(axis)
<tf.Tensor: shape=(1, 2, 12), dtype=int32, numpy=
array([[[0, 1, 2, 0, 1, 2, 0, 1, 2, 0, 1, 2],
        [3, 4, 5, 3, 4, 5, 3, 4, 5, 3, 4, 5]]], dtype=int32)>
Аргументы
value экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, для объединения в один тензор. Он также может быть обычным тензором при использовании с tf.distribute.OneDeviceStrategy или по умолчанию. Тензоры, составляющие DistributedValues, могут быть только плотными тензорами с ненулевым рангом, а НЕ tf.IndexedSlices.
axis 0-мерный тензор int32. Размерность для сбора. Должно быть в диапазоне [0, rank(value)).
Возвращаемые значения
Tensor, который представляет собой конкатенацию value по репликам вдоль axis размерности.

reduce

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

reduce(
    reduce_op, value, axis
)

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Аргументы
reduce_op значение tf.distribute.ReduceOp, определяющее, как должны комбинироваться значения. Допускается использование строкового представления перечисления, такого как «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 для явного распределения вычисления. См. пример в руководстве по обучению с пользовательским циклом custom training loop.

Что должно находиться в области действия, а что вне её?

Существует ряд требований к тому, что должно происходить в рамках области видимости. Однако в тех местах, где у нас есть информация о том, какая стратегия используется, мы часто входим в область видимости для пользователя, чтобы он не должен был делать это явно (т.е. вызов этих функций внутри или вне области видимости допустим).

  • Любой код, создающий переменные, которые должны быть распределенными, должен вызываться в 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/MirroredStrategy

Spec-Zone.ru

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