Spec-Zone.ru › TensorFlow 2.9

tf.compat.v1.distribute.experimental.MultiWorkerMirroredStrategy

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

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

tf.compat.v1.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.
Атрибуты
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 для запроса типа элементов, возвращаемых итератором. Это можно использовать для установки свойства input_signature функции tf.function. См. tf.distribute.DistributedDataset.element_spec для примера.

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

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

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

experimental_distribute_dataset

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

experimental_distribute_dataset(
    dataset, options=None
)

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

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

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

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

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

В приведенном выше фрагменте кода dataset разделяются на пакеты методом global_batch_size, и вызов experimental_distribute_dataset на нем перегруппировывает dataset в новый размер пакета, равный глобальному размеру пакета, деленному на количество реплик в синхронизации. Мы перебираем его с помощью цикла for в стиле Python. x является tf.distribute.DistributedValues, содержащим данные для всех реплик, и каждая реплика получает данные нового размера пакета. tf.distribute.Strategy.run позаботится о передаче правильных данных для каждой реплики в x нужным replica_fn , выполняемым на каждой реплике.

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

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

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

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

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

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

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

experimental_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,).

experimental_make_numpy_dataset

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

experimental_make_numpy_dataset(
    numpy_input, session=None
)

Создаёт tf.data.Dataset для ввода, предоставленного через массив NumPy.

Это позволяет избежать добавления numpy_input как большой константы в граф и копирует данные на машину или машины, которые будут обрабатывать входные данные.

Обратите внимание, что вам, скорее всего, потребуется использовать tf.distribute.Strategy.experimental_distribute_dataset с возвращённым набором данных для его дальнейшего распределения с помощью стратегии.

Пример:

numpy_input = np.ones([10], dtype=np.float32)
dataset = strategy.experimental_make_numpy_dataset(numpy_input)
dist_dataset = strategy.experimental_distribute_dataset(dataset)
Аргументы
numpy_input Вложенный набор массивов NumPy, которые будут преобразованы в набор данных. Обратите внимание, что списки массивов NumPy складываются, так как это нормальное поведение tf.data.Dataset.
session (Только для выполнения графа TensorFlow v1.x) Сессия, используемая для инициализации.
Возвращаемое значение
tf.data.Dataset, представляющий numpy_input.

experimental_run

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

experimental_run(
    fn, input_iterator=None
)

Выполняет операции в fn на каждой реплике с входными данными из input_iterator. (устарело)

Устарело: ЭТА ФУНКЦИЯ УСТАРЕЛА. Она будет удалена в будущей версии. Инструкции по обновлению: Этот метод недоступен в TF 2.x. Переключитесь на использование run вместо этого.
Устарело: Этот метод недоступен в TF 2.x. Переключитесь на использование run вместо этого.

При включенном режиме выполнения Eager, выполняет операции, указанные в fn на каждой реплике. В противном случае создает граф для выполнения операций на каждой реплике.

Каждая реплика будет получать один, разный вход из входов, предоставленных одним вызовом get_next для итератора ввода.

fn может вызвать tf.distribute.get_replica_context() для доступа к элементам, таким как replica_id_in_sync_group.

Важно: В зависимости от реализации tf.distribute.Strategy, используемой, и от того, включен ли режим выполнения Eager, fn может быть вызвано один или несколько раз (один раз для каждой реплики).
Аргументы
fn Функция для выполнения. Входы в функцию должны соответствовать выходам input_iterator.get_next(). Выход должен быть tf.nest из Tensors.
input_iterator (Необязательно) итератор ввода, из которого берутся входные данные.
Возвращаемое значение
Объединённое возвращаемое значение fn по всем репликам. Структура возвращаемого значения такая же, как и возвращаемого значения от fn. Каждый элемент структуры может быть PerReplica (если значения не синхронизированы), Mirrored (если значения синхронизированы), или Tensor (если выполняется на одной реплике).

make_dataset_iterator

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

make_dataset_iterator(
    dataset
)

Создаёт итератор для входных данных, предоставленных через dataset.

Устаревшее: Этот метод недоступен в TF 2.x.

Данные из заданного набора данных будут равномерно распределены по всем вычислительным репликам. Мы будем предполагать, что входной набор данных сгруппирован по глобальному размеру пакета. При этом предположении мы сделаем всё возможное, чтобы разделить каждый пакет по всем репликам (одному или нескольким работникам). Если эта попытка не удастся, будет выброшено исключение, и пользователь должен вместо этого использовать make_input_fn_iterator, который предоставляет больше контроля пользователю и не пытается разделить пакет по репликам.

Пользователь также может использовать make_input_fn_iterator, если хочет настроить, какие входные данные подаются на какую реплику/работника и т. д.

Аргументы
dataset tf.data.Dataset, который будет равномерно распределён по всем репликам.
Возвращаемое значение
Объект tf.distribute.InputIterator, который возвращает входные данные для каждого шага вычисления. Пользователь должен вызвать initialize на возвращённом итераторе.

make_input_fn_iterator

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

make_input_fn_iterator(
    input_fn,
    replication_mode=tf.distribute.InputReplicationMode.PER_WORKER
)

Возвращает итератор, разделённый по репликам, созданный из функции ввода.

Устаревшее: Этот метод недоступен в TF 2.x.

Функция input_fn должна принимать объект tf.distribute.InputContext, где можно получить информацию о пакетировании и фрагментации ввода:

def input_fn(input_context):
  batch_size = input_context.get_per_replica_batch_size(global_batch_size)
  d = tf.data.Dataset.from_tensors([[1.]]).repeat().batch(batch_size)
  return d.shard(input_context.num_input_pipelines,
                 input_context.input_pipeline_id)
with strategy.scope():
  iterator = strategy.make_input_fn_iterator(input_fn)
  replica_results = strategy.experimental_run(replica_fn, iterator)

Возвращаемый tf.data.Dataset функцией input_fn должен иметь размер пакета на реплику, который можно вычислить с помощью input_context.get_per_replica_batch_size.

Аргументы
input_fn Функция, принимающая объект tf.distribute.InputContext и возвращающая tf.data.Dataset.
replication_mode Значение перечисления tf.distribute.InputReplicationMode. В настоящее время поддерживается только PER_WORKER, что означает, что вызов input_fn будет выполняться один раз на каждом работнике. Реплики будут извлекать элементы из локального tf.data.Dataset на своих работниках.
Возвращаемое значение
Объект итератора, который нужно предварительно обработать с помощью .initialize()-ed. Затем его можно передать в strategy.experimental_run() или получить следующее значение с помощью iterator.get_next() для передачи в strategy.extended.call_for_each_replica().

reduce

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

reduce(
    reduce_op, value, axis=None
)

Сведение 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, скажем, потерю на пример, пакет будет разделен по всем репликам. Эта функция позволяет агрегировать по репликам и необязательно по элементам пакета, задавая параметр оси соответствующим образом.

Например, если у вас есть глобальный размер пакета 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 для сокращения только по репликам (например, если тензор не имеет размерности пакета).
Возвращаемое значение
A Tensor.

run

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

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

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

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

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

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

Важно: В зависимости от реализации tf.distribute.Strategy и включения режима eager execution, fn может быть вызван один или несколько раз. Если fn аннотирован с помощью tf.function или tf.distribute.Strategy.run вызван внутри tf.function (режим eager execution по умолчанию отключён внутри tf.function), fn вызывается один раз на каждую реплику для генерации графа TensorFlow, который затем будет повторно использован для выполнения с новыми входными данными. В противном случае, если режим eager execution включен, 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 захватывает информацию об области действия. При последующих вызовах методов высокоуровневого фреймворка обучения, таких как 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 иногда может потребоваться находиться внутри области действия, если оно создаёт переменные.
Возвращаемое значение
Менеджер контекста.

update_config_proto

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

update_config_proto(
    config_proto
)

Возвращает копию config_proto с изменениями, пригодными для использования с этой стратегией.

Устаревший: Этот метод недоступен в TF 2.x.

Обновлённая конфигурация содержит что-то необходимое для запуска стратегии, например, конфигурацию для запуска коллективных операций или фильтры устройств для повышения производительности распределённого обучения.

Аргументы
config_proto Объект tf.ConfigProto.
Возвращаемое значение
Обновлённая копия config_proto.

© 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/compat/v1/distribute/experimental/MultiWorkerMirroredStrategy

Spec-Zone.ru

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