Spec-Zone.ru › TensorFlow

tf.distribute.TPUStrategy

Синхронное обучение на TPUs и TPU-кластерах.

Наследуется от: Strategy

tf.distribute.TPUStrategy(
    tpu_cluster_resolver=None,
    experimental_device_assignment=None,
    experimental_spmd_xla_partitioning=False
)

Используется в ноутбуках

Используется в руководстве Используется в учебниках
  • Миграция с TPU embedding_columns на слой TPUEmbedding
  • Миграция с TPUEstimator на TPUStrategy
  • Использование TPUs
  • Обучение с Orbit
  • Быстрый старт с TPUEmbeddingLayer в TensorFlow 2
  • Решение задач GLUE с помощью BERT на TPU

Для создания объекта 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 будет вызываться на устройстве 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 для примера.

Важно: 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-процесса.

Для получения дополнительных примеров использования и свойств этого метода обратитесь к учебнику по распределённому вводу. Если вы интересуетесь обработкой последней частичной партии, прочитайте эту секцию.

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, содержащий значение для каждой реплики.

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

  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>)
        
  2. Распределение значений массива в зависимости от идентификатора реплики:

    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)
        
  3. Указание значений с использованием 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)
        
  4. Размещение значений на устройствах и их распределение:

    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 или 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)>
Аргументы
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)
END_OF_DOCUMENT_MARKER

Если есть последняя частичная партия, вам необходимо указать ось, чтобы результирующая форма была согласованной на всех репликах. Таким образом, если последняя партия имеет размер 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" вводится на каждом работнике.
Примечание: Вход в блок не автоматически распределяет вычисление, за исключением случаев использования высокоуровневой обучающей рамки, такой как keras model.fit. Если вы не используете model.fit, вам необходимо использовать API strategy.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

Spec-Zone.ru

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