tf.distribute.MirroredStrategy
| Просмотреть исходный код на GitHub |
Синхронное обучение на нескольких репликах на одной машине.
Наследуется от: Strategy
tf.distribute.MirroredStrategy(
devices=None, cross_device_ops=None
)
Эта стратегия обычно используется для обучения на одной машине с несколькими графическими процессорами. Для TPUs используйте tf.distribute.TPUStrategy. Чтобы использовать MirroredStrategy с несколькими рабочими узлами, обратитесь к tf.distribute.experimental.MultiWorkerMirroredStrategy.
Например, переменная, созданная в рамках MirroredStrategy, является MirroredVariable. Если в конструкторе стратегии не указаны устройства, она будет использовать все доступные графические процессоры. Если графических процессоров не найдено, она будет использовать доступные центральные процессоры. Обратите внимание, что TensorFlow обрабатывает все центральные процессоры на машине как одно устройство и использует потоки внутри для параллелизма.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
with strategy.scope():
x = tf.Variable(1.)
x
MirroredVariable:{
0: <tf.Variable ... shape=() dtype=float32, numpy=1.0>,
1: <tf.Variable ... shape=() dtype=float32, numpy=1.0>
}
При использовании стратегий распределения все создание переменных должно выполняться в области действия стратегии. Это позволит продублировать переменные на всех репликах и поддерживать их синхронизацию с помощью алгоритма all-reduce.
Переменные, созданные внутри MirroredStrategy, который обернут в tf.function, по-прежнему MirroredVariables.
x = []
@tf.function # Wrap the function with tf.function.
def create_variable():
if not x:
x.append(tf.Variable(1.))
return x[0]
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
with strategy.scope():
_ = create_variable()
print(x[0])
MirroredVariable:{
0: <tf.Variable ... shape=() dtype=float32, numpy=1.0>,
1: <tf.Variable ... shape=() dtype=float32, numpy=1.0>
}
experimental_distribute_dataset можно использовать для распределения набора данных по репликам при написании собственного цикла обучения. Если вы используете .fit и .compile методы, доступные в tf.keras, то tf.keras будет обрабатывать распределение за вас.
Например:
my_strategy = tf.distribute.MirroredStrategy()
with my_strategy.scope():
@tf.function
def distribute_train_epoch(dataset):
def replica_fn(input):
# process input and return result
return result
total_result = 0
for x in dataset:
per_replica_result = my_strategy.run(replica_fn, args=(x,))
total_result += my_strategy.reduce(tf.distribute.ReduceOp.SUM,
per_replica_result, axis=None)
return total_result
dist_dataset = my_strategy.experimental_distribute_dataset(dataset)
for _ in range(EPOCHS):
train_result = distribute_train_epoch(dist_dataset)
| Аргументы | |
|---|---|
devices | список строк устройств, таких как ['/gpu:0', '/gpu:1']. Если None, используются все доступные графические процессоры. Если графические процессоры не найдены, используется центральный процессор. |
cross_device_ops | необязательно, потомок CrossDeviceOps. Если это не задано, по умолчанию будет использоваться NcclAllReduce(). Это нужно настраивать, если NCCL недоступен или если доступна специальная реализация, которая использует конкретное оборудование. |
| Атрибуты | |
|---|---|
cluster_resolver | Возвращает решатель кластера, связанный с этой стратегией. В общем случае, при использовании многоузловой Стратегии, которые намерены иметь связанный Стратегии с одним рабочим узлом обычно не имеют
os.environ['TF_CONFIG'] = json.dumps({
'cluster': {
'worker': ["localhost:12345", "localhost:23456"],
'ps': ["localhost:34567"]
},
'task': {'type': 'worker', 'index': 0}
})
# This implicitly uses TF_CONFIG for the cluster and current task info.
strategy = tf.distribute.experimental.MultiWorkerMirroredStrategy()
...
if strategy.cluster_resolver.task_type == 'worker':
# Perform something that's only applicable on workers. Since we set this
# as a worker above, this block will run on this particular instance.
elif strategy.cluster_resolver.task_type == 'ps':
# Perform something that's only applicable on parameter servers. Since we
# set this as a worker above, this block will not run on this particular
# instance.
Для получения дополнительной информации, пожалуйста, обратитесь к документации API |
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 до нового размера пакета, равного глобальному размеру пакета, деленному на количество реплик в синхронизации. Мы перебираем его с помощью Python-цикла for. x — это tf.distribute.DistributedValues, содержащий данные для всех реплик, и каждая реплика получает данные нового размера пакета. tf.distribute.Strategy.run позаботится о подаче правильных данных на реплику в x в соответствующую replica_fn , выполняемую на каждой реплике.
Разбиение на фрагменты (sharding) включает автоматическое разбиение на фрагменты (autosharding) по множеству рабочих процессов и внутри каждого рабочего процесса. Во-первых, в распределённом обучении с несколькими рабочими процессами (то есть, когда вы используете tf.distribute.experimental.MultiWorkerMirroredStrategy или tf.distribute.TPUStrategy), автоматическое разбиение на фрагменты набора данных по набору рабочих процессов означает, что каждому рабочему процессу назначается подмножество всего набора данных (если установлен правильный tf.data.experimental.AutoShardPolicy). Это гарантирует, что на каждой итерации каждый рабочий процесс будет обрабатывать глобальный размер пакета непересекающихся элементов набора данных. Автоматическое разбиение на фрагменты имеет несколько различных вариантов, которые можно указать с помощью tf.data.experimental.DistributeOptions. Затем, разбиение на фрагменты внутри каждого рабочего процесса означает, что метод разделит данные среди всех устройств рабочего процесса (если их более одного). Это произойдёт независимо от автоматического разбиения на фрагменты с использованием нескольких рабочих процессов.
Примечание: для автоматического разбиения на фрагменты с использованием нескольких рабочих процессов, режим по умолчанию —tf.data.experimental.AutoShardPolicy.AUTO. Этот режим попытается разбить входной набор данных по файлам, если набор данных создаётся на основе наборов данных для чтения (например,tf.data.TFRecordDataset,tf.data.TextLineDatasetи т. д.) или иначе разбить набор данных по данным, где каждый из рабочих процессов будет читать весь набор данных и обрабатывать только фрагмент, назначенный ему. Однако, если у вас менее одного входного файла на рабочий процесс, мы рекомендуем отключить автоматическое разбиение на фрагменты набора данных для рабочих процессов, установивtf.data.experimental.DistributeOptions.auto_shard_policyна значениеtf.data.experimental.AutoShardPolicy.OFF.
По умолчанию этот метод добавляет преобразование предварительной выборки в конце предоставленного пользователем экземпляра tf.data.Dataset. Аргумент преобразования предварительной выборки, который является buffer_size, равен количеству реплик в синхронизации.
Если описанная выше логика разделения пакетов и разбиения на фрагменты набора данных нежелательна, используйте tf.distribute.Strategy.distribute_datasets_from_function, который не выполняет автоматическое формирование пакетов или разбиение на фрагменты.
Примечание: Если вы используете TPUStrategy, порядок обработки данных рабочими процессами при использованииtf.distribute.Strategy.experimental_distribute_datasetилиtf.distribute.Strategy.distribute_datasets_from_functionне гарантируется. Это обычно необходимо, если вы используетеtf.distributeдля масштабирования прогнозирования. Однако вы можете вставить индекс для каждого элемента в пакете и упорядочить результаты соответственно. Обратитесь к этому фрагменту для примера того, как упорядочить результаты.
Примечание: В настоящее время преобразования набора данных со состоянием не поддерживаются сtf.distribute.experimental_distribute_datasetилиtf.distribute.distribute_datasets_from_function. Любые операторы со состоянием, которые может содержать набор данных, в настоящее время игнорируются. Например, если ваш набор данных содержитmap_fn, который используетtf.random.uniformдля поворота изображения, то у вас есть граф набора данных, зависящий от состояния (например, от случайного зерна) на локальной машине, где выполняется процесс Python.
Для получения дополнительных руководств и свойств этого метода обратитесь к руководству по распределённому вводу. Если вас интересует обработка последнего частичного пакета, прочитайте эту секцию.
| Аргументы | |
|---|---|
dataset | tf.data.Dataset, который будет разделен по фрагментам на всех репликах в соответствии с вышеуказанными правилами. |
options | tf.distribute.InputOptions, используемый для управления параметрами распределения этого набора данных. |
| Возвращаемые значения | |
|---|---|
tf.distribute.DistributedDataset. |
experimental_distribute_values_from_function
experimental_distribute_values_from_function(
value_fn
)
Генерирует tf.distribute.DistributedValues из value_fn.
Эта функция предназначена для генерации tf.distribute.DistributedValues для передачи в run, reduce, или другие методы, принимающие распределённые значения, когда не используются наборы данных.
| Аргументы | |
|---|---|
value_fn | Функция для запуска генерации значений. Она вызывается для каждой реплики с tf.distribute.ValueContext в качестве единственного аргумента. Она должна возвращать тензор или тип, который можно преобразовать в тензор. |
| Возвращаемые значения | |
|---|---|
tf.distribute.DistributedValues, содержащий значение для каждой реплики. |
Пример использования:
- Возвратить постоянное значение для каждой реплики:
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, rank(value)). |
| Возвращаемые значения | |
|---|---|
Tensor, который представляет собой конкатенацию value по репликам вдоль axis размерности. |
reduce
reduce(
reduce_op, value, axis
)
Свести value по всем репликам и вернуть результат на текущем устройстве.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
def step_fn():
i = tf.distribute.get_replica_context().replica_id_in_sync_group
return tf.identity(i)
per_replica_result = strategy.run(step_fn)
total = strategy.reduce("SUM", per_replica_result, axis=None)
total
<tf.Tensor: shape=(), dtype=int32, numpy=1>
Чтобы увидеть, как это будет выглядеть с несколькими репликами, рассмотрите тот же пример с MirroredStrategy с 2 графическими процессорами:
strategy = tf.distribute.MirroredStrategy(devices=["GPU:0", "GPU:1"])
def step_fn():
i = tf.distribute.get_replica_context().replica_id_in_sync_group
return tf.identity(i)
per_replica_result = strategy.run(step_fn)
# Check devices on which per replica result is:
strategy.experimental_local_results(per_replica_result)[0].device
# /job:localhost/replica:0/task:0/device:GPU:0
strategy.experimental_local_results(per_replica_result)[1].device
# /job:localhost/replica:0/task:0/device:GPU:1
total = strategy.reduce("SUM", per_replica_result, axis=None)
# Check device on which reduced result is:
total.device
# /job:localhost/replica:0/task:0/device:CPU:0
Этот API обычно используется для агрегирования результатов, возвращаемых различными репликами, для отчётности и т. д. Например, потерю, вычисленную на разных репликах, можно усреднить с помощью этого API перед выводом.
Примечание: Результат копируется на «текущее» устройство — обычно это процессор узла, на котором выполняется программа. ДляTPUStrategy, это первый хост TPU. Для многоклиентскогоMultiWorkerMirroredStrategy, это процессор каждого узла.
Существует ряд разных API tf.distribute для сведения значений по всем репликам:
-
tf.distribute.ReplicaContext.all_reduce: Это отличается отStrategy.reduceтем, что оно предназначено для контекста реплики и не копирует результаты на устройство узла.all_reduceобычно используется для вычислений внутри этапа обучения, таких как градиенты. -
tf.distribute.StrategyExtended.reduce_toиtf.distribute.StrategyExtended.batch_reduce_to: Эти API являются более продвинутыми версиямиStrategy.reduce, поскольку они позволяют настраивать место назначения результата. Они также вызываются в контексте между репликами.
Каким должен быть ось?
Учитывая значение на реплике, возвращаемое run, например, потерю на пример, пакет будет разделен между всеми репликами. Эта функция позволяет агрегировать по репликам и, по желанию, также по элементам пакета, указав параметр axis соответствующим образом.
Например, если у вас есть глобальный размер пакета 8 и 2 реплики, значения для примеров [0, 1, 2, 3] будут на реплике 0, а [4, 5, 6, 7] — на реплике 1. С axis=None, reduce будет агрегировать только по репликам, возвращая [0+4, 1+5, 2+6, 3+7]. Это полезно, когда каждая реплика вычисляет скаляр или какое-то другое значение без размерности «пакет» (например, градиент или потерю).
strategy.reduce("sum", per_replica_result, axis=None)
Иногда нужно агрегировать как по глобальному пакету, так и по всем репликам. Это поведение можно получить, указав размер пакета как axis, обычно axis=0. В этом случае будет возвращен скаляр 0+1+2+3+4+5+6+7.
strategy.reduce("sum", per_replica_result, axis=0)
Если есть последний частичный пакет, необходимо указать ось, чтобы форма результата была согласована между репликами. Так, если последний пакет имеет размер 6 и разделен на [0, 1, 2, 3] и [4, 5], вы получите несовпадение форм, если не укажете axis=0. Если вы укажете tf.distribute.ReduceOp.MEAN, используя axis=0 будет использоваться правителььная знаменатель 6. Сравните это с вычислением reduce_mean для получения скалярного значения на каждой реплике, и этой функцией для усреднения этих средних значений, что приведет к разному взвешиванию значений 1/8 и 1/4.
| Аргументы | |
|---|---|
reduce_op | значение tf.distribute.ReduceOp, определяющее, как должны комбинироваться значения. Допускается использование строкового представления перечисления, такого как «SUM», «MEAN». |
value | экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, для объединения в один тензор. Он также может быть обычным тензором, когда используется с OneDeviceStrategy или стратегией по умолчанию. |
axis | указывает размерность для сведения по каждой реплике тензора. Обычно должен быть установлен на размерность пакета или None для сведения только по репликам (например, если тензор не имеет размерности пакета). |
| Возвращаемое значение | |
|---|---|
Tensor. |
run
run(
fn, args=(), kwargs=None, options=None
)
Вызывает fn на каждой реплике с заданными аргументами.
Этот метод является основным способом распределения вычислений с объектом tf.distribute. Он вызывает fn на каждой реплике. Если args или kwargs имеют tf.distribute.DistributedValues, такие как те, что получены от tf.distribute.DistributedDataset из tf.distribute.Strategy.experimental_distribute_dataset или tf.distribute.Strategy.distribute_datasets_from_function, когда fn выполняется на определенной реплике, он будет выполняться с компонентом tf.distribute.DistributedValues, соответствующим этой реплике.
fn вызывается в контексте реплики. fn может вызвать tf.distribute.get_replica_context() для доступа к членам, таким как all_reduce. Подробнее о контексте реплики см. в документе tf.distribute.
Все аргументы в args или kwargs могут быть вложенной структурой тензоров, например, списком тензоров, в этом случае args и kwargs будут переданы в fn , вызываемый на каждой реплике. Или args или kwargs могут быть tf.distribute.DistributedValues, содержащими тензоры или составные тензоры, т. е. tf.compat.v1.TensorInfo.CompositeTensor, в этом случае каждый вызов fn получит компонент tf.distribute.DistributedValues, соответствующий его реплике. Обратите внимание, что произвольные значения Python, не являющиеся указанных типов, не поддерживаются.
Пример использования:
- Входной тензор-константа.
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для явного распределения вычисления. См. пример в руководстве по обучению с пользовательским циклом custom training loop.
Что должно находиться в области действия, а что вне её?
Существует ряд требований к тому, что должно происходить в рамках области видимости. Однако в тех местах, где у нас есть информация о том, какая стратегия используется, мы часто входим в область видимости для пользователя, чтобы он не должен был делать это явно (т.е. вызов этих функций внутри или вне области видимости допустим).
- Любой код, создающий переменные, которые должны быть распределенными, должен вызываться в
strategy.scope. Это можно сделать, либо непосредственно вызвав функцию создания переменной в контексте области видимости, либо используя другой API, например,strategy.runилиkeras.Model.fit, чтобы автоматически войти в него. Любые переменные, созданные вне области видимости, не будут распределены и могут привести к проблемам производительности. Некоторые общие объекты, создающие переменные в TF, — это модели, оптимизаторы, метрики. Такие объекты всегда должны инициализироваться в области видимости, и любые функции, которые могут лениво создавать переменные (например,Model.__call__(), отслеживаниеtf.functionи т. д.) также должны вызываться в рамках области видимости. Другим источником создания переменных может быть восстановление контрольной точки — когда переменные создаются лениво. Обратите внимание, что любая переменная, созданная внутри стратегии, содержит информацию о стратегии. Таким образом, чтение и запись в эти переменные внеstrategy.scopeтакже могут работать без проблем, без необходимости для пользователя входить в область видимости. - Некоторые API стратегии (например,
strategy.runиstrategy.reduce) которые должны находиться в области видимости стратегии, автоматически входят в нее, что означает, что при использовании этих API вам не нужно явно входить в область видимости. - Когда
tf.keras.Modelсоздается внутриstrategy.scope, объект модели получает информацию о области видимости. При вызове методов высокого уровня для обучения, таких какmodel.compile,model.fit, и т. д., захваченная область видимости будет автоматически входить в нее, и соответствующая стратегия будет использоваться для распределения обучения и т. д. Подробный пример см. в учебнике по распределенному Keras. ВНИМАНИЕ: простой вызовmodel(..)не автоматически входит в захваченную область видимости — только API высокого уровня для обучения поддерживают это поведение:model.compile,model.fit,model.evaluate,model.predictиmodel.saveмогут быть вызваны внутри или вне области видимости. - Следующее может находиться как внутри, так и вне области видимости:
- Создание наборов входных данных
- Определение
tf.functionпредставляющих ваш шаг обучения - API сохранения, такие как
tf.saved_model.save. Загрузка создает переменные, поэтому это нужно выполнять в области видимости, если вы хотите обучать модель в распределенном режиме. - Сохранение контрольных точек. Как упоминалось выше -
checkpoint.restoreиногда может потребоваться находиться в области видимости, если оно создает переменные.
| Возвращаемое значение | |
|---|---|
| Менеджер контекста. |
© 2022 The TensorFlow Authors. All rights reserved.
Licensed under the Creative Commons Attribution License 4.0.
Code samples licensed under the Apache 2.0 License.
https://www.tensorflow.org/versions/r2.9/api_docs/python/tf/distribute/MirroredStrategy