tf.distribute.TPUStrategy
Синхронное обучение на TPUs и TPU Pods.
Наследуется от: Strategy
tf.distribute.TPUStrategy(
tpu_cluster_resolver=None, experimental_device_assignment=None
)
Для создания объекта 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 можно использовать API .compile и .fit в 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),))
| Аргументы | |
|---|---|
tpu_cluster_resolver | tf.distribute.cluster_resolver.TPUClusterResolver, который предоставляет информацию о кластере TPU. Если None, предполагается запуск на локальном TPU-узле. |
experimental_device_assignment | Необязательный tf.tpu.experimental.DeviceAssignment для указания размещения реплик в кластере TPU. |
| Атрибуты | |
|---|---|
cluster_resolver | Возвращает решатель кластера, связанный с этой стратегией. В общем случае, при использовании стратегии распределения многоузлового Стратегии, которые намереваются иметь связанный Стратегии одноузлового типа обычно не имеют Решатель кластера
os.environ['TF_CONFIG'] = json.dumps({
'cluster': {
'worker': ["localhost:12345", "localhost:23456"],
'ps': ["localhost:34567"]
},
'task': {'type': 'worker', 'index': 0}
})
# This implicitly uses TF_CONFIG for the cluster and current task info.
strategy = tf.distribute.experimental.MultiWorkerMirroredStrategy()
...
if strategy.cluster_resolver.task_type == 'worker':
# Perform something that's only applicable on workers. Since we set this
# as a worker above, this block will run on this particular instance.
elif strategy.cluster_resolver.task_type == 'ps':
# Perform something that's only applicable on parameter servers. Since we
# set this as a worker above, this block will not run on this particular
# instance.
Для получения дополнительной информации, обратитесь к документации API |
extended | tf.distribute.StrategyExtended с дополнительными методами. |
num_replicas_in_sync | Возвращает количество реплик, по которым агрегируются градиенты. |
Методы
distribute_datasets_from_function
distribute_datasets_from_function(
dataset_fn, options=None
)
Распределяет экземпляры tf.data.Dataset, созданные вызовами к dataset_fn.
Аргумент dataset_fn , который пользователи передают, является функцией ввода, у которой есть аргумент tf.distribute.InputContext, и она возвращает экземпляр tf.data.Dataset. Ожидается, что возвращаемый из dataset_fn набор данных уже разбит на пакеты по размеру пакета на реплику (т. е. глобальный размер пакета, деленный на количество реплик в синхронизации) и фрагментирован. tf.distribute.Strategy.distribute_datasets_from_function не разбивает и не фрагментирует экземпляр tf.data.Dataset, возвращаемый функцией ввода. dataset_fn будет вызываться на процессорном устройстве каждого из рабочих узлов, и каждый генерирует набор данных, где каждая реплика на этом рабочем узле будет извлекать по одной партии входных данных (т. е. если на рабочем узле есть две реплики, две партии будут извлекаться из Dataset на каждом шаге).
Этот метод может быть использован для нескольких целей. Во-первых, он позволяет вам указать собственную логику разбивки на пакеты и фрагментацию. (В отличие от tf.distribute.experimental_distribute_dataset, который выполняет разбивку на пакеты и фрагментацию за вас.) Например, в тех случаях, когда experimental_distribute_dataset не может фрагментировать входные файлы, этот метод может быть использован для ручного фрагментирования набора данных (чтобы избежать медленной обратной совместимости в experimental_distribute_dataset). В случаях, когда набор данных бесконечен, эта фрагментация может быть выполнена путем создания реплик набора данных, которые различаются только своим случайным зерном.
Функция dataset_fn должна принимать экземпляр tf.distribute.InputContext, где можно получить доступ к информации о разбивке на пакеты и репликации ввода.
Вы можете использовать свойство element_spec возвращаемого этим API tf.distribute.DistributedDataset, чтобы запросить tf.TypeSpec элементов, возвращаемых итератором. Это можно использовать для установки свойства input_signature функции tf.function. Следуйте tf.distribute.DistributedDataset.element_spec, чтобы увидеть пример.
Примечание: При использовании TPUStrategy порядок обработки данных рабочими узлами при использованииtf.distribute.Strategy.experimental_distribute_datasetилиtf.distribute.Strategy.distribute_datasets_from_functionне гарантируется. Это обычно необходимо, если вы используетеtf.distributeдля масштабирования прогнозирования. Тем не менее, вы можете вставить индекс для каждого элемента в пачке и упорядочить выводы соответственно. Обратитесь к этому фрагменту для примера того, как упорядочить выводы.
Примечание: Состоятельные преобразования наборов данных в настоящее время не поддерживаются сtf.distribute.experimental_distribute_datasetилиtf.distribute.distribute_datasets_from_function. Любые состоятельные операции, которые может иметь набор данных, в настоящее время игнорируются. Например, если ваш набор данных содержитmap_fn, использующийtf.random.uniformдля вращения изображения, то у вас есть граф набора данных, который зависит от состояния (т. е. случайного зерна) на локальной машине, где выполняется процесс Python.
Дополнительные сведения о методе и его свойствах см. в учебнике по распределенному вводу. Если вас интересует обработка последней частичной партии, прочитайте эту часть.
| Аргументы | |
|---|---|
dataset_fn | Функция, принимающая экземпляр tf.distribute.InputContext и возвращающая tf.data.Dataset. |
options | tf.distribute.InputOptions, используемый для управления параметрами распределения этого набора данных. |
| Возвращаемое значение | |
|---|---|
tf.distribute.DistributedDataset. |
experimental_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 | Идентификатор логического устройства не соответствует общему количеству разделов, указанных в задании устройства. |
| Возвращаемое значение | |
|---|---|
Тензор с аннотациями, имеющий то же значение, что и 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(), extended.call_for_each_replica(), или переменная, созданная в scope . |
| Возвращает | |
|---|---|
Кортеж значений, содержащихся в value . Если 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 . |
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)
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. |
gather
gather(
value, axis
)
Сборка value по репликам вдоль axis на текущее устройство.
Учитывая объект типа tf.distribute.DistributedValues или похожий на tf.Tensor value, этот API собирает и конкатенирует value по репликам вдоль axis-й размерности. Результат копируется на "текущее" устройство
- что обычно является процессором узла, на котором выполняется программа. Для
tf.distribute.TPUStrategyэто первый хост TPU. Для многоклиентскихMultiWorkerMirroredStrategy, это процессор каждого узла.
Этот API может быть вызван только в контексте между репликами. Для аналога в контексте реплики см. tf.distribute.ReplicaContext.all_gather.
Примечание: Для всех стратегий, кромеtf.distribute.TPUStrategy, входнойvalueна разных репликах должен иметь одинаковый ранг, и их формы должны быть одинаковыми во всех измерениях, кромеaxis-й размерности. Другими словами, их формы не могут отличаться по размерностиd, гдеdне равно аргументуaxis. Например, дляtf.distribute.DistributedValuesс тензорами компонентов формы(1, 2, 3)и(1, 3, 3)на двух репликах вы можете вызватьgather(..., axis=1, ...), но неgather(..., axis=0, ...)илиgather(..., axis=2, ...). Однако дляtf.distribute.TPUStrategy.gather, все тензоры должны иметь точно такой же ранг и форму.
Примечание: Учитываяtf.distribute.DistributedValuesvalue, её компонентные тензоры должны иметь ранг не равный нулю. В противном случае, рассмотрите использованиеtf.expand_dimsперед их сбором.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
# A DistributedValues with component tensor of shape (2, 1) on each replica
distributed_values = strategy.experimental_distribute_values_from_function(lambda _: tf.identity(tf.constant([[1], [2]])))
@tf.function
def run():
return strategy.gather(distributed_values, axis=0)
run()
<tf.Tensor: shape=(4, 1), dtype=int32, numpy=
array([[1],
[2],
[1],
[2]], dtype=int32)>
Рассмотрите следующий пример для других комбинаций:
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1", "GPU:2", "GPU:3"])
single_tensor = tf.reshape(tf.range(6), shape=(1,2,3))
distributed_values = strategy.experimental_distribute_values_from_function(lambda _: tf.identity(single_tensor))
@tf.function
def run(axis):
return strategy.gather(distributed_values, axis=axis)
axis=0
run(axis)
<tf.Tensor: shape=(4, 2, 3), dtype=int32, numpy=
array([[[0, 1, 2],
[3, 4, 5]],
[[0, 1, 2],
[3, 4, 5]],
[[0, 1, 2],
[3, 4, 5]],
[[0, 1, 2],
[3, 4, 5]]], dtype=int32)>
axis=1
run(axis)
<tf.Tensor: shape=(1, 8, 3), dtype=int32, numpy=
array([[[0, 1, 2],
[3, 4, 5],
[0, 1, 2],
[3, 4, 5],
[0, 1, 2],
[3, 4, 5],
[0, 1, 2],
[3, 4, 5]]], dtype=int32)>
axis=2
run(axis)
<tf.Tensor: shape=(1, 2, 12), dtype=int32, numpy=
array([[[0, 1, 2, 0, 1, 2, 0, 1, 2, 0, 1, 2],
[3, 4, 5, 3, 4, 5, 3, 4, 5, 3, 4, 5]]], dtype=int32)>
| Аргументы | |
|---|---|
value | экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который необходимо объединить в один тензор. Он также может быть обычным тензором, когда используется с tf.distribute.OneDeviceStrategy или по умолчанию. Тензоры, составляющие DistributedValues, могут быть только плотными тензорами с ненулевым рангом, НЕ tf.IndexedSlices. |
axis | 0-мерный тензор типа int32. Измерение, по которому выполняется сбор. Должен быть в диапазоне [0, rank(value)). |
| Возвращает | |
|---|---|
А Tensor, который представляет собой конкатенацию value по репликам вдоль axis размерности. |
reduce
reduce(
reduce_op, value, axis
)
Сведение value по репликам и возврат результата на текущее устройство.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
def step_fn():
i = tf.distribute.get_replica_context().replica_id_in_sync_group
return tf.identity(i)
per_replica_result = strategy.run(step_fn)
total = strategy.reduce("SUM", per_replica_result, axis=None)
total
<tf.Tensor: shape=(), dtype=int32, numpy=1>
Чтобы увидеть, как это будет выглядеть с несколькими репликами, рассмотрите тот же пример с MirroredStrategy с 2 GPU:
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 | Функция для выполнения. Выход должен быть вложен из 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илиmodel.fit, для входа в него за вас. Любая переменная, созданная за пределами области действия, не будет распределена и может привести к проблемам производительности. Типичные вещи, создающие переменные в TF: модели, оптимизаторы, метрики. Они всегда должны создаваться внутри области действия. Еще одним источником создания переменных может быть восстановление контрольной точки - когда переменные создаются лениво. Обратите внимание, что любая переменная, созданная внутри стратегии, фиксирует информацию о стратегии. Поэтому чтение и запись этих переменных за пределамиstrategy.scopeтакже могут работать бесперебойно, без необходимости для пользователя входить в область действия. - Некоторые API стратегий (такие как
strategy.runиstrategy.reduce) которые требуют нахождения в области действия стратегии, автоматически входят в область действия, что означает, что при использовании этих API вам не нужно входить в область действия самостоятельно. - Когда
tf.keras.Modelсоздаётся внутриstrategy.scope, эта информация фиксируется. Когда вызываются методы высокоуровневых фреймворков обучения, такие какmodel.compile,model.fitи т.д., на этой модели, мы автоматически входим в область действия, а также используем эту стратегию для распределения обучения и т.д. См. подробный пример в руководстве по распределенному Keras distributed keras tutorial. Обратите внимание, что простой вызовmodel(..)не затронут - только API высокоуровневых фреймворков обучения.model.compile,model.fit,model.evaluate,model.predictиmodel.saveможно вызывать как внутри, так и вне области действия. - Следующее может быть либо внутри, либо вне области действия:
- Создание наборов данных для входных данных
- Определение
tf.functionдля представления шага обучения - API сохранения, такие как
tf.saved_model.save. Загрузка создает переменные, поэтому это должно происходить внутри области действия, если вы хотите обучить модель распределённым способом. - Сохранение контрольных точек. Как упоминалось выше,
checkpoint.restoreиногда может потребоваться внутри области действия, если она создает переменные.
| Возвращаемое значение | |
|---|---|
| Менеджер контекста. |
© 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/distribute/TPUStrategy