Spec-Zone.ru › TensorFlow 2.9

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 Возвращает резолвер кластера, связанный со стратегией.

tf.distribute.TPUStrategy предоставляет связанный tf.distribute.cluster_resolver.ClusterResolver. Если пользователь предоставил его в __init__, возвращается этот экземпляр; если нет, предоставляется по умолчанию tf.distribute.cluster_resolver.TPUClusterResolver.

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.

Важно: Набор данных tf.data.Dataset, возвращаемый dataset_fn должен иметь размер пакета на реплику, в отличие от experimental_distribute_dataset, использующего глобальный размер пакета. Это можно вычислить, используя input_context.get_per_replica_batch_size.
Примечание: Если вы используете 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 Идентификатор логического устройства не соответствует общему количеству разделов, заданных назначением устройства, или 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, содержащий значение для каждой реплики.

Пример использования:

  1. Возврат постоянного значения для каждой реплики:
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>)
  1. Распределение значений в массиве на основе replica_id:
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
array_value = np.array([3., 2., 1.])
def value_fn(ctx):
  return array_value[ctx.replica_id_in_sync_group]
distributed_values = (
     strategy.experimental_distribute_values_from_function(
       value_fn))
local_result = strategy.experimental_local_results(distributed_values)
local_result
(3.0, 2.0)
  1. Указание значений с помощью 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)
  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(), 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.

experimental_split_to_logical_devices

Просмотр исходного кода

experimental_split_to_logical_devices(
    tensor, partition_dimensions
)

Добавляет аннотацию, что tensor будет разделен по логическим устройствам.

Это добавляет аннотацию к тензору tensor, указывающую, что операции над tensor будут разделены между несколькими логическими устройствами. Тензор tensor будет разделён по измерениям, указанным в partition_dimensions. Измерения tensor должны быть делимы на соответствующее значение в partition_dimensions.

Например, для системы с 8 логическими устройствами, если tensor — это тензор изображения с формой (batch_size, width, height, channel), и partition_dimensions — [1, 2, 4, 1], то tensor будет разделён на 2 в измерении width и 4 способами в измерении height, и значения разделенного тензора будут поданы на 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),))

Args: tensor: Входной тензор для аннотирования. partition_dimensions: Неупорядоченный список целых чисел размером, равным рангу tensor , определяющий, как будет разбит tensor. Произведение всех элементов в partition_dimensions должно быть равно общему количеству логических устройств на реплику.

Возможные ошибки
ValueError

1) Если размер partition_dimensions не равен рангу tensor или 2) если произведение элементов partition_dimensions не соответствует количеству логических устройств на реплику, определённому спецификацией устройства реализуемой стратегией распределения, или 3) если известный размер tensor не делится на соответствующее значение в 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.DistributedValues value, его компонентные тензоры должны иметь ранг не равный нулю. В противном случае, рассмотрите использование 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)>
Args
value экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который необходимо объединить в один тензор. Он также может быть обычным тензором, когда используется с tf.distribute.OneDeviceStrategy или по умолчанию. Тензоры, составляющие DistributedValues, могут быть только плотными тензорами с ненулевым рангом, НЕ tf.IndexedSlices.
axis 0-мерный тензор int32. Измерение, вдоль которого происходит сбор. Должно быть в диапазоне [0, rank(value)).
Возвращает
A 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, поскольку они позволяют настраивать место назначения результата. Они также вызываются в контексте между репликами.

Каким должен быть параметр axis?

Учитывая значение на реплике, возвращённое run, скажем, потерю на пример, партия будет разделена между всеми репликами. Эта функция позволяет агрегировать по репликам и необязательно также по элементам партии, указав параметр axis соответственно.

Например, если у вас глобальный размер партии 8 и 2 реплики, значения для примеров [0, 1, 2, 3] будут на реплике 0, а [4, 5, 6, 7] на реплике 1. С axis=None, reduce будет агрегировать только по репликам, возвращая [0+4, 1+5, 2+6, 3+7]. Это полезно, когда каждая реплика вычисляет скаляр или какое-то другое значение, у которого нет размерности «партия» (например, градиент или потеря).

strategy.reduce("sum", per_replica_result, axis=None)

Иногда вам нужно агрегировать как по глобальной партии, так и по всем репликам. Вы можете получить это поведение, задав размерность партии в качестве axis, обычно axis=0. В этом случае это вернёт скаляр 0+1+2+3+4+5+6+7.

strategy.reduce("sum", per_replica_result, axis=0)

Если есть последняя частичная партия, вам нужно указать ось, чтобы полученная форма была согласованной между репликами. Итак, если последняя партия имеет размер 6 и делится на [0, 1, 2, 3] и [4, 5], у вас возникнет несоответствие форм, если вы не укажете axis=0. Если вы укажите tf.distribute.ReduceOp.MEAN, используя axis=0 будет использоваться правильный знаменатель 6. Противопоставьте это вычислению reduce_mean для получения скалярного значения на каждой реплике и этой функции для усреднения этих средних значений, которые будут взвешивать некоторые значения 1/8 и другие 1/4.

Args
reduce_op значение tf.distribute.ReduceOp, определяющее, как должны быть объединены значения. Разрешает использовать строковое представление перечисления, например, "SUM", "MEAN".
value экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который необходимо объединить в один тензор. Он также может быть обычным тензором, если используется с OneDeviceStrategy или стратегией по умолчанию.
axis указывает размерность, по которой нужно уменьшить тензор каждой реплики. Обычно следует устанавливать в размерность партии или None для уменьшения только по репликам (например, если тензор не имеет размерности партии).
Возвращает
A 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".
Примечание: Вход в область не автоматически распределяет вычисления, за исключением случаев использования высокоуровневого фреймворка обучения, например, keras model.fit. Если вы не используете model.fit, вам нужно использовать API strategy.run для явного распределения вычислений. См. пример в руководстве по созданию пользовательского цикла обучения https://www.tensorflow.org/tutorials/distribute/custom_training.

Что должно находиться в области действия, а что — вне её?

Существует ряд требований к тому, что должно происходить внутри области действия. Однако в местах, где у нас есть информация о используемой стратегии, мы часто входим в область действия для пользователя, поэтому ему не нужно делать это явно (то есть вызов внутри или вне области действия допустим).

  • Всё, что создаёт переменные, которые должны быть распределёнными переменными, должно вызываться в области strategy.scope. Этого можно добиться, либо напрямую вызвав функцию создания переменных внутри контекста области действия, либо полагаясь на другой API, например, strategy.run или keras.Model.fit, чтобы он автоматически вошёл в него за вас. Любая переменная, созданная вне области действия, не будет распределена и может иметь последствия для производительности. Некоторые общие объекты, создающие переменные в TF, — это модели, оптимизаторы, метрики. Такие объекты всегда должны быть инициализированы в области действия, и любые функции, которые могут лениво создавать переменные (например, Model.__call__(), трассировка tf.function и т. д.), аналогичным образом должны вызываться внутри области действия. Ещё одним источником создания переменных может быть восстановление контрольной точки — когда переменные создаются лениво. Обратите внимание, что любая переменная, созданная внутри стратегии, сохраняет информацию о стратегии. Таким образом, чтение и запись в эти переменные вне strategy.scope также могут работать без проблем, без необходимости ввода пользователя в область действия.
  • Некоторые API стратегии (например, strategy.run и strategy.reduce). которые должны находиться в области действия стратегии, автоматически входят в область действия, а значит, при использовании этих API вам не нужно входить в область действия явно.
  • Когда tf.keras.Model создаётся внутри strategy.scope, объект модели сохраняет информацию об области действия. При вызове методов высокоуровневого фреймворка обучения, таких как model.compile, model.fit и т. д., захваченная область действия будет автоматически введена, и связанная стратегия будет использоваться для распределения обучения и т. д. См. подробный пример в руководстве по распределённому Keras. ПРЕДУПРЕЖДЕНИЕ: Просто вызов model(..) не автоматически входит в захваченную область действия — только высокоуровневые API обучения поддерживают это поведение: model.compile, model.fit, model.evaluate, model.predict и model.save могут вызываться внутри или вне области действия.
  • Следующее может быть как внутри, так и вне области действия:
    • Создание наборов данных для входных данных
    • Определение tf.functionов, которые представляют ваш шаг обучения
    • API сохранения, такие как tf.saved_model.save. Загрузка создаёт переменные, поэтому это должно находиться в области действия, если вы хотите обучить модель в распределённом режиме.
    • Сохранение контрольных точек. Как упоминалось выше — checkpoint.restore иногда может потребоваться в области действия, если оно создаёт переменные.
Возвращаемое значение
Менеджер контекста.

© 2022 The TensorFlow Authors. All rights reserved.
Licensed under the Creative Commons Attribution License 4.0.
Code samples licensed under the Apache 2.0 License.
https://www.tensorflow.org/versions/r2.9/api_docs/python/tf/distribute/TPUStrategy

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API