tf.compat.v1.distribute.experimental.MultiWorkerMirroredStrategy
Стратегия распределения для синхронного обучения на нескольких рабочих узлах.
Наследуется от: Strategy
tf.compat.v1.distribute.experimental.MultiWorkerMirroredStrategy(
communication=tf.distribute.experimental.CollectiveCommunication.AUTO,
cluster_resolver=None
)
Эта стратегия реализует синхронное распределение обучения на нескольких рабочих узлах, каждый из которых может иметь несколько графических процессоров. Подобно tf.distribute.MirroredStrategy, она дублирует все переменные и вычисления на каждом локальном устройстве. Разница заключается в том, что она использует распределённую реализацию коллективных операций (например, all-reduce), чтобы несколько рабочих узлов могли работать вместе.
Необходимо запустить программу на каждом рабочем узле и правильно настроить cluster_resolver. Например, если вы используете tf.distribute.cluster_resolver.TFConfigClusterResolver, каждый рабочий узел должен иметь соответствующие значения task_type и task_id в переменной окружения TF_CONFIG. Пример TF_CONFIG для рабочего узла-0 кластера из двух рабочих узлов:
TF_CONFIG = '{"cluster": {"worker": ["localhost:12345", "localhost:23456"]}, "task": {"type": "worker", "index": 0} }'
Ваша программа запускается на каждом рабочем узле как есть. Обратите внимание, что коллективные операции требуют участия каждого рабочего узла. Все tf.distribute и не-tf.distribute API могут использовать коллективные операции внутри, например, для создания контрольных точек и сохранения, поскольку чтение tf.Variable с tf.VariableSynchronization.ON_READ выполняет all-reduce для значения. Поэтому рекомендуется запускать на каждом рабочем узле точно такую же программу. Распределение на основе task_type или task_id рабочего узла чревато ошибками.
cluster_resolver.num_accelerators() определяет количество графических процессоров, используемых стратегией. Если значение равно нулю, стратегия использует процессор. Все рабочие узлы должны использовать одинаковое количество устройств; в противном случае поведение неопределённо.
Эта стратегия не предназначена для TPU. Используйте tf.distribute.TPUStrategy вместо неё.
После настройки TF_CONFIG использование этой стратегии аналогично использованию tf.distribute.MirroredStrategy и tf.distribute.TPUStrategy.
strategy = tf.distribute.MultiWorkerMirroredStrategy()
with strategy.scope():
model = tf.keras.Sequential([
tf.keras.layers.Dense(2, input_shape=(5,)),
])
optimizer = tf.keras.optimizers.SGD(learning_rate=0.1)
def dataset_fn(ctx):
x = np.random.random((2, 5)).astype(np.float32)
y = np.random.randint(2, size=(2, 1))
dataset = tf.data.Dataset.from_tensor_slices((x, y))
return dataset.repeat().batch(1, drop_remainder=True)
dist_dataset = strategy.distribute_datasets_from_function(dataset_fn)
model.compile()
model.fit(dist_dataset)
Вы также можете написать свой собственный цикл обучения:
@tf.function
def train_step(iterator):
def step_fn(inputs):
features, labels = inputs
with tf.GradientTape() as tape:
logits = model(features, training=True)
loss = tf.keras.losses.sparse_categorical_crossentropy(
labels, logits)
grads = tape.gradient(loss, model.trainable_variables)
optimizer.apply_gradients(zip(grads, model.trainable_variables))
strategy.run(step_fn, args=(next(iterator),))
for _ in range(NUM_STEP):
train_step(iterator)
Подробный учебник см. в статье Обучение на нескольких рабочих узлах с помощью Keras.
Сохранение
Необходимо сохранять и создавать контрольные точки на всех рабочих узлах, а не только на одном. Это связано с тем, что переменные, для которых synchronization=ON_READ, вызывают агрегацию во время сохранения. Рекомендуется сохранять в разные пути на каждом рабочем узле, чтобы избежать гонок. Каждый рабочий узел сохраняет одно и то же. Примеры см. в руководстве Обучение на нескольких рабочих узлах с помощью Keras.
Известные проблемы
-
tf.distribute.cluster_resolver.TFConfigClusterResolverне возвращает правильное количество ускорителей. Стратегия использует все доступные графические процессоры, еслиcluster_resolvertf.distribute.cluster_resolver.TFConfigClusterResolverилиNone. - В режиме eager стратегия должна быть создана до вызова любого другого API TensorFlow.
| Атрибуты | |
|---|---|
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 будет вызываться на устройстве CPU каждого из рабочих узлов, и каждый из них сгенерирует набор данных, где каждая реплика на этом рабочем узле будет извлекать одну партию входных данных (т. е. если на рабочем узле две реплики, то с 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 в стиле Python. x — это tf.distribute.DistributedValues, содержащий данные для всех реплик, и каждая реплика получает данные нового размера пакета. tf.distribute.Strategy.run позаботится о подаче правильных данных на реплику в x в правильный replica_fn метод, выполняемый на каждой реплике.
Фрагментация включает в себя автоматическую фрагментацию по нескольким рабочим узлам и внутри каждого рабочего узла. Во-первых, при распределённом обучении на нескольких рабочих узлах (то есть, когда вы используете tf.distribute.experimental.MultiWorkerMirroredStrategy или tf.distribute.TPUStrategy), автоматическая фрагментация набора данных по набору рабочих узлов означает, что каждому рабочему узлу назначается подмножество всего набора данных (если установлен правильный tf.data.experimental.AutoShardPolicy). Это необходимо для обеспечения того, что на каждом шаге глобальный размер пакета непересекающихся элементов набора данных будет обрабатываться каждым рабочим узлом. Автоматическая фрагментация имеет несколько различных вариантов, которые можно указать с помощью tf.data.experimental.DistributeOptions. Затем фрагментация внутри каждого рабочего узла означает, что метод разделит данные между всеми устройствами рабочего узла (если их более одного). Это произойдёт независимо от автоматической фрагментации по нескольким рабочим узлам.
Примечание: для автоматической фрагментации по нескольким рабочим узлам режим по умолчанию —tf.data.experimental.AutoShardPolicy.AUTO. Этот режим попытается фрагментировать входной набор данных по файлам, если набор данных создаётся из наборов данных для чтения (например,tf.data.TFRecordDataset,tf.data.TextLineDatasetи т. д.) или иначе фрагментировать набор данных по данным, где каждый из рабочих узлов будет читать весь набор данных и обрабатывать только назначенный ему фрагмент. Однако, если у вас меньше одного входного файла на рабочий узел, мы рекомендуем отключить автоматическую фрагментацию набора данных по рабочим узлам, установивtf.data.experimental.DistributeOptions.auto_shard_policyв значениеtf.data.experimental.AutoShardPolicy.OFF.
По умолчанию этот метод добавляет преобразование предварительной выборки в конец предоставленного пользователем экземпляра tf.data.Dataset. Аргумент преобразования предварительной выборки, который является buffer_size, равен количеству реплик в синхронизации.
Если описанная выше логика разделения пакетов и фрагментации набора данных нежелательна, используйте tf.distribute.Strategy.distribute_datasets_from_function вместо этого, который не выполняет автоматической пакетной обработки или фрагментации.
Примечание: Если вы используете TPUStrategy, порядок обработки данных рабочими узлами при использованииtf.distribute.Strategy.experimental_distribute_datasetилиtf.distribute.Strategy.distribute_datasets_from_functionне гарантируется. Это обычно необходимо, если вы используетеtf.distributeдля масштабирования прогнозирования. Тем не менее, вы можете вставить индекс для каждого элемента в пакете и упорядочить выводы соответственно. Обратитесь к этому фрагменту для примера того, как упорядочить выводы.
Примечание: Состоятельные преобразования набора данных в настоящее время не поддерживаются сtf.distribute.experimental_distribute_datasetилиtf.distribute.distribute_datasets_from_function. Любые состоятельные операции, которые может иметь набор данных, в настоящее время игнорируются. Например, если ваш набор данных имеетmap_fn, использующийtf.random.uniformдля поворота изображения, то у вас есть граф набора данных, зависящий от состояния (то есть от случайного зерна) на локальной машине, где выполняется процесс Python.
Для получения учебного пособия по более подробному использованию и свойствам этого метода см. учебное пособие по распределенному вводу. Если вас интересует обработка последнего частичного пакета, прочитайте эту секцию.
| Аргументы | |
|---|---|
dataset | tf.data.Dataset, который будет фрагментирован по всем репликам в соответствии с описанными выше правилами. |
options | tf.distribute.InputOptions, используемый для управления параметрами распределения этого набора данных. |
| Возвращаемые значения | |
|---|---|
tf.distribute.DistributedDataset. |
experimental_local_results
experimental_local_results(
value
)
Возвращает список всех локальных значений на реплику, содержащихся в value.
Примечание: Это возвращает только значения на рабочем узле, инициированном этим клиентом. При использованииtf.distribute.Strategy, например,tf.distribute.experimental.MultiWorkerMirroredStrategy, каждый рабочий узел будет своим клиентом, и эта функция вернёт только значения, вычисленные на этом рабочем узле.
| Аргументы | |
|---|---|
value | Значение, возвращённое experimental_run(), run(), extended.call_for_each_replica(), или переменной, созданной в scope . |
| Возвращаемые значения | |
|---|---|
Кортеж значений, содержащихся в value . Если value представляет одно значение, возвращается (value,). |
experimental_make_numpy_dataset
experimental_make_numpy_dataset(
numpy_input, session=None
)
Создаёт tf.data.Dataset для входных данных, предоставленных через массив NumPy.
Это позволяет избежать добавления numpy_input в виде большой константы в граф и копирует данные на машину или машины, которые будут обрабатывать вход.
Обратите внимание, что, вероятно, вам потребуется использовать tf.distribute.Strategy.experimental_distribute_dataset с возвращаемым набором данных для его дальнейшего распределения с помощью стратегии.
Пример:
numpy_input = np.ones([10], dtype=np.float32) dataset = strategy.experimental_make_numpy_dataset(numpy_input) dist_dataset = strategy.experimental_distribute_dataset(dataset)
| Аргументы | |
|---|---|
numpy_input | Вложенный массив NumPy входных массивов, который будет преобразован в набор данных. Обратите внимание, что списки массивов NumPy складываются, так как это обычное поведение tf.data.Dataset. |
session | (Только для выполнения графов TensorFlow v1.x) Сессия, используемая для инициализации. |
| Возвращаемые значения | |
|---|---|
tf.data.Dataset, представляющий numpy_input . |
experimental_run
experimental_run(
fn, input_iterator=None
)
Выполняет операции в fn на каждой реплике, используя входные данные из input_iterator.
УСТАРЕЛО: Этот метод недоступен в TF 2.x. Пожалуйста, переключитесь на использование run вместо этого.
При включённом выполнении Eager выполняет операции, заданные fn на каждой реплике. В противном случае создаёт граф для выполнения операций на каждой реплике.
Каждая реплика получит один отдельный вход из входных данных, предоставленных одним вызовом get_next на итератор входных данных.
fn может вызвать tf.distribute.get_replica_context() для доступа к членам, таким как replica_id_in_sync_group.
| Аргументы | |
|---|---|
fn | Функция для выполнения. Входные данные для функции должны соответствовать выходам input_iterator.get_next() . Выход должен быть tf.nest из Tensor . |
input_iterator | (Необязательно) итератор входных данных, из которого берутся входные данные. |
| Возвращаемое значение | |
|---|---|
Объединенное возвращаемое значение fn по всем репликам. Структура возвращаемого значения такая же, как и возвращаемое значение от fn. Каждый элемент структуры может быть PerReplica (если значения не синхронизированы), Mirrored (если значения синхронизированы), или Tensor (если выполняется на одной реплике). |
make_dataset_iterator
make_dataset_iterator(
dataset
)
Создаёт итератор для ввода, предоставленного через dataset.
УСТАРЕВШЕЕ: Этот метод недоступен в TF 2.x.
Данные из заданного набора данных будут равномерно распределены по всем вычислительным репликам. Мы будем исходить из того, что входной набор данных сгруппирован по глобальному размеру пакета. С этим предположением, мы сделаем всё возможное, чтобы разделить каждый пакет по всем репликам (один или несколько рабочих узлов). Если эта попытка не удастся, будет выброшено исключение, и пользователь должен вместо этого использовать make_input_fn_iterator, который предоставляет пользователю больше контроля и не пытается разделить пакет между репликами.
Пользователь также может использовать make_input_fn_iterator, если хочет настроить, какой ввод подаётся на какую реплику/рабочий узел и т.д.
| Аргументы | |
|---|---|
dataset | tf.data.Dataset, который будет равномерно распределён по всем репликам. |
| Возвращаемое значение | |
|---|---|
Итератор, который возвращает ввод для каждого шага вычисления. Пользователь должен вызвать initialize для возвращённого итератора. |
make_input_fn_iterator
make_input_fn_iterator(
input_fn, replication_mode=tf.distribute.InputReplicationMode.PER_WORKER
)
Возвращает итератор, разделённый между репликами, созданный из функции ввода.
УСТАРЕВШЕЕ: Этот метод недоступен в TF 2.x.
Функция input_fn должна принимать объект tf.distribute.InputContext, где можно получить информацию о группировке по пакетам и фрагментации ввода:
def input_fn(input_context):
batch_size = input_context.get_per_replica_batch_size(global_batch_size)
d = tf.data.Dataset.from_tensors([[1.]]).repeat().batch(batch_size)
return d.shard(input_context.num_input_pipelines,
input_context.input_pipeline_id)
with strategy.scope():
iterator = strategy.make_input_fn_iterator(input_fn)
replica_results = strategy.experimental_run(replica_fn, iterator)
Возвращаемый tf.data.Dataset функцией input_fn должен иметь размер пакета для каждой реплики, который может быть вычислен с помощью input_context.get_per_replica_batch_size.
| Аргументы | |
|---|---|
input_fn | Функция, которая принимает объект tf.distribute.InputContext и возвращает tf.data.Dataset. |
replication_mode | значение перечисления tf.distribute.InputReplicationMode. В настоящее время поддерживается только PER_WORKER, что означает, что будет один вызов input_fn на каждый рабочий узел. Реплики будут извлекать данные из локального tf.data.Dataset на своих рабочих узлах. |
| Возвращаемое значение | |
|---|---|
Объект итератора, который должен быть предварительно .initialize()-н. Затем его можно передать в strategy.experimental_run() или получить следующее значение для передачи strategy.extended.call_for_each_replica() с помощью iterator.get_next(). |
reduce
reduce(
reduce_op, value, axis=None
)
Свести value по репликам и вернуть результат на текущем устройстве.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
def step_fn():
i = tf.distribute.get_replica_context().replica_id_in_sync_group
return tf.identity(i)
per_replica_result = strategy.run(step_fn)
total = strategy.reduce("SUM", per_replica_result, axis=None)
total
<tf.Tensor: shape=(), dtype=int32, numpy=1>
Чтобы увидеть, как это будет выглядеть с несколькими репликами, рассмотрите тот же пример с MirroredStrategy с 2 графическими процессорами:
strategy = tf.distribute.MirroredStrategy(devices=["GPU:0", "GPU:1"])
def step_fn():
i = tf.distribute.get_replica_context().replica_id_in_sync_group
return tf.identity(i)
per_replica_result = strategy.run(step_fn)
# Check devices on which per replica result is:
strategy.experimental_local_results(per_replica_result)[0].device
# /job:localhost/replica:0/task:0/device:GPU:0
strategy.experimental_local_results(per_replica_result)[1].device
# /job:localhost/replica:0/task:0/device:GPU:1
total = strategy.reduce("SUM", per_replica_result, axis=None)
# Check device on which reduced result is:
total.device
# /job:localhost/replica:0/task:0/device:CPU:0
Этот API обычно используется для агрегирования результатов, возвращаемых различными репликами, например, для отчётов. Например, потерю, вычисленную из различных реплик, можно усреднить с помощью этого API перед выводом.
Примечание: Результат копируется на «текущее» устройство — обычно это процессор рабочего узла, на котором выполняется программа. ДляTPUStrategy, это первый хост TPU. Для многоклиентскогоMultiWorkerMirroredStrategy, это процессор каждого рабочего узла.
Существует ряд различных API tf.distribute для сворачивания значений по репликам:
-
tf.distribute.ReplicaContext.all_reduce: Это отличается отStrategy.reduceтем, что предназначено для контекста реплики и не копирует результаты на устройство хоста.all_reduceобычно используется для сворачиваний внутри шага обучения, например, для градиентов. -
tf.distribute.StrategyExtended.reduce_toиtf.distribute.StrategyExtended.batch_reduce_to: Эти API являются более продвинутыми версиямиStrategy.reduce, так как они позволяют настроить место назначения результата. Они также вызываются в контексте между репликами.
Что должно быть значением 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 для сворачивания только по репликам (например, если тензор не имеет размерности пакета). |
| Возвращаемое значение | |
|---|---|
| Тензор. |
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 должны быть либо Python-значениями вложенной структуры тензоров, например, списком тензоров, в этом случае args и kwargs будут переданы в вызываемое на каждой реплике fn. Или args или kwargs могут быть tf.distribute.DistributedValues, содержащими тензоры или составные тензоры, т.е. tf.compat.v1.TensorInfo.CompositeTensor, в этом случае каждый вызов fn получит компонент tf.distribute.DistributedValues, соответствующий его реплике.
Пример использования:
- Ввод тензора-константы.
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. Его элементы могут быть Python-значением, тензором или tf.distribute.DistributedValues. |
kwargs | Необязательные именованные аргументы для fn. Его элементы могут быть Python-значением, тензором или 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 tutorial.
Что должно быть в области и что за её пределами?
Существует ряд требований к тому, что должно происходить внутри области. Однако в местах, где у нас есть информация о используемой стратегии, мы часто входим в область для пользователя, чтобы он не должен был делать это явно (т. е. вызов внутри или вне области допустим).
- Всё, что создаёт переменные, которые должны быть распределёнными переменными, должно находиться в
strategy.scope. Это может быть сделано либо путём прямого размещения в области, либо путём использования другого API, такого какstrategy.runилиmodel.fitдля входа в неё. Любая переменная, созданная вне области, не будет распределена и может иметь последствия для производительности. Типичные вещи, создающие переменные в TF: модели, оптимизаторы, метрики. Они всегда должны создаваться внутри области. Ещё одним источником создания переменных может быть восстановление контрольной точки — когда переменные создаются лениво. Обратите внимание, что любая переменная, созданная внутри стратегии, сохраняет информацию о стратегии. Чтение и запись в эти переменные вне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иногда может потребоваться внутри области, если оно создаёт переменные.
| Возвращаемое значение | |
|---|---|
| Контекстный менеджер. |
update_config_proto
update_config_proto(
config_proto
)
Возвращает копию config_proto, модифицированную для использования с данной стратегией.
УСТАРЕЛО: Этот метод недоступен в TF 2.x.
Обновлённая конфигурация содержит необходимые данные для работы с стратегией, например, настройки для выполнения коллективных операций или фильтры устройств для повышения производительности распределённого обучения.
| Аргументы | |
|---|---|
config_proto | Объект tf.ConfigProto. |
| Возвращаемое значение | |
|---|---|
Обновлённая копия config_proto. |
© 2020 The TensorFlow Authors. All rights reserved.
Licensed under the Creative Commons Attribution License 3.0.
Code samples licensed under the Apache 2.0 License.
https://www.tensorflow.org/versions/r2.4/api_docs/python/tf/compat/v1/distribute/experimental/MultiWorkerMirroredStrategy