tf.distribute.TPUStrategy
Синхронное обучение на TPUs и TPU-кластерах.
Наследуется от: Strategy
tf.distribute.TPUStrategy(
tpu_cluster_resolver=None,
experimental_device_assignment=None,
experimental_spmd_xla_partitioning=False
)
Используется в ноутбуках
| Используется в руководстве | Используется в учебниках |
|---|---|
Для создания объекта TPUStrategy необходимо выполнить код инициализации, как показано ниже:
resolver = tf.distribute.cluster_resolver.TPUClusterResolver(tpu='') tf.config.experimental_connect_to_cluster(resolver) tf.tpu.experimental.initialize_tpu_system(resolver) strategy = tf.distribute.TPUStrategy(resolver)
При использовании стратегий распределения переменные, созданные в рамках области действия стратегии, будут дублироваться на всех репликах и могут поддерживаться синхронизированными с помощью алгоритмов all-reduce.
Для запуска программ TF2 на TPUs можно использовать .compile и .fit API в tf.keras со стратегией TPUStrategy или написать собственную настраиваемую петлю обучения, вызвав strategy.run напрямую. Обратите внимание, что TPUStrategy не поддерживает чистое выполнение eager, поэтому убедитесь, что функция, переданная в strategy.run, является tf.function или что strategy.run вызывается внутри tf.function, если включено eager-поведение. Дополнительные сведения см. в https://www.tensorflow.org/guide/tpu.
distribute_datasets_from_function и experimental_distribute_dataset API могут использоваться для распределения набора данных по TPU-узлам при написании собственной петли обучения. Если вы используете fit и compile методы, доступные в tf.keras.Model, Keras будет обрабатывать распределение за вас.
Пример написания настраиваемой петли обучения на TPUs:
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)
iterator = iter(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),))train_step(iterator)
Для расширенных случаев использования, таких как параллелизм моделей, можно установить аргумент experimental_device_assignment при создании TPUStrategy, чтобы указать количество реплик и количество логических устройств. Ниже приведен пример инициализации TPU-системы с 2 логическими устройствами и 1 репликой.
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=1)
strategy = tf.distribute.TPUStrategy(
resolver, experimental_device_assignment=device_assignment)Затем можно запустить операцию tf.add только на логическом устройстве 0.
@tf.function()
def step_fn(inputs):
features, _ = inputs
output = tf.add(features, features)
# Add operation will be executed on logical device 0.
output = strategy.experimental_assign_to_logical_device(output, 0)
return output
dist_dataset = strategy.distribute_datasets_from_function(
dataset_fn)
iterator = iter(dist_dataset)
strategy.run(step_fn, args=(next(iterator),))experimental_spmd_xla_partitioning включает экспериментальную функцию XLA SPMD для параллелизма моделей. Этот флаг может сократить время компиляции и требования к HBM. При работе в этом режиме каждый тензор входных данных должен быть либо разбит (с помощью strategy.experimental_split_to_logical_devices), либо полностью дублирован (с помощью strategy.experimental_replicate_to_logical_devices) на всех логических устройствах. И вызов strategy.experimental_assign_to_logical_device приведет к ошибке ValueError в этом режиме.
| Аргументы | |
|---|---|
tpu_cluster_resolver | Экземпляр tf.distribute.cluster_resolver.TPUClusterResolver, который предоставляет информацию о кластере TPU. Если None, предполагается, что он работает на локальном TPU-узле. |
experimental_device_assignment | Необязательный tf.tpu.experimental.DeviceAssignment для указания размещения реплик в кластере TPU. |
experimental_spmd_xla_partitioning | Если True, включает режим SPMD (Single Program Multiple Data) в компиляторе XLA. Этот флаг влияет только на производительность компиляции XLA и требования к HBM скомпилированной программы TPU. Предупреждение: если этот флаг установлен в True, вызов tf.distribute.TPUStrategy.experimental_assign_to_logical_device приведет к ошибке ValueError. |
| Атрибуты | |
|---|---|
cluster_resolver | Возвращает решатель кластера, связанный с этой стратегией.
|
extended | tf.distribute.StrategyExtended с дополнительными методами. |
num_replicas_in_sync | Возвращает количество реплик, по которым агрегируются градиенты. |
Методы
distribute_datasets_from_function
distribute_datasets_from_function(
dataset_fn, options=None
)
Распределяет экземпляры tf.data.Dataset, созданные при вызовах dataset_fn.
Аргумент dataset_fn, который передаёт пользователь, является функцией ввода, которая имеет аргумент tf.distribute.InputContext и возвращает экземпляр tf.data.Dataset. Ожидается, что возвращаемый набор данных из dataset_fn уже сгруппирован по размеру пакета на реплику (т. е. глобальный размер пакета, деленный на количество реплик в синхронизации) и разделен. tf.distribute.Strategy.distribute_datasets_from_function не группирует и не разделяет экземпляр tf.data.Dataset, возвращаемый функцией ввода. dataset_fn будет вызываться на устройстве 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-процесса.
Для получения дополнительных примеров использования и свойств этого метода обратитесь к учебнику по распределённому вводу. Если вы интересуетесь обработкой последней частичной партии, прочитайте эту секцию.
END_OF_DOCUMENT_MARKER| Аргументы | |
|---|---|
dataset_fn | Функция, принимающая экземпляр tf.distribute.InputContext и возвращающая tf.data.Dataset. |
options | tf.distribute.InputOptions, используемые для управления параметрами распределения этого набора данных. |
| Возвращаемое значение | |
|---|---|
tf.distribute.DistributedDataset. |
experimental_assign_to_logical_device
experimental_assign_to_logical_device(
tensor, logical_device_id
)
Добавляет аннотацию, что tensor будет назначен логическому устройству.
Это добавляет аннотацию к 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 | Представленный идентификатор логического устройства не соответствует общему количеству разделов, указанных в назначении устройства, или TPUStrategy был создан с использованием experimental_spmd_xla_partitioning=True. |
| Возвращаемое значение | |
|---|---|
Анотированный тензор с идентичным значением, что и tensor. |
experimental_distribute_dataset
experimental_distribute_dataset(
dataset, options=None
)
Создаёт tf.distribute.DistributedDataset из tf.data.Dataset.
Возвращённый tf.distribute.DistributedDataset можно перебирать, как и обычные наборы данных. ВАЖНО: пользователь не может добавить дополнительные преобразования в tf.distribute.DistributedDataset. Вы можете только создать итератор или проверить tf.TypeSpec генерируемых данных. Подробнее см. документацию API tf.distribute.DistributedDataset.
Вот пример:
global_batch_size = 2
# Passing the devices is optional.
strategy = tf.distribute.MirroredStrategy(devices=["GPU:0", "GPU:1"])
# Create a dataset
dataset = tf.data.Dataset.range(4).batch(global_batch_size)
# Distribute that dataset
dist_dataset = strategy.experimental_distribute_dataset(dataset)
@tf.function
def replica_fn(input):
return input*2
result = []
# Iterate over the `tf.distribute.DistributedDataset`
for x in dist_dataset:
# process dataset elements
result.append(strategy.run(replica_fn, args=(x,)))
print(result)
[PerReplica:{
0: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([0])>,
1: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([2])>
}, PerReplica:{
0: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([4])>,
1: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([6])>
}]Три ключевых действия, происходящие за кулисами этого метода, это пакетное формирование, фрагментация и предварительная обработка.
В приведённом фрагменте кода dataset формируется пакетами размером global_batch_size, и вызов experimental_distribute_dataset на нём перегруппировывает dataset в новый размер пакета, равный общему размеру пакета, делённому на количество реплик в синхронизации. Мы перебираем его с помощью цикла for в стиле Python. x - это tf.distribute.DistributedValues, содержащий данные для всех реплик, и каждая реплика получает данные нового размера пакета. tf.distribute.Strategy.run будет обрабатывать передачу правильных данных для каждой реплики в x в соответствующую replica_fn, выполняемую на каждой реплике.
Фрагментация включает в себя автоматическую фрагментацию по нескольким узлам и внутри каждого узла. Во-первых, при распределённом обучении по нескольким узлам (то есть при использовании tf.distribute.experimental.MultiWorkerMirroredStrategy или tf.distribute.TPUStrategy) автоматическая фрагментация набора данных по узлам означает, что каждому узлу назначается подмножество всего набора данных (если установлен правильный tf.data.experimental.AutoShardPolicy). Это гарантирует, что на каждом шаге каждый узел обработает глобальный пакет элементов набора данных без пересечений. Автоматическая фрагментация имеет несколько разных вариантов, которые можно указать с помощью tf.data.experimental.DistributeOptions. Затем фрагментация внутри каждого узла означает, что метод разделит данные между всеми устройствами узла (если их несколько). Это произойдёт независимо от автоматической фрагментации по нескольким узлам.
Примечание: по умолчанию режим автоматической фрагментации по нескольким узлам -tf.data.experimental.AutoShardPolicy.AUTO. Этот режим попытается фрагментировать входной набор данных по файлам, если набор данных создаётся из наборов данных, читаемых из файлов (например,tf.data.TFRecordDataset,tf.data.TextLineDatasetи т. д.), или фрагментировать набор данных по данным, где каждый из узлов прочитает весь набор данных и обработает только назначенный ему фрагмент. Однако, если у вас меньше одного входного файла на узел, рекомендуется отключить автоматическую фрагментацию набора данных по узлам, установивtf.data.experimental.DistributeOptions.auto_shard_policyв значениеtf.data.experimental.AutoShardPolicy.OFF.
По умолчанию этот метод добавляет преобразование предварительной обработки в конце экземпляра входного набора данных, предоставленного пользователем tf.data.Dataset. Аргумент для преобразования предварительной обработки, который является buffer_size, равен количеству реплик в синхронизации.
Если описанная выше логика разделения пакетов и фрагментации набора данных нежелательна, используйте вместо этого tf.distribute.Strategy.distribute_datasets_from_function, который не выполняет автоматической пакетной обработки или фрагментации.
Примечание: Если вы используете TPUStrategy, порядок обработки данных узлами при использованииtf.distribute.Strategy.experimental_distribute_datasetилиtf.distribute.Strategy.distribute_datasets_from_functionне гарантируется. Это обычно требуется, если вы используетеtf.distributeдля масштабирования прогнозирования. Однако вы можете вставить индекс для каждого элемента в пакете и упорядочить выводы соответственно. См. этот фрагмент для примера того, как упорядочить выводы.
Примечание: Трансформации состояний набора данных в настоящее время не поддерживаются сtf.distribute.experimental_distribute_datasetилиtf.distribute.distribute_datasets_from_function. Любые операторы состояния, которые может иметь набор данных, в настоящее время игнорируются. Например, если ваш набор данных имеет операциюmap_fn, которая используетtf.random.uniformдля поворота изображения, то у вас есть граф набора данных, зависящий от состояния (т. е. от начального значения случайного числа) на локальной машине, где выполняется процесс Python.
Дополнительную информацию об использовании и свойствах этого метода можно найти в руководстве по распределённому вводу. Если вас интересует обработка последнего частичного пакета, прочтите этот раздел .
| Аргументы | |
|---|---|
dataset | tf.data.Dataset, который будет распределён по всем репликам в соответствии с вышеуказанными правилами. |
options | tf.distribute.InputOptions, используемые для управления параметрами распределения этого набора данных. |
| Возвращаемое значение | |
|---|---|
tf.distribute.DistributedDataset. |
experimental_distribute_values_from_function
experimental_distribute_values_from_function(
value_fn
)
Генерирует tf.distribute.DistributedValues из value_fn.
Эта функция генерирует tf.distribute.DistributedValues для передачи в run, reduce или другие методы, принимающие распределённые значения, когда не используются наборы данных.
| Аргументы | |
|---|---|
value_fn | Функция, выполняемая для генерации значений. Она вызывается для каждой реплики с tf.distribute.ValueContext в качестве единственного аргумента. Она должна возвращать тензор или тип, который может быть преобразован в тензор. |
| Возвращаемое значение | |
|---|---|
tf.distribute.DistributedValues, содержащий значение для каждой реплики. |
Пример использования:
-
Возвращение постоянного значения для каждой реплики:
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"]) def value_fn(ctx): return tf.constant(1.) distributed_values = ( strategy.experimental_distribute_values_from_function( value_fn)) local_result = strategy.experimental_local_results( distributed_values) local_result (<tf.Tensor: shape=(), dtype=float32, numpy=1.0>, <tf.Tensor: shape=(), dtype=float32, numpy=1.0>) -
Распределение значений массива в зависимости от идентификатора реплики:
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,). |
experimental_replicate_to_logical_devices
experimental_replicate_to_logical_devices(
tensor
)
Добавляет аннотацию, что tensor будет дублирован на все логические устройства.
Это добавляет аннотацию к тензору 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 — входной тензор для аннотации.
| Возвращаемое значение | |
|---|---|
Аннотированный тензор с идентичным значением, как tensor. |
experimental_split_to_logical_devices
experimental_split_to_logical_devices(
tensor, partition_dimensions
)
Добавляет аннотацию, что tensor будет разделен между логическими устройствами.
Это добавляет аннотацию к тензору 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)
# Construct the TPUStrategy. Since we are going to split the image across
# logical devices, here we set `experimental_spmd_xla_partitioning=True`
# so that the partitioning can be compiled in SPMD mode, which usually
# results in faster compilation and smaller HBM requirement if the size of
# input and activation tensors are much bigger than that of the model
# parameters. Note that this flag is suggested but not a hard requirement
# for `experimental_split_to_logical_devices`.
strategy = tf.distribute.TPUStrategy(
resolver, experimental_device_assignment=device_assignment,
experimental_spmd_xla_partitioning=True)
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),))
Аргументы: tensor — входной тензор для аннотации. partition_dimensions — список целых чисел без вложенности, равный размеру ранга tensor, определяющий, как будет разделяться tensor. Произведение всех элементов в partition_dimensions должно быть равно общему количеству логических устройств на реплику.
| Исключения | |
|---|---|
ValueError | 1) Если размер partition_dimensions не равен рангу |
| Возвращаемое значение | |
|---|---|
Аннотированный тензор с идентичным значением, как tensor. |
gather
gather(
value, axis
)
Собрать value по репликам вдоль axis на текущее устройство.
Учитывая объект типа tf.distribute.DistributedValues или tf.Tensor value, этот API собирает и конкатенирует value по репликам вдоль axis-й размерности. Результат копируется на "текущее" устройство, которое обычно является процессором на рабочем узле, на котором выполняется программа. Для tf.distribute.TPUStrategy — это первый хост TPU. Для многоклиентской tf.distribute.MultiWorkerMirroredStrategy — это процессор каждого рабочего узла.
Этот API можно вызывать только в контексте между репликами. Для аналога в контексте реплики см. tf.distribute.ReplicaContext.all_gather.
Примечание: Для всех стратегий, кромеtf.distribute.TPUStrategy, входvalueна разных репликах должен иметь одинаковый ранг, а их формы должны быть одинаковыми во всех измерениях, кромеaxis-й размерности. Другими словами, их формы не могут отличаться по размеру в измеренииd, гдеdне равно аргументуaxis. Например, учитываяtf.distribute.DistributedValuesс компонентами тензоров формы(1, 2, 3)и(1, 3, 3)на двух репликах, вы можете вызватьgather(..., axis=1, ...)на нём, но неgather(..., axis=0, ...)илиgather(..., axis=2, ...). Однако дляtf.distribute.TPUStrategy.gatherвсе тензоры должны иметь точно такой же ранг и форму.
Примечание: Учитываяtf.distribute.DistributedValuesvalue, его компоненты тензоров должны иметь ненулевой ранг. В противном случае рассмотрите возможность использованияtf.expand_dimsперед их сбором.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
# A DistributedValues with component tensor of shape (2, 1) on each replica
distributed_values = strategy.experimental_distribute_values_from_function(lambda _: tf.identity(tf.constant([[1], [2]])))
@tf.function
def run():
return strategy.gather(distributed_values, axis=0)
run()
<tf.Tensor: shape=(4, 1), dtype=int32, numpy=
array([[1],
[2],
[1],
[2]], dtype=int32)>Рассмотрим следующий пример для получения дополнительных сочетаний:
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1", "GPU:2", "GPU:3"])
single_tensor = tf.reshape(tf.range(6), shape=(1,2,3))
distributed_values = strategy.experimental_distribute_values_from_function(lambda _: tf.identity(single_tensor))
@tf.function
def run(axis):
return strategy.gather(distributed_values, axis=axis)
axis=0
run(axis)
<tf.Tensor: shape=(4, 2, 3), dtype=int32, numpy=
array([[[0, 1, 2],
[3, 4, 5]],
[[0, 1, 2],
[3, 4, 5]],
[[0, 1, 2],
[3, 4, 5]],
[[0, 1, 2],
[3, 4, 5]]], dtype=int32)>
axis=1
run(axis)
<tf.Tensor: shape=(1, 8, 3), dtype=int32, numpy=
array([[[0, 1, 2],
[3, 4, 5],
[0, 1, 2],
[3, 4, 5],
[0, 1, 2],
[3, 4, 5],
[0, 1, 2],
[3, 4, 5]]], dtype=int32)>
axis=2
run(axis)
<tf.Tensor: shape=(1, 2, 12), dtype=int32, numpy=
array([[[0, 1, 2, 0, 1, 2, 0, 1, 2, 0, 1, 2],
[3, 4, 5, 3, 4, 5, 3, 4, 5, 3, 4, 5]]], dtype=int32)>| Аргументы | |
|---|---|
value | Объект tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который следует объединить в один тензор. Он также может быть обычным тензором при использовании с tf.distribute.OneDeviceStrategy или стратегией по умолчанию. Тензоры, составляющие DistributedValues, могут быть только плоскими тензорами с ненулевым рангом, а НЕ tf.IndexedSlices. |
axis | 0-мерный тензор int32. Размерность, по которой выполнять сбор. Должно быть в диапазоне [0, ранг(значение)). |
| Возвращаемое значение | |
|---|---|
Объект Tensor, являющийся конкатенацией value по репликам вдоль axis измерения. |
reduce
reduce(
reduce_op, value, axis
)
Свести value по репликам и вернуть результат на текущем устройстве.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
def step_fn():
i = tf.distribute.get_replica_context().replica_id_in_sync_group
return tf.identity(i)
per_replica_result = strategy.run(step_fn)
total = strategy.reduce("SUM", per_replica_result, axis=None)
total
<tf.Tensor: shape=(), dtype=int32, numpy=1>Чтобы увидеть, как это будет выглядеть с несколькими репликами, рассмотрите тот же пример с MirroredStrategy с 2 графическими процессорами:
strategy = tf.distribute.MirroredStrategy(devices=["GPU:0", "GPU:1"])
def step_fn():
i = tf.distribute.get_replica_context().replica_id_in_sync_group
return tf.identity(i)
per_replica_result = strategy.run(step_fn)
# Check devices on which per replica result is:
strategy.experimental_local_results(per_replica_result)[0].device
# /job:localhost/replica:0/task:0/device:GPU:0
strategy.experimental_local_results(per_replica_result)[1].device
# /job:localhost/replica:0/task:0/device:GPU:1
total = strategy.reduce("SUM", per_replica_result, axis=None)
# Check device on which reduced result is:
total.device
# /job:localhost/replica:0/task:0/device:CPU:0
Этот API обычно используется для агрегирования результатов, возвращаемых из разных реплик, для отчётности и т. п. Например, потерю, вычисленную из разных реплик, можно усреднить с помощью этого API перед печатью.
Примечание: Результат копируется на "текущее" устройство — это обычно процессор рабочего узла, на котором выполняется программа. ДляTPUStrategy— это первый хост TPU. Для многоклиентской стратегииMultiWorkerMirroredStrategy— это процессор каждого рабочего узла.
Существует ряд различных API tf.distribute для сведения значений по репликам:
-
tf.distribute.ReplicaContext.all_reduce: Это отличается отStrategy.reduceтем, что предназначено для контекста реплик и не копирует результаты на устройство хоста.all_reduceобычно используется для сокращений внутри шага обучения, таких как градиенты. -
tf.distribute.StrategyExtended.reduce_toиtf.distribute.StrategyExtended.batch_reduce_to: Эти API являются более продвинутыми версиямиStrategy.reduce, так как они позволяют настраивать место назначения результата. Они также вызываются в контексте между репликами.
Каким должен быть ось?
Учитывая значение на реплику, возвращаемое run, например, потерю на пример, пакет будет разделен между всеми репликами. Эта функция позволяет вам агрегировать по репликам и, по желанию, также по элементам пакета, указав параметр оси соответственно.
Например, если у вас есть глобальный размер пакета 8 и 2 реплики, значения для примеров [0, 1, 2, 3] будут на реплике 0, а [4, 5, 6, 7] — на реплике 1. С axis=None, reduce будет агрегировать только по репликам, возвращая [0+4, 1+5, 2+6, 3+7]. Это полезно, когда каждая реплика вычисляет скаляр или какое-то другое значение, у которого нет «размерности пакета» (например, градиент или потеря).
strategy.reduce("sum", per_replica_result, axis=None)
Иногда вы захотите агрегировать как по глобальному пакету, так и по всем репликам. Вы можете получить это поведение, указав размерность пакета как axis, обычно axis=0. В этом случае он вернёт скаляр 0+1+2+3+4+5+6+7.
strategy.reduce("sum", per_replica_result, axis=0)
Если есть последняя частичная партия, вам необходимо указать ось, чтобы результирующая форма была согласованной на всех репликах. Таким образом, если последняя партия имеет размер 6 и разделена на [0, 1, 2, 3] и [4, 5], вы получите несовпадение форм, если не укажете axis=0. Если вы укажете tf.distribute.ReduceOp.MEAN, используя axis=0 будет использован правильный знаменатель 6. Противопоставьте это вычислению reduce_mean для получения скалярного значения на каждой реплике и этой функции для усреднения этих средних значений, которая будет взвешивать некоторые значения 1/8, а другие 1/4.
| Аргументы | |
|---|---|
reduce_op | значение tf.distribute.ReduceOp, определяющее, как следует комбинировать значения. Разрешает использовать строковое представление перечисления, например "SUM", "MEAN". |
value | экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который необходимо объединить в один тензор. Он также может быть обычным тензором, когда используется с OneDeviceStrategy или стратегией по умолчанию. |
axis | определяет размерность, по которой следует производить сокращение внутри тензора каждой реплики. Обычно следует устанавливать на размерность пакетной обработки или None, чтобы сократить только по репликам (например, если тензор не имеет размерности пакетной обработки). |
| Возвращаемое значение | |
|---|---|
Tensor. |
run
run(
fn, args=(), kwargs=None, options=None
)
Выполните вычисление, определенное fn на каждой реплике TPU.
Выполняет операции, указанные в 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 может вызвать tf.distribute.get_replica_context() для доступа к элементам, таким как all_reduce.
Все аргументы в args или kwargs должны быть либо вложенными тензорами, либо tf.distribute.DistributedValues, содержащими тензоры или составные тензоры.
Пример использования:
resolver = tf.distribute.cluster_resolver.TPUClusterResolver(tpu='')
tf.config.experimental_connect_to_cluster(resolver)
tf.tpu.experimental.initialize_tpu_system(resolver)
strategy = tf.distribute.TPUStrategy(resolver)
@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_fn(input):
return input * 2
return strategy.run(replica_fn, args=(distributed_values,))
result = run()| Аргументы | |
|---|---|
fn | Функция для выполнения. Выход должен быть tf.nest из Tensor. |
args | (Необязательно) Позиционные аргументы для fn. |
kwargs | (Необязательно) Именованные аргументы для fn. |
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илиkeras.Model.fit, чтобы он был автоматически для вас. Любые переменные, созданные за пределами блока, не будут распределены и могут иметь последствия для производительности. Некоторые распространенные объекты, создающие переменные в TF, это модели, оптимизаторы, метрики. Такие объекты всегда должны инициализироваться в блоке, и любые функции, которые могут создавать переменные лениво (например,Model.call(), отслеживаяtf.functionи т.д.), также должны вызываться внутри блока. Другим источником создания переменных может быть восстановление контрольной точки - когда переменные создаются лениво. Обратите внимание, что любая переменная, созданная внутри стратегии, сохраняет информацию о стратегии. Таким образом, чтение и запись в эти переменные за пределамиstrategy.scopeтакже могут работать без проблем, без необходимости ввода пользователем блока. - Некоторые API стратегии (такие как
strategy.runиstrategy.reduce), которые требуют находиться в области действия стратегии, автоматически вводят область действия, что означает, что при использовании этих API вам не нужно явным образом вводить область действия. - Когда
tf.keras.Modelсоздается внутриstrategy.scope, объект Model захватывает информацию о области действия. Когда вызываются методы высокоуровневой обучающей рамки, такие какmodel.compile,model.fitи т.д., захваченная область действия будет автоматически введена, и используемая стратегия будет использоваться для распределения обучения и т.д. Смотрите подробный пример в tutorial по распределенному keras. ПРЕДУПРЕЖДЕНИЕ: Просто вызовmodel(..)не автоматически вводит захваченную область действия - только API высокоуровневых обучающих рамок поддерживают это поведение:model.compile,model.fit,model.evaluate,model.predictиmodel.saveмогут быть вызваны внутри или вне области действия. - Следующее может быть как внутри, так и вне области действия:
- Создание наборов данных для входа
- Определение
tf.function, представляющих ваш шаг обучения - API сохранения, такие как
tf.saved_model.save. Загрузка создает переменные, поэтому это должно происходить внутри области действия, если вы хотите обучить модель в распределенном режиме. - Сохранение контрольных точек. Как упоминалось выше -
checkpoint.restoreиногда может потребоваться внутри области действия, если оно создает переменные.
| Возвращаемое значение | |
|---|---|
| Блок кода. |
© 2022 The TensorFlow Authors. All rights reserved.
Licensed under the Creative Commons Attribution License 4.0.
Code samples licensed under the Apache 2.0 License.
https://www.tensorflow.org/api_docs/python/tf/distribute/TPUStrategy