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