tf.distribute.MultiWorkerMirroredStrategy
Стратегия распределения для синхронного обучения на нескольких рабочих узлах.
Наследуется от: Strategy
tf.distribute.MultiWorkerMirroredStrategy(
cluster_resolver=None, communication_options=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.
| Аргументы | |
|---|---|
cluster_resolver | необязательный tf.distribute.cluster_resolver.ClusterResolver. Если None, используется tf.distribute.cluster_resolver.TFConfigClusterResolver. |
communication_options | необязательный tf.distribute.experimental.CommunicationOptions. Это настраивает опции по умолчанию для меж-устройств коммуникаций. Их можно переопределить, используя опции, предоставленные API коммуникаций, например, tf.distribute.ReplicaContext.all_reduce. Подробнее см. в tf.distribute.experimental.CommunicationOptions. |
| Атрибуты | |
|---|---|
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 будет вызываться на устройстве CPU каждого из рабочих узлов, и каждый из них создаст набор данных, где каждая реплика на этом рабочем узле будет извлекать по одному пакету ввода (т. е. если у рабочего узла две реплики, с каждого шага из 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.
Примечание: Если вы используете 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_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, содержащий значение для каждой реплики. |
Пример использования:
-
Возвращение постоянного значения для каждой реплики:
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>) -
Распределение значений в массиве на основе 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) -
Указание значений с помощью 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, каждый рабочий узел будет своим клиентом, и эта функция вернёт только значения, вычисленные на этом рабочем узле.
| Аргументы | |
|---|---|
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.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 перед печатью.
Примечание: Результат копируется на «текущее» устройство, которое обычно является процессором узла, на котором выполняется программа. Для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, определяющее, как следует объединять значения. Разрешает использование строкового представления перечисления, такого как «СУММА», «СРЕДНЕЕ». |
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/MultiWorkerMirroredStrategy