tf.distribute.experimental.MultiWorkerMirroredStrategy
Стратегия распределения для синхронного обучения на нескольких рабочих узлах.
Наследуется от: MultiWorkerMirroredStrategy, Strategy
tf.distribute.experimental.MultiWorkerMirroredStrategy(
communication=tf.distribute.experimental.CollectiveCommunication.AUTO,
cluster_resolver=None
)
Используется в ноутбуках
| Используется в учебниках |
|---|
Эта стратегия реализует синхронное распределение обучения на нескольких рабочих узлах, каждый из которых может иметь несколько GPU. Подобно tf.distribute.MirroredStrategy, она дублирует все переменные и вычисления на каждом локальном устройстве. Разница в том, что она использует реализацию распределенных коллективных операций (например, all-reduce), чтобы несколько рабочих узлов могли работать вместе.
Вам необходимо запустить свою программу на каждом рабочем узле и правильно настроить cluster_resolver. Например, если вы используете tf.distribute.cluster_resolver.TFConfigClusterResolver, каждый рабочий узел должен иметь соответствующие task_type и task_id, установленные в переменной окружения TF_CONFIG. Пример TF_CONFIG на рабочем узле-0 кластера из двух рабочих узлов:
TF_CONFIG = '{"cluster": {"worker": ["localhost:12345", "localhost:23456"]}, "task": {"type": "worker", "index": 0} }'
Ваша программа выполняется на каждом рабочем узле как есть. Обратите внимание, что коллективные операции требуют участия каждого рабочего узла. Все tf.distribute и не tf.distribute API могут использовать коллективные операции внутри, например, сохранение и создание контрольных точек, поскольку чтение tf.Variable с tf.VariableSynchronization.ON_READ выполняет all-reduce значения. Поэтому рекомендуется запускать точно такую же программу на каждом рабочем узле. Распределение задач на основе task_type или task_id рабочего узла чревато ошибками.
cluster_resolver.num_accelerators() определяет количество GPU, используемых стратегией. Если оно равно нулю, стратегия использует ЦП. Все рабочие узлы должны использовать одинаковое количество устройств, в противном случае поведение не определено.
Эта стратегия не предназначена для TPU. Используйте tf.distribute.TPUStrategy вместо этого.
После настройки TF_CONFIG использование этой стратегии аналогично использованию tf.distribute.MirroredStrategy и tf.distribute.TPUStrategy.
strategy = tf.distribute.MultiWorkerMirroredStrategy()
with strategy.scope():
model = tf.keras.Sequential([
tf.keras.layers.Dense(2, input_shape=(5,)),
])
optimizer = tf.keras.optimizers.SGD(learning_rate=0.1)
def dataset_fn(ctx):
x = np.random.random((2, 5)).astype(np.float32)
y = np.random.randint(2, size=(2, 1))
dataset = tf.data.Dataset.from_tensor_slices((x, y))
return dataset.repeat().batch(1, drop_remainder=True)
dist_dataset = strategy.distribute_datasets_from_function(dataset_fn)
model.compile()
model.fit(dist_dataset)
Вы также можете написать свою собственную петлю обучения:
@tf.function
def train_step(iterator):
def step_fn(inputs):
features, labels = inputs
with tf.GradientTape() as tape:
logits = model(features, training=True)
loss = tf.keras.losses.sparse_categorical_crossentropy(
labels, logits)
grads = tape.gradient(loss, model.trainable_variables)
optimizer.apply_gradients(zip(grads, model.trainable_variables))
strategy.run(step_fn, args=(next(iterator),))
for _ in range(NUM_STEP):
train_step(iterator)
См. Обучение на нескольких рабочих узлах с помощью Keras для подробного учебника.
Сохранение
Вам нужно сохранять и создавать контрольные точки на всех рабочих узлах, а не только на одном. Это связано с тем, что переменные, у которых synchronization=ON_READ, вызывают агрегацию при сохранении. Рекомендуется сохранять в разные пути на каждом рабочем узле, чтобы избежать гонок. Каждый рабочий узел сохраняет одно и то же. См. учебник Обучение на нескольких рабочих узлах с помощью Keras для примеров.
Известные проблемы
-
tf.distribute.cluster_resolver.TFConfigClusterResolverне возвращает правильное количество ускорителей. Стратегия использует все доступные GPU, еслиcluster_resolver—tf.distribute.cluster_resolver.TFConfigClusterResolverилиNone. - В режиме eager стратегия должна быть создана до вызова любого другого API TensorFlow.
| Аргументы | |
|---|---|
communication | необязательная tf.distribute.experimental.CommunicationImplementation. Это подсказка о предпочтительной реализации коллективной коммуникации. Возможные значения включают AUTO, RING и NCCL. |
cluster_resolver | необязательный tf.distribute.cluster_resolver.ClusterResolver. Если None, используется tf.distribute.cluster_resolver.TFConfigClusterResolver. |
| Атрибуты | |
|---|---|
cluster_resolver | Возвращает решатель кластера, связанный с этой стратегией. В качестве стратегии для нескольких рабочих узлов |
extended | tf.distribute.StrategyExtended с дополнительными методами. |
num_replicas_in_sync | Возвращает количество реплик, по которым агрегируются градиенты. |
Методы
distribute_datasets_from_function
distribute_datasets_from_function(
dataset_fn, options=None
)
Распределяет экземпляры tf.data.Dataset, созданные вызовами dataset_fn.
Аргумент dataset_fn, который передается пользователями, — это функция-вход, которая имеет аргумент tf.distribute.InputContext и возвращает экземпляр tf.data.Dataset. Ожидается, что возвращаемый из dataset_fn набор данных уже разбит по размеру пакета на реплику (т. е. глобальный размер пакета, деленный на количество реплик в синхронизации) и распределен. tf.distribute.Strategy.distribute_datasets_from_function не группирует и не распределяет экземпляр набора данных, возвращаемый функцией входных данных. dataset_fn будет вызываться на устройстве ЦП каждого из рабочих узлов, и каждый создаст набор данных, в котором каждая реплика на этом рабочем узле будет извлекать один пакет входных данных (т. е. если на рабочем узле две реплики, с каждого шага из Dataset будут извлечены два пакета).
Этот метод может использоваться для нескольких целей. Во-первых, он позволяет указать свою логику группирования и распределения. (В отличие от tf.distribute.experimental_distribute_dataset, который выполняет группирование и распределение за вас.) Например, где experimental_distribute_dataset не может распределить входные файлы, этот метод может использоваться для ручного распределения набора данных (избегая медленного поведения по умолчанию в experimental_distribute_dataset). В случаях, когда набор данных бесконечен, это распределение можно выполнить, создав копии набора данных, которые различаются только случайным зерном.
dataset_fn должен принимать экземпляр tf.distribute.InputContext, где можно получить информацию о группировании и репликации входных данных.
Вы можете использовать свойство element_spec возвращаемого этим API tf.distribute.DistributedDataset для запроса tf.TypeSpec элементов, возвращаемых итератором. Это можно использовать для установки свойства input_signature tf.function. Следуйте tf.distribute.DistributedDataset.element_spec, чтобы увидеть пример.
Примечание: Если вы используете TPUStrategy, порядок обработки данных рабочими узлами при использованииtf.distribute.Strategy.experimental_distribute_datasetилиtf.distribute.Strategy.distribute_datasets_from_functionне гарантируется. Это обычно требуется, если вы используетеtf.distributeдля масштабирования прогнозирования. Однако вы можете вставить индекс для каждого элемента в пакете и упорядочить выходы соответственно. См. этот фрагмент для примера упорядочивания выходов.
Примечание: Состоятельные преобразования наборов данных в настоящее время не поддерживаются сtf.distribute.experimental_distribute_datasetилиtf.distribute.distribute_datasets_from_function. Любые состоятельные операции, которые может иметь набор данных, в настоящее время игнорируются. Например, если ваш набор данных имеетmap_fn, который используетtf.random.uniformдля поворота изображения, то у вас есть граф набора данных, который зависит от состояния (т. е. случайного зерна) на локальной машине, где выполняется процесс Python.
Для получения дополнительной информации об использовании и свойствах этого метода обратитесь к учебнику по распределенным входным данным. Если вас интересует обработка последнего частичного пакета, прочитайте эту часть.
END_OF_DOCUMENT_MARKER ```| Args | |
|---|---|
dataset_fn | Функция, принимающая экземпляр tf.distribute.InputContext и возвращающая tf.data.Dataset. |
options | tf.distribute.InputOptions, используемая для управления параметрами распределения набора данных. |
| Returns | |
|---|---|
tf.distribute.DistributedDataset. |
experimental_distribute_dataset
experimental_distribute_dataset(
dataset, options=None
)
Создаёт tf.distribute.DistributedDataset из tf.data.Dataset.
Возвращённый tf.distribute.DistributedDataset можно итерировать, как обычные наборы данных. ПРИМЕЧАНИЕ: Пользователь не может добавить больше преобразований к tf.distribute.DistributedDataset. Вы можете только создать итератор или изучить tf.TypeSpec генерируемых им данных. Подробнее см. в документации tf.distribute.DistributedDataset.
Вот пример:
global_batch_size = 2
# Passing the devices is optional.
strategy = tf.distribute.MirroredStrategy(devices=["GPU:0", "GPU:1"])
# Create a dataset
dataset = tf.data.Dataset.range(4).batch(global_batch_size)
# Distribute that dataset
dist_dataset = strategy.experimental_distribute_dataset(dataset)
@tf.function
def replica_fn(input):
return input*2
result = []
# Iterate over the `tf.distribute.DistributedDataset`
for x in dist_dataset:
# process dataset elements
result.append(strategy.run(replica_fn, args=(x,)))
print(result)
[PerReplica:{
0: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([0])>,
1: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([2])>
}, PerReplica:{
0: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([4])>,
1: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([6])>
}]Внутри этого метода происходит три ключевых действия: группировка, фрагментация и предварительная загрузка.
В приведённом фрагменте кода dataset группируется по global_batch_size, а вызов experimental_distribute_dataset на нём перегруппировывает dataset в новый размер пакета, равный глобальному размеру пакета, делённому на число реплик в синхронизации. Мы итерируем через него с помощью цикла Python for. x — tf.distribute.DistributedValues, содержащий данные для всех реплик, и каждая реплика получает данные нового размера пакета. tf.distribute.Strategy.run позаботится о передаче правильных данных каждой реплике в x в правильную replica_fn, выполняемую на каждой реплике.
Фрагментация включает автоматическую фрагментацию по нескольким рабочим узлам и внутри каждого рабочего узла. Во-первых, при распределённом обучении по нескольким рабочим узлам (т. е. при использовании tf.distribute.experimental.MultiWorkerMirroredStrategy или tf.distribute.TPUStrategy) автоматическая фрагментация набора данных по набору рабочих узлов означает, что каждому рабочему узлу назначается подмножество всего набора данных (если установлен соответствующий tf.data.experimental.AutoShardPolicy). Это гарантирует, что на каждом шаге глобальный размер пакета непересекающихся элементов набора данных будет обработан каждым рабочим узлом. Автоматическая фрагментация имеет несколько разных параметров, которые можно указать с помощью tf.data.experimental.DistributeOptions. Затем фрагментация внутри каждого рабочего узла означает, что метод разделит данные между всеми устройствами рабочего узла (если их несколько). Это произойдёт независимо от автоматической фрагментации по нескольким рабочим узлам.
Примечание: для автоматической фрагментации по нескольким рабочим узлам режим по умолчанию —tf.data.experimental.AutoShardPolicy.AUTO. Этот режим попытается фрагментировать входной набор данных по файлам, если набор данных создаётся на основе наборов данных считывателя (например,tf.data.TFRecordDataset,tf.data.TextLineDatasetи т. д.) или фрагментировать набор данных по данным, где каждый из рабочих узлов прочитает весь набор данных и обработает только назначенный ему фрагмент. Однако, если у вас меньше одного входного файла на рабочий узел, мы рекомендуем отключить автоматическую фрагментацию набора данных по рабочим узлам, установивtf.data.experimental.DistributeOptions.auto_shard_policyвtf.data.experimental.AutoShardPolicy.OFF.
По умолчанию этот метод добавляет преобразование предварительной загрузки в конец предоставленного пользователем экземпляра tf.data.Dataset. Аргумент преобразования предварительной загрузки, который составляет buffer_size, равен числу реплик в синхронизации.
Если описанная выше логика разделения пакетов и фрагментации наборов данных нежелательна, используйте tf.distribute.Strategy.distribute_datasets_from_function вместо этого, так как он не выполняет автоматическую группировку или фрагментацию.
Примечание: Если вы используете TPUStrategy, порядок обработки данных рабочими узлами при использованииtf.distribute.Strategy.experimental_distribute_datasetилиtf.distribute.Strategy.distribute_datasets_from_functionне гарантируется. Это обычно необходимо, если вы используетеtf.distributeдля масштабирования прогнозирования. Однако вы можете вставить индекс для каждого элемента в пакете и упорядочить результаты соответственно. Обратитесь к этому фрагменту для примера упорядочивания результатов.
Примечание: В настоящее время преобразования наборов данных с состоянием не поддерживаются сtf.distribute.experimental_distribute_datasetилиtf.distribute.distribute_datasets_from_function. Любые операторы с состоянием, которые может иметь набор данных, в настоящее время игнорируются. Например, если ваш набор данных имеетmap_fn, который используетtf.random.uniformдля поворота изображения, то у вас есть граф набора данных, зависящий от состояния (т. е. от случайного зерна) на локальном компьютере, где выполняется процесс Python.
Для получения более подробной информации об использовании и свойствах этого метода обратитесь к уроку по распределённому вводу. Если вас интересует обработка последнего частичного пакета, прочитайте эту часть.
| Args | |
|---|---|
dataset | tf.data.Dataset, который будет фрагментирован по всем репликам в соответствии с описанными выше правилами. |
options | tf.distribute.InputOptions, используемая для управления параметрами распределения набора данных. |
| Returns | |
|---|---|
tf.distribute.DistributedDataset. |
experimental_distribute_values_from_function
experimental_distribute_values_from_function(
value_fn
)
Генерирует tf.distribute.DistributedValues из value_fn.
Эта функция предназначена для генерации tf.distribute.DistributedValues для передачи в run, reduce или другие методы, принимающие распределённые значения, когда не используются наборы данных.
| Args | |
|---|---|
value_fn | Функция для выполнения генерации значений. Она вызывается для каждой реплики с tf.distribute.ValueContext в качестве единственного аргумента. Она должна возвращать тензор или тип, который можно преобразовать в тензор. |
| Returns | |
|---|---|
tf.distribute.DistributedValues, содержащий значение для каждой реплики. |
Пример использования:
-
Возврат постоянного значения для каждой реплики:
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>) -
Распределение значений в массиве на основе идентификатора реплики:
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) -
Указание значений с использованием 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) -
Размещение значений на устройствах и распределение:
strategy = tf.distribute.TPUStrategy() worker_devices = strategy.extended.worker_devices multiple_values = [] for i in range(strategy.num_replicas_in_sync): with tf.device(worker_devices[i]): multiple_values.append(tf.constant(1.0)) def value_fn(ctx): return multiple_values[ctx.replica_id_in_sync_group] distributed_values = strategy. experimental_distribute_values_from_function( value_fn)
experimental_local_results
experimental_local_results(
value
)
Возвращает список всех локальных значений каждой реплики, содержащихся в value.
Примечание: Это возвращает только значения на рабочем узле, запущенном этим клиентом. При использованииtf.distribute.Strategy(например,tf.distribute.experimental.MultiWorkerMirroredStrategy), каждый рабочий узел будет своим клиентом, и эта функция будет возвращать только значения, вычисленные на этом рабочем узле.
| Args | |
|---|---|
value | Значение, возвращённое experimental_run(), run(), or a variable created inscope`. |
| Returns | |
|---|---|
Кортеж значений, содержащихся в value, где i-й элемент соответствует i-й реплике. Если value представляет одно значение, возвращается (value,). |
gather
gather(
value, axis
)
Собирает value по репликам вдоль axis на текущее устройство.
При заданном объекте tf.distribute.DistributedValues или tf.Tensor подобного типа value, этот API собирает и конкатенирует value по всем репликам вдоль axis-й размерности. Результат копируется на «текущее» устройство, которое обычно является процессором (CPU) рабочего узла, на котором выполняется программа. Для tf.distribute.TPUStrategy это первый хост TPU. Для многоклиентского tf.distribute.MultiWorkerMirroredStrategy это процессор (CPU) каждого рабочего узла.
Этот API может вызываться только в контексте межрепликации. Для аналога в контексте реплики см. tf.distribute.ReplicaContext.all_gather.
Примечание: Для всех стратегий, кромеtf.distribute.TPUStrategy, входныеvalueна разных репликах должны иметь одинаковый ранг, а их формы должны быть одинаковыми во всех измерениях, кромеaxis-й размерности. Другими словами, их формы не могут отличаться в измеренииd, гдеdне равно аргументуaxis. Например, при заданномtf.distribute.DistributedValuesс тензорами компонентов формы(1, 2, 3)и(1, 3, 3)на двух репликах вы можете вызватьgather(..., axis=1, ...), но неgather(..., axis=0, ...)илиgather(..., axis=2, ...). Однако дляtf.distribute.TPUStrategy.gatherвсе тензоры должны иметь ровно одинаковый ранг и форму.
Примечание: При заданномtf.distribute.DistributedValuesvalue, его составные тензоры должны иметь ненулевой ранг. В противном случае, рассмотрите использованиеtf.expand_dimsперед их сбором.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
# A DistributedValues with component tensor of shape (2, 1) on each replica
distributed_values = strategy.experimental_distribute_values_from_function(lambda _: tf.identity(tf.constant([[1], [2]])))
@tf.function
def run():
return strategy.gather(distributed_values, axis=0)
run()
<tf.Tensor: shape=(4, 1), dtype=int32, numpy=
array([[1],
[2],
[1],
[2]], dtype=int32)>Рассмотрим следующий пример для получения дополнительных комбинаций:
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1", "GPU:2", "GPU:3"])
single_tensor = tf.reshape(tf.range(6), shape=(1,2,3))
distributed_values = strategy.experimental_distribute_values_from_function(lambda _: tf.identity(single_tensor))
@tf.function
def run(axis):
return strategy.gather(distributed_values, axis=axis)
axis=0
run(axis)
<tf.Tensor: shape=(4, 2, 3), dtype=int32, numpy=
array([[[0, 1, 2],
[3, 4, 5]],
[[0, 1, 2],
[3, 4, 5]],
[[0, 1, 2],
[3, 4, 5]],
[[0, 1, 2],
[3, 4, 5]]], dtype=int32)>
axis=1
run(axis)
<tf.Tensor: shape=(1, 8, 3), dtype=int32, numpy=
array([[[0, 1, 2],
[3, 4, 5],
[0, 1, 2],
[3, 4, 5],
[0, 1, 2],
[3, 4, 5],
[0, 1, 2],
[3, 4, 5]]], dtype=int32)>
axis=2
run(axis)
<tf.Tensor: shape=(1, 2, 12), dtype=int32, numpy=
array([[[0, 1, 2, 0, 1, 2, 0, 1, 2, 0, 1, 2],
[3, 4, 5, 3, 4, 5, 3, 4, 5, 3, 4, 5]]], dtype=int32)>| Аргументы | |
|---|---|
value | экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который должен быть объединен в один тензор. Он также может быть обычным тензором при использовании с tf.distribute.OneDeviceStrategy или стратегией по умолчанию. Тензоры, составляющие DistributedValues, могут быть только плотными тензорами с ненулевым рангом, а НЕ tf.IndexedSlices. |
axis | 0-мерный тензор int32. Измерение, вдоль которого происходит сбор. Должно быть в диапазоне [0, ранг(значение)). |
| Возвращаемое значение | |
|---|---|
Tensor, являющийся конкатенацией value по всем репликам вдоль axis-го измерения. |
reduce
reduce(
reduce_op, value, axis
)
Суммирование value по всем репликам и возврат результата на текущем устройстве.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
def step_fn():
i = tf.distribute.get_replica_context().replica_id_in_sync_group
return tf.identity(i)
per_replica_result = strategy.run(step_fn)
total = strategy.reduce("SUM", per_replica_result, axis=None)
total
<tf.Tensor: shape=(), dtype=int32, numpy=1>Чтобы увидеть, как это будет выглядеть с несколькими репликами, рассмотрите тот же пример с MirroredStrategy с 2 графическими процессорами (GPU):
strategy = tf.distribute.MirroredStrategy(devices=["GPU:0", "GPU:1"])
def step_fn():
i = tf.distribute.get_replica_context().replica_id_in_sync_group
return tf.identity(i)
per_replica_result = strategy.run(step_fn)
# Check devices on which per replica result is:
strategy.experimental_local_results(per_replica_result)[0].device
# /job:localhost/replica:0/task:0/device:GPU:0
strategy.experimental_local_results(per_replica_result)[1].device
# /job:localhost/replica:0/task:0/device:GPU:1
total = strategy.reduce("SUM", per_replica_result, axis=None)
# Check device on which reduced result is:
total.device
# /job:localhost/replica:0/task:0/device:CPU:0
Этот API обычно используется для агрегирования результатов, возвращаемых различными репликами, для отчётности и т.п. Например, потерю, вычисленную на разных репликах, можно усреднить с помощью этого API перед выводом.
Примечание: Результат копируется на «текущее» устройство — обычно это процессор (CPU) рабочего узла, на котором выполняется программа. ДляTPUStrategyэто первый хост TPU. Для многоклиентскогоMultiWorkerMirroredStrategyэто процессор (CPU) каждого рабочего узла.
Существует ряд различных API tf.distribute для уменьшения значений по всем репликам:
-
tf.distribute.ReplicaContext.all_reduce: Это отличается отStrategy.reduceтем, что предназначено для контекста реплики и не копирует результаты на устройство хоста.all_reduceобычно используется для уменьшения внутри шага обучения, таких как градиенты. -
tf.distribute.StrategyExtended.reduce_toиtf.distribute.StrategyExtended.batch_reduce_to: Эти API являются более продвинутыми версиямиStrategy.reduce, так как они позволяют настраивать место назначения результата. Они также вызываются в контексте межрепликации.
Что должно быть значением axis?
Учитывая значение на каждой реплике, возвращаемое run, например, потерю на пример, партия будет разделена между всеми репликами. Эта функция позволяет агрегировать по репликам и, при необходимости, также по элементам партии, указав параметр axis соответствующим образом.
Например, если у вас глобальный размер партии 8 и 2 реплики, значения для примеров [0, 1, 2, 3] будут на реплике 0, а [4, 5, 6, 7] будут на реплике 1. С axis=None, reduce будет агрегировать только по репликам, возвращая [0+4, 1+5, 2+6, 3+7]. Это полезно, когда каждая реплика вычисляет скаляр или какое-либо другое значение, не имеющее «размерности партии» (например, градиент или потерю).
strategy.reduce("sum", per_replica_result, axis=None)
Иногда вам нужно агрегировать как по глобальной партии, так и по всем репликам. Этого можно добиться, указав размер партии в качестве axis, обычно axis=0. В этом случае он вернёт скалярное значение 0+1+2+3+4+5+6+7.
strategy.reduce("sum", per_replica_result, axis=0)
Если есть последняя частичная партия, вам необходимо указать ось, чтобы результирующая форма была согласована между репликами. Так что, если последняя партия имеет размер 6 и разделена на [0, 1, 2, 3] и [4, 5], у вас возникнет несоответствие формы, если не указать axis=0. Если вы укажете tf.distribute.ReduceOp.MEAN, используя axis=0, будет использоваться правильный знаменатель 6. Противопоставьте это вычислению reduce_mean для получения скалярного значения на каждой реплике и этой функции для усреднения этих средних значений, что будет учитывать некоторые значения 1/8 и другие 1/4.
| Аргументы | |
|---|---|
reduce_op | значение tf.distribute.ReduceOp, указывающее, как должны комбинироваться значения. Допускается использование строкового представления перечисления, такого как "SUM", "MEAN". |
value | экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который должен быть объединен в один тензор. Он также может быть обычным тензором при использовании с OneDeviceStrategy или стратегией по умолчанию. |
axis | определяет измерение, по которому следует уменьшить тензор каждой реплики. Обычно нужно установить значение размерности партии или None, чтобы уменьшить только по репликам (например, если у тензора нет размерности партии). |
| Возвращаемое значение | |
|---|---|
Tensor. |
run
run(
fn, args=(), kwargs=None, options=None
)
Вызывает fn на каждой реплике с заданными аргументами.
Этот метод является основным способом распределения вычислений с объектом tf.distribute. Он вызывает fn на каждой реплике. Если args или kwargs содержат tf.distribute.DistributedValues, например, такие, которые генерируются tf.distribute.DistributedDataset из tf.distribute.Strategy.experimental_distribute_dataset или tf.distribute.Strategy.distribute_datasets_from_function, когда fn выполняется на определённой реплике, она будет выполняться с компонентом tf.distribute.DistributedValues, соответствующим этой реплике.
fn вызывается в контексте реплики. fn может вызвать tf.distribute.get_replica_context() для доступа к членам, таким как all_reduce. См. строку документации tf.distribute на концепцию контекста реплики.
Все аргументы в args или kwargs могут быть вложенной структурой тензоров, например, списком тензоров, в этом случае args и kwargs будут переданы в вызываемую на каждой реплике функцию fn. Или args или kwargs могут быть tf.distribute.DistributedValues, содержащими тензоры или составные тензоры, т.е. tf.compat.v1.TensorInfo.CompositeTensor, в этом случае каждый вызов fn получит компонент tf.distribute.DistributedValues, соответствующий его реплике. Обратите внимание, что произвольные Python-значения, не являющиеся типами выше, не поддерживаются.
Пример использования:
-
Постоянный тензорный вход.
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> } -
Вход 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> -
Используйте
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" на каждом рабочем узле.
Примечание: Вход в область не автоматически распределяет вычисление, за исключением случаев использования высокоуровневых фреймворков обучения, таких как Kerasmodel.fit. Если вы не используетеmodel.fit, необходимо использовать APIstrategy.runдля явного распределения вычислений. См. пример в учебнике по пользовательскому циклу обучения .
Что должно быть внутри и что снаружи области?
Существует ряд требований к тому, что должно происходить внутри области. Однако там, где у нас есть информация о используемой стратегии, мы часто входим в область за пользователя, чтобы он не должен был делать это явно (т. е. вызов внутри или вне области является допустимым).
- Все, что создает переменные, которые должны быть распределенными переменными, должно вызываться в
strategy.scope. Это можно сделать, либо вызвав функцию создания переменных непосредственно в контексте области, либо используя другой API, например,strategy.runилиkeras.Model.fitдля автоматического входа в него. Любая переменная, созданная вне области, не будет распределена и может иметь последствия для производительности. Некоторые общие объекты, создающие переменные в TF, - это модели, оптимизаторы, метрики. Такие объекты всегда должны инициализироваться в области, а любые функции, которые могут создавать переменные ленивым образом (например,Model.call(), отслеживаниеtf.functionи т. д.) также должны вызываться в области. Другим источником создания переменных может быть восстановление контрольной точки - когда переменные создаются ленивым образом. Обратите внимание, что любая переменная, созданная внутри стратегии, сохраняет информацию о стратегии. Таким образом, чтение и запись этих переменных внеstrategy.scopeтакже могут работать без проблем, без необходимости ввода пользователем в область. - Некоторые API стратегии (например,
strategy.runиstrategy.reduce), которые требуют наличия в области стратегии, автоматически входят в область, что означает, что при использовании этих API вам не нужно явно входить в область. - Когда
tf.keras.Modelсоздается внутриstrategy.scope, объект модели сохраняет информацию о области. При последующем вызове методов высокоуровневого обучения, таких какmodel.compile,model.fitи т. д., захваченная область будет автоматически введена, и связанная стратегия будет использована для распределения обучения и т. д. См. подробный пример в учебнике по распределенному Keras. ПРЕДУПРЕЖДЕНИЕ: простой вызовmodel(..)не приводит к автоматическому входу в захваченную область - только высокоуровневые API обучения поддерживают это поведение:model.compile,model.fit,model.evaluate,model.predictиmodel.saveмогут быть вызваны как внутри, так и вне области. - Следующее может быть как внутри, так и вне области:
- Создание наборов данных входных данных
- Определение
tf.function, представляющих ваш шаг обучения - API сохранения, такие как
tf.saved_model.save. Загрузка создает переменные, поэтому это должно быть внутри области, если вы хотите обучить модель в распределенном режиме. - Сохранение контрольных точек. Как упоминалось выше -
checkpoint.restoreиногда может потребоваться находиться внутри области, если оно создает переменные.
| Возвращаемое значение | |
|---|---|
| Менеджер контекста. |
© 2022 The TensorFlow Authors. All rights reserved.
Licensed under the Creative Commons Attribution License 4.0.
Code samples licensed under the Apache 2.0 License.
https://www.tensorflow.org/api_docs/python/tf/distribute/experimental/MultiWorkerMirroredStrategy