tf.distribute.experimental.MultiWorkerMirroredStrategy
| Просмотреть исходный код на GitHub |
Стратегия распределения для синхронного обучения на нескольких рабочих узлах.
Наследуется от: Strategy
tf.distribute.experimental.MultiWorkerMirroredStrategy(
communication=tf.distribute.experimental.CollectiveCommunication.AUTO,
cluster_resolver=None
)
Данная стратегия реализует синхронное распределённое обучение на нескольких рабочих узлах, каждый из которых может иметь несколько графических процессоров. Подобно tf.distribute.MirroredStrategy, она создаёт копии всех переменных модели на каждом устройстве на всех рабочих узлах.
Она использует реализацию multi-worker all-reduce коллективных операций для синхронизации переменных. Коллективная операция — это единственная операция в графе TensorFlow, которая может автоматически выбрать алгоритм all-reduce в среде выполнения TensorFlow в соответствии с аппаратным обеспечением, топологией сети и размерами тензоров.
По умолчанию она использует все локальные графические процессоры или процессор для обучения на одном рабочем узле.
Когда переменная окружения 'TF_CONFIG' установлена, она анализирует cluster_spec, task_type и task_id из 'TF_CONFIG' и превращает их в стратегию для нескольких рабочих узлов, которая дублирует модели на графических процессорах всех машин в кластере. В текущей реализации используются все графические процессоры в кластере, и предполагается, что все рабочие узлы имеют одинаковое количество графических процессоров.
Вы также можете передать экземпляр distribute.cluster_resolver.ClusterResolver при инициализации стратегии. Тип задачи, идентификатор задачи и т.д. будут анализироваться из экземпляра решателя вместо переменной окружения TF_CONFIG.
Поддерживает как режим eager, так и режим графа. Однако для режима eager она должна настроить контекст eager в своём конструкторе, и поэтому все операции в режиме eager должны выполняться после создания объекта стратегии.
| Аргументы | |
|---|---|
communication | необязательный перечисление типа distribute.experimental.CollectiveCommunication. Это предоставляет способ для пользователя переопределить выбор коллективной операции связи. Возможные значения включают AUTO, RING, и NCCL. |
cluster_resolver | необязательный объект distribute.cluster_resolver.ClusterResolver. По умолчанию используется TFConfigClusterResolver, который инициализируется из переменной окружения TF_CONFIG. |
| Атрибуты | |
|---|---|
cluster_resolver | Возвращает решатель кластера, связанный с данной стратегией. Как стратегия для нескольких рабочих узлов, |
extended | tf.distribute.StrategyExtended с дополнительными методами. |
num_replicas_in_sync | Возвращает количество реплик, по которым агрегируются градиенты. |
Методы
experimental_assign_to_logical_device
experimental_assign_to_logical_device(
tensor, logical_device_id
)
Добавляет аннотацию, что tensor будет назначено логическому устройству.
Примечание: Этот API поддерживается только в TPUStrategy на данный момент. Это добавляет аннотацию кtensor, указывая, что операции сtensorбудут вызваны на логическом устройстве с идентификаторомlogical_device_id. При использовании параллелизма моделей, по умолчанию все операции размещаются на нулевом логическом устройстве.
# Initializing TPU system with 2 logical devices and 4 replicas.
resolver = tf.distribute.cluster_resolver.TPUClusterResolver(tpu='')
tf.config.experimental_connect_to_cluster(resolver)
topology = tf.tpu.experimental.initialize_tpu_system(resolver)
device_assignment = tf.tpu.experimental.DeviceAssignment.build(
topology,
computation_shape=[1, 1, 1, 2],
num_replicas=4)
strategy = tf.distribute.TPUStrategy(
resolver, experimental_device_assignment=device_assignment)
iterator = iter(inputs)
@tf.function()
def step_fn(inputs):
output = tf.add(inputs, inputs)
# Add operation will be executed on logical device 0.
output = strategy.experimental_assign_to_logical_device(output, 0)
return output
strategy.run(step_fn, args=(next(iterator),))
| Аргументы | |
|---|---|
tensor | Входной тензор для аннотации. |
logical_device_id | Идентификатор логического ядра, которому будет назначен тензор. |
| Исключения | |
|---|---|
ValueError | Указанный идентификатор логического устройства не соответствует общему количеству разделов, указанных в задании устройств. |
| Возвращаемое значение | |
|---|---|
Аннотированный тензор с идентичным значением, как у tensor. |
experimental_distribute_dataset
experimental_distribute_dataset(
dataset, options=None
)
Создаёт tf.distribute.DistributedDataset из tf.data.Dataset.
Возвращаемый tf.distribute.DistributedDataset можно итерировать аналогично обычным наборам данных. ПРИМЕЧАНИЕ: Пользователь не может добавить больше преобразований к tf.distribute.DistributedDataset.
Вот пример:
strategy = tf.distribute.MirroredStrategy() # Create a dataset dataset = dataset_ops.Dataset.TFRecordDataset([ "/a/1.tfr", "/a/2.tfr", "/a/3.tfr", "/a/4.tfr"]) # Distribute that dataset dist_dataset = strategy.experimental_distribute_dataset(dataset) # Iterate over the `tf.distribute.DistributedDataset` for x in dist_dataset: # process dataset elements strategy.run(replica_fn, args=(x,))
В приведенном фрагменте кода tf.distribute.DistributedDataset dist_dataset группируется по GLOBAL_BATCH_SIZE, и мы итерируемся по нему с помощью for x in dist_dataset. x a tf.distribute.DistributedValues, содержащий данные для всех реплик, которые агрегируются в пакет из GLOBAL_BATCH_SIZE. tf.distribute.Strategy.run позаботится о подаче правильных данных для каждой реплики в x на правильные replica_fn , выполненные на каждой реплике.
Что происходит «под капотом» этого метода, когда мы говорим, что экземпляр tf.data.Dataset - dataset - распределяется? Это зависит от того, как вы установили tf.data.experimental.AutoShardPolicy через tf.data.experimental.DistributeOptions. По умолчанию он установлен на tf.data.experimental.AutoShardPolicy.AUTO. В многоузловой настройке мы сначала попытаемся распределить dataset путем определения, создаётся ли dataset из наборов данных для чтения (например, tf.data.TFRecordDataset, tf.data.TextLineDataset и т.д.) и, если это так, попытаемся разбить входные файлы. Обратите внимание, что должно быть по крайней мере один входной файл на каждый рабочий узел. Если у вас меньше одного входного файла на каждый рабочий узел, мы рекомендуем отключить разделение наборов данных между рабочими узлами, установив tf.data.experimental.DistributeOptions.auto_shard_policy на tf.data.experimental.AutoShardPolicy.OFF.
Если попытка разбить по файлам не удалась (т.е. набор данных не читается из файлов), мы разделим набор данных равномерно в конце, добавив операцию .shard в конец цепочки обработки. Это заставит всю цепочку предобработки для всех данных выполняться на каждом рабочем узле, и каждый рабочий узел будет выполнять избыточную работу. Мы выведем предупреждение, если этот путь будет выбран.
Как упоминалось ранее, внутри каждого рабочего узла мы также разделим данные между всеми устройствами рабочего узла (если их несколько). Это произойдёт, даже если многоузловое разделение отключено.
Если вышеописанная логика разделения пакетов и разбиения наборов данных нежелательна, используйте tf.distribute.Strategy.experimental_distribute_datasets_from_function вместо этого, которая не выполняет автоматического разделения или разбиения.
Вы также можете использовать свойство element_spec экземпляра tf.distribute.DistributedDataset, возвращаемого этим API, для запроса типа элементов, возвращаемых итератором. Это можно использовать для установки свойства input_signature функции tf.function.
strategy = tf.distribute.MirroredStrategy() # Create a dataset dataset = dataset_ops.Dataset.TFRecordDataset([ "/a/1.tfr", "/a/2.tfr", "/a/3.tfr", "/a/4.tfr"]) # Distribute that dataset dist_dataset = strategy.experimental_distribute_dataset(dataset) @tf.function(input_signature=[dist_dataset.element_spec]) def train_step(inputs): # train model with inputs return # Iterate over the `tf.distribute.DistributedDataset` for x in dist_dataset: # process dataset elements strategy.run(train_step, args=(x,))
Примечание: Порядок обработки данных рабочими узлами при использованииtf.distribute.Strategy.experimental_distribute_datasetилиtf.distribute.Strategy.experimental_distribute_datasets_from_functionне гарантируется. Это обычно требуется, если вы используетеtf.distributeдля масштабирования прогнозирования. Тем не менее, вы можете вставить индекс для каждого элемента в пакет и упорядочить выходы соответственно. Обратитесь к этой ссылке для примера того, как упорядочить выходы.
| Аргументы | |
|---|---|
dataset | tf.data.Dataset, который будет разделен между всеми репликами по вышеуказанным правилам. |
options | tf.distribute.InputOptions для управления параметрами распределения данного набора данных. |
| Возвращаемое значение | |
|---|---|
tf.distribute.DistributedDataset. |
experimental_distribute_datasets_from_function
experimental_distribute_datasets_from_function(
dataset_fn, options=None
)
Распределяет экземпляры tf.data.Dataset, созданные вызовами к dataset_fn.
dataset_fn будет вызван один раз для каждого рабочего узла в стратегии. Каждая реплика на этом рабочем узле будет извлекать одну партию входных данных из локального Dataset (т.е., если у рабочего узла две реплики, две партии будут извлечены из Dataset на каждом шаге).
Этот метод можно использовать для нескольких целей. Например, в случаях, когда experimental_distribute_dataset не может разделить входные файлы, этот метод может быть использован для ручного разделения набора данных (чтобы избежать медленного поведения по умолчанию в experimental_distribute_dataset). В случаях, когда набор данных бесконечен, это разделение можно выполнить, создав копии набора данных, которые отличаются только своим случайным начальным значением. experimental_distribute_dataset также иногда может не удаваться разделить пакет между репликами на рабочем узле. В этом случае этот метод может быть использован, так как в нём нет такого ограничения.
dataset_fn должен принимать экземпляр tf.distribute.InputContext, где можно получить информацию о пакетировании и репликации ввода.
Также можно использовать свойство element_spec объекта tf.distribute.DistributedDataset, возвращаемого этим API, для запроса tf.TypeSpec элементов, возвращаемых итератором. Это можно использовать для установки свойства input_signature объекта tf.function.
global_batch_size = 8
def dataset_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)
strategy = tf.distribute.MirroredStrategy() ds = strategy.experimental_distribute_datasets_from_function(dataset_fn)
def train(ds):
@tf.function(input_signature=[ds.element_spec])
def step_fn(inputs):
# train the model with inputs
return inputs
... for batch in ds: ... replica_results = strategy.run(replica_fn, args=(batch,))
train(ds)
Примечание: Порядок обработки данных рабочими узлами при использованииtf.distribute.Strategy.experimental_distribute_datasetилиtf.distribute.Strategy.experimental_distribute_datasets_from_functionне гарантируется. Это обычно требуется, если вы используетеtf.distributeдля масштабирования предсказания. Однако вы можете вставить индекс для каждого элемента в пакете и упорядочить выходные данные соответственно. См. этот фрагмент для примера того, как упорядочить выходные данные.
| Аргументы | |
|---|---|
dataset_fn | Функция, принимающая экземпляр tf.distribute.InputContext и возвращающая 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()
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>,)
- Распределение значений в массиве на основе идентификатора реплики:
strategy = tf.distribute.MirroredStrategy()
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,)
- Указание значений с использованием num_replicas_in_sync:
strategy = tf.distribute.MirroredStrategy()
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
(1,)
- Размещение значений на устройствах и их распределение:
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(), extended.call_for_each_replica(), или переменная, созданная в scope |
| Возвращаемое значение | |
|---|---|
Кортеж значений, содержащихся в value Если value представляет одно значение, это возвращает (value,). |
experimental_make_numpy_dataset
experimental_make_numpy_dataset(
numpy_input
)
Создаёт tf.data.Dataset из массива NumPy. (устаревший)
Это позволяет избежать добавления numpy_input в виде большой константы в граф и копирует данные на машину или машины, которые будут обрабатывать входные данные.
Обратите внимание, что, вероятно, вам нужно будет использовать experimental_distribute_dataset с возвращённым набором данных, чтобы далее распределить его с помощью стратегии.
Пример:
strategy = tf.distribute.MirroredStrategy() numpy_input = np.ones([10], dtype=np.float32) dataset = strategy.experimental_make_numpy_dataset(numpy_input) dataset <TensorSliceDataset shapes: (), types: tf.float32> dataset = dataset.batch(2) dist_dataset = strategy.experimental_distribute_dataset(dataset)
| Аргументы | |
|---|---|
numpy_input | Вложенный массив массивов NumPy, который будет преобразован в набор данных. Обратите внимание, что массивы NumPy складываются, так как это стандартное поведение tf.data.Dataset. |
| Возвращаемое значение | |
|---|---|
tf.data.Dataset представляющий numpy_input. |
experimental_replicate_to_logical_devices
experimental_replicate_to_logical_devices(
tensor
)
Добавляет аннотацию, что tensor будет дублирован на все логические устройства.
Примечание: Этот API в настоящее время поддерживается только в TPUStrategy. Это добавляет аннотацию к тензоруtensor, указывая, что операции надtensorбудут вызваны на всех логических устройствах.
# Initializing TPU system with 2 logical devices and 4 replicas.
resolver = tf.distribute.cluster_resolver.TPUClusterResolver(tpu='')
tf.config.experimental_connect_to_cluster(resolver)
topology = tf.tpu.experimental.initialize_tpu_system(resolver)
device_assignment = tf.tpu.experimental.DeviceAssignment.build(
topology,
computation_shape=[1, 1, 1, 2],
num_replicas=4)
strategy = tf.distribute.TPUStrategy(
resolver, experimental_device_assignment=device_assignment)
iterator = iter(inputs)
@tf.function()
def step_fn(inputs):
images, labels = inputs
images = strategy.experimental_split_to_logical_devices(
inputs, [1, 2, 4, 1])
# model() function will be executed on 8 logical devices with `inputs`
# split 2 * 4 ways.
output = model(inputs)
# For loss calculation, all logical devices share the same logits
# and labels.
labels = strategy.experimental_replicate_to_logical_devices(labels)
output = strategy.experimental_replicate_to_logical_devices(output)
loss = loss_fn(labels, output)
return loss
strategy.run(step_fn, args=(next(iterator),))
Аргументы: тензор: входной тензор для аннотации.
| Возвращаемое значение | |
|---|---|
Анотированный тензор с идентичным значением, как у tensor. |
experimental_split_to_logical_devices
experimental_split_to_logical_devices(
tensor, partition_dimensions
)
Добавляет аннотацию, что tensor будет разделен на логические устройства.
Примечание: Этот API в настоящее время поддерживается только в TPUStrategy. Это добавляет аннотацию к тензоруtensor, указывая, что операции надtensorбудут разделены между несколькими логическими устройствами. Тензорtensorбудет разделён по измерениям, указанным вpartition_dimensions. Размерыtensorдолжны быть кратны соответствующим значениям вpartition_dimensions.
Например, для системы с 8 логическими устройствами, если tensor - это тензор изображения с формой (размер_пакета, ширина, высота, канал) и partition_dimensions - это [1, 2, 4, 1], то tensor будет разделён на 2 части по ширине и на 4 части по высоте, а значения разделенного тензора будут поданы на 8 логических устройств.
# Initializing TPU system with 8 logical devices and 1 replica.
resolver = tf.distribute.cluster_resolver.TPUClusterResolver(tpu='')
tf.config.experimental_connect_to_cluster(resolver)
topology = tf.tpu.experimental.initialize_tpu_system(resolver)
device_assignment = tf.tpu.experimental.DeviceAssignment.build(
topology,
computation_shape=[1, 2, 2, 2],
num_replicas=1)
strategy = tf.distribute.TPUStrategy(
resolver, experimental_device_assignment=device_assignment)
iterator = iter(inputs)
@tf.function()
def step_fn(inputs):
inputs = strategy.experimental_split_to_logical_devices(
inputs, [1, 2, 4, 1])
# model() function will be executed on 8 logical devices with `inputs`
# split 2 * 4 ways.
output = model(inputs)
return output
strategy.run(step_fn, args=(next(iterator),))
Аргументы: тензор: входной тензор для аннотации. partition_dimensions: список целых чисел без вложенности, размер которого равен рангу tensor, определяющий, как tensor будет разделен. Произведение всех элементов в partition_dimensions должно быть равно общему количеству логических устройств на реплику.
| Исключения | |
|---|---|
ValueError | 1) Если размер partition_dimensions не равен рангу |
| Возвращаемое значение | |
|---|---|
Анотированный тензор с идентичным значением, как у tensor. |
reduce
reduce(
reduce_op, value, axis
)
Уменьшить value по всем репликам.
Учитывая значение на реплику, возвращаемое run, например, потерю на пример, пакет будет разделен по всем репликам. Эта функция позволяет агрегировать по репликам и, необязательно, по элементам пакета. Например, если у вас глобальный размер пакета 8 и 2 реплики, значения для примеров [0, 1, 2, 3] будут на реплике 0, а [4, 5, 6, 7] — на реплике 1. По умолчанию reduce просто агрегирует по репликам, возвращая [0+4, 1+5, 2+6, 3+7]. Это полезно, когда каждая реплика вычисляет скаляр или какое-либо другое значение без измерения «пакет» (например, градиент). Чаще всего вам нужно будет агрегировать по глобальному пакету, что можно сделать, указав измерение пакета как axis, обычно axis=0. В этом случае она вернёт скаляр 0+1+2+3+4+5+6+7.
Если есть последний частичный пакет, вам нужно будет указать ось, чтобы форма результата была согласованной по всем репликам. Таким образом, если последний пакет имеет размер 6 и делится на [0, 1, 2, 3] и [4, 5], у вас возникнет несоответствие формы, если вы не укажете axis=0. Если вы укажете tf.distribute.ReduceOp.MEAN, используя axis=0 будет использоваться правитель литель 6. Противопоставьте это вычислению reduce_mean, чтобы получить скалярное значение на каждой реплике, и этой функции для усреднения этих средних значений, которая будет взвешивать некоторые значения 1/8 и другие 1/4.
| Аргументы | |
|---|---|
reduce_op | Значение tf.distribute.ReduceOp, определяющее, как следует объединять значения. |
value | Значение «по реплике», например, возвращаемое run для объединения в один тензор. |
axis | Указывает измерение для уменьшения вдоль тензора каждой реплики. Обычно следует установить для измерения пакета или None для уменьшения только по репликам (например, если тензор не имеет измерения пакета). |
| Возвращаемое значение | |
|---|---|
Tensor. |
run
run(
fn, args=(), kwargs=None, options=None
)
Выполнить fn на каждой реплике с указанными аргументами.
Выполняет операции, заданные fn на каждой реплике. Если args или kwargs имеют tf.distribute.DistributedValues, например, те, которые создаются tf.distribute.DistributedDataset из tf.distribute.Strategy.experimental_distribute_dataset или tf.distribute.Strategy.experimental_distribute_datasets_from_function, когда fn выполняется на конкретной реплике, она будет выполняться с компонентом tf.distribute.DistributedValues, соответствующим этой реплике.
fn может вызвать tf.distribute.get_replica_context() для доступа к элементам, таким как all_reduce.
Все аргументы в args или kwargs должны быть либо вложены в тензоры, либо tf.distribute.DistributedValues, содержащие тензоры или составные тензоры.
Пример использования:
- Ввод тензора со постоянными значениями.
strategy = tf.distribute.MirroredStrategy() 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 <tf.Tensor: shape=(), dtype=float32, numpy=6.0>
- Вход DistributedValues.
strategy = tf.distribute.MirroredStrategy()
@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=2>
| Аргументы | |
|---|---|
fn | Функция для выполнения. Выход должен быть tf.nest Tensors. |
args | (Необязательно) Позиционные аргументы для fn. |
kwargs | (Необязательно) Именованные аргументы для fn. |
options | (Необязательно) Экземпляр tf.distribute.RunOptions, определяющий параметры выполнения fn. |
| Возвращаемое значение | |
|---|---|
Объединённое возвращаемое значение fn по репликам. Структура возвращаемого значения такая же, как и возвращаемого значение от fn. Каждый элемент структуры может быть tf.distribute.DistributedValues, объектами Tensor, или Tensor (например, при выполнении на одной реплике). |
scope
scope()
Возвращает контекстный менеджер, выбирающий эту стратегию в качестве текущей.
Внутри блока кода with strategy.scope():, этот поток будет использовать создатель переменных, заданный strategy, и войдёт в свой «межрепликационный контекст».
В MultiWorkerMirroredStrategy, все переменные, созданные внутри `strategy.scope()`, будут дублироваться на всех репликах каждого рабочего узла. Кроме того, она также устанавливает область действия устройства по умолчанию, поэтому операции без указанных устройств будут оказываться на правильном рабочем узле.
| Возвращаемое значение | |
|---|---|
| Контекстный менеджер для создания переменных с этой стратегией. |
© 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.3/api_docs/python/tf/distribute/experimental/MultiWorkerMirroredStrategy