Spec-Zone.ru › TensorFlow 2.3

tf.distribute.MirroredStrategy

Просмотреть исходный код на GitHub

Синхронное обучение на нескольких репликах на одном компьютере.

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

tf.distribute.MirroredStrategy(
    devices=None, cross_device_ops=None
)

Эта стратегия обычно используется для обучения на одном компьютере с несколькими графическими процессорами. Для TPUs используйте tf.distribute.TPUStrategy. Чтобы использовать MirroredStrategy с несколькими рабочими процессами, обратитесь к tf.distribute.experimental.MultiWorkerMirroredStrategy.

Например, переменная, созданная в рамках MirroredStrategy, является MirroredVariable. Если в аргументе конструктора стратегии устройства не указаны, она будет использовать все доступные графические процессоры. Если графические процессоры не найдены, она будет использовать доступные центральные процессоры. Обратите внимание, что TensorFlow обрабатывает все центральные процессоры на компьютере как одно устройство и использует потоки внутри для параллелизма.

strategy = tf.distribute.MirroredStrategy()
with strategy.scope():
  x = tf.Variable(1.)
x
MirroredVariable:{
    0: <tf.Variable 'Variable:0' shape=() dtype=float32, numpy=1.0>
  }

При использовании стратегий распределения все создание переменных должно выполняться в рамках области действия стратегии. Это позволит дублировать переменные на всех репликах и поддерживать их синхронизацию с помощью алгоритма all-reduce.

Переменные, созданные внутри MirroredStrategy, который заключен в tf.function, все еще MirroredVariables.

x = []
@tf.function  # Wrap the function with tf.function.
def create_variable():
  if not x:
    x.append(tf.Variable(1.))
strategy = tf.distribute.MirroredStrategy()
with strategy.scope():
  create_variable()
  print (x[0])
MirroredVariable:{
    0: <tf.Variable 'Variable:0' shape=() dtype=float32, numpy=1.0>
  }

experimental_distribute_dataset может использоваться для распределения набора данных по репликам при написании собственного цикла обучения. Если вы используете .fit и .compile методы, доступные в tf.keras, то tf.keras будет обрабатывать распределение за вас.

Например:

my_strategy = tf.distribute.MirroredStrategy()
with my_strategy.scope():
  @tf.function
  def distribute_train_epoch(dataset):
    def replica_fn(input):
      # process input and return result
      return result

    total_result = 0
    for x in dataset:
      per_replica_result = my_strategy.run(replica_fn, args=(x,))
      total_result += my_strategy.reduce(tf.distribute.ReduceOp.SUM,
                                         per_replica_result, axis=None)
    return total_result

  dist_dataset = my_strategy.experimental_distribute_dataset(dataset)
  for _ in range(EPOCHS):
    train_result = distribute_train_epoch(dist_dataset)
Аргументы
devices список строк устройств, таких как ['/gpu:0', '/gpu:1']. Если None, используются все доступные графические процессоры. Если графические процессоры не найдены, используется ЦП.
cross_device_ops необязательно, наследник CrossDeviceOps. Если это не задано, по умолчанию используется NcclAllReduce(). Вы бы настраивали это, если NCCL недоступен или если доступна специальная реализация, использующая определённое оборудование.
Атрибуты
cluster_resolver Возвращает решатель кластера, связанный с этой стратегией.

В целом, при использовании многопроцессорной стратегии распределения, такой как tf.distribute или tf.distribute.experimental.MultiWorkerMirroredStrategy, существует связанный с используемой стратегией tf.distribute.experimental.TPUStrategy(), и такая инстанция возвращается этим свойством.

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

Однопроцессорные стратегии обычно не имеют None, и в этих случаях это свойство вернёт tf.distribute.cluster_resolver.ClusterResolver.

Получаемый None может быть полезен, когда пользователю необходимо получить информацию, такую как спектр кластера, тип задачи или идентификатор задачи. Например,

tf.distribute.cluster_resolver.ClusterResolver

Для получения дополнительной информации, пожалуйста, обратитесь к документации API

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.
.

tf.distribute.cluster_resolver.ClusterResolver extended с дополнительными методами.
tf.distribute.StrategyExtended Возвращает количество реплик, по которым агрегируются градиенты.

Методы

num_replicas_in_sync

Просмотреть исходный код

experimental_assign_to_logical_device

Добавляет аннотацию, что

experimental_assign_to_logical_device(
    tensor, logical_device_id
)
будет назначено логическому устройству.

Примечание: Этот API поддерживается только в TPUStrategy на данный момент. Это добавляет аннотацию к 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

В фрагменте кода выше,

strategy = tf.distribute.MirroredStrategy()

# Create a dataset
dataset = dataset_ops.Dataset.TFRecordDataset([
  "/a/1.tfr", "/a/2.tfr", "/a/3.tfr", "/a/4.tfr"])

# Distribute that dataset
dist_dataset = strategy.experimental_distribute_dataset(dataset)

# Iterate over the `tf.distribute.DistributedDataset`
for x in dist_dataset:
  # process dataset elements
  strategy.run(replica_fn, args=(x,))
tf.distribute.DistributedDataset разбивается на пакеты с размером dist_dataset, и мы перебираем его, используя GLOBAL_BATCH_SIZE. for x in dist_dataset x содержащий данные для всех реплик, которые агрегируются в пакет из tf.distribute.DistributedValues. GLOBAL_BATCH_SIZE будет заботиться о подаче правильных данных для каждой реплики в tf.distribute.Strategy.run к соответствующим x , выполняемым на каждой реплике.

Что происходит под капотом этого метода, когда мы говорим, что экземпляр replica_fn - tf.data.Dataset - распределяется? Это зависит от того, как вы установили dataset через tf.data.experimental.AutoShardPolicy. По умолчанию он установлен на tf.data.experimental.DistributeOptions. В многопроцессорной среде мы сначала попытаемся распределить tf.data.experimental.AutoShardPolicy.AUTO путём определения, создаётся ли dataset из наборов данных чтения (например, dataset, tf.data.TFRecordDataset и т.д.) и если да, то попытаемся разбить входные файлы. Обратите внимание, что должно быть как минимум один входной файл на каждый рабочий процесс. Если у вас меньше одного входного файла на каждый рабочий процесс, рекомендуется отключить фрагментацию наборов данных между рабочими процессами, установив tf.data.TextLineDataset на tf.data.experimental.DistributeOptions.auto_shard_policy.

Если попытка разбить по файлам неуспешна (т.е. набор данных не читается из файлов), мы разделим набор данных равномерно в конце, добавив операцию tf.data.experimental.AutoShardPolicy.OFF в конец потока обработки. Это заставит весь предобработочный поток для всех данных выполняться на каждом рабочем процессе, и каждый рабочий процесс будет выполнять избыточную работу. Мы выведем предупреждение, если этот вариант будет выбран.

Как уже упоминалось, внутри каждого рабочего процесса мы также разделим данные между всеми устройствами рабочего процесса (если их более одного). Это произойдёт даже если фрагментация между рабочими процессами отключена.

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

Вы также можете использовать свойство tf.distribute.Strategy.experimental_distribute_datasets_from_function экземпляра element_spec, возвращаемого этим API, чтобы запросить tf.distribute.DistributedDataset элементов, возвращаемых итератором. Это может быть использовано для установки свойства tf.TypeSpec input_signature.

tf.function
Примечание: Порядок обработки данных рабочими процессами при использовании
strategy = tf.distribute.MirroredStrategy()

# Create a dataset
dataset = dataset_ops.Dataset.TFRecordDataset([
  "/a/1.tfr", "/a/2.tfr", "/a/3.tfr", "/a/4.tfr"])

# Distribute that dataset
dist_dataset = strategy.experimental_distribute_dataset(dataset)

@tf.function(input_signature=[dist_dataset.element_spec])
def train_step(inputs):
  # train model with inputs
  return

# Iterate over the `tf.distribute.DistributedDataset`
for x in dist_dataset:
  # process dataset elements
  strategy.run(train_step, args=(x,))
или tf.distribute.Strategy.experimental_distribute_dataset не гарантируется. Это обычно требуется, если вы используете tf.distribute.Strategy.experimental_distribute_datasets_from_function для масштабирования предсказания. Однако вы можете вставить индекс для каждого элемента в пакете и упорядочить выводы соответственно. Обратитесь к этому фрагменту для примера того, как упорядочить выводы.
Аргументы
tf.distribute dataset, который будет разделен между всеми репликами в соответствии с описанными выше правилами.
tf.data.Dataset options, используемый для управления параметрами распределения этого набора данных.
Возвращает
tf.distribute.InputOptions.

tf.distribute.DistributedDataset

Просмотреть исходный код

experimental_distribute_datasets_from_function(
    dataset_fn, options=None
)

Распределяет экземпляры tf.data.Dataset, созданные вызовами dataset_fn.

dataset_fn будет вызван один раз для каждого рабочего узла в стратегии. Каждая реплика на этом рабочем узле будет извлекать одну партию входных данных из локального Dataset (то есть, если у рабочего узла две реплики, две партии будут извлечены из Dataset на каждом шаге).

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

dataset_fn должен принимать экземпляр tf.distribute.InputContext, где можно получить информацию о разбиении на партии и репликации ввода.

Вы также можете использовать свойство element_spec tf.distribute.DistributedDataset, возвращаемого этим API, для запроса tf.TypeSpec элементов, возвращаемых итератором. Это может быть использовано для установки свойства input_signature tf.function.

global_batch_size = 8
def dataset_fn(input_context):
  batch_size = input_context.get_per_replica_batch_size(
                   global_batch_size)
  d = tf.data.Dataset.from_tensors([[1.]]).repeat().batch(batch_size)
  return d.shard(
      input_context.num_input_pipelines,
      input_context.input_pipeline_id)
strategy = tf.distribute.MirroredStrategy()
ds = strategy.experimental_distribute_datasets_from_function(dataset_fn)
def train(ds):
  @tf.function(input_signature=[ds.element_spec])
  def step_fn(inputs):
    # train the model with inputs
    return inputs

... for batch in ds: ... replica_results = strategy.run(replica_fn, args=(batch,))

train(ds)

Ключевой момент: tf.data.Dataset, возвращаемый dataset_fn, должен иметь размер партии на реплику, в отличие от experimental_distribute_dataset, которое использует глобальный размер партии. Это можно вычислить с помощью input_context.get_per_replica_batch_size.
Примечание: Порядок обработки данных рабочими узлами при использовании tf.distribute.Strategy.experimental_distribute_dataset или tf.distribute.Strategy.experimental_distribute_datasets_from_function не гарантируется. Это обычно требуется, если вы используете tf.distribute для масштабирования прогнозирования. Однако вы можете вставить индекс для каждого элемента в партии и упорядочить выводы соответственно. Обратитесь к этому фрагменту для примера, как упорядочить выводы.
Аргументы
dataset_fn Функция, принимающая экземпляр tf.distribute.InputContext и возвращающая tf.data.Dataset.
options tf.distribute.InputOptions, используемый для управления параметрами распределения этого набора данных.
Возвращаемые значения
tf.distribute.DistributedDataset.

experimental_distribute_values_from_function

Просмотреть исходный код

experimental_distribute_values_from_function(
    value_fn
)

Генерирует tf.distribute.DistributedValues из value_fn.

Эта функция предназначена для генерации tf.distribute.DistributedValues для передачи в run, reduce, или другие методы, принимающие распределённые значения, когда не используются наборы данных.

Аргументы
value_fn Функция для выполнения для генерации значений. Она вызывается для каждой реплики с tf.distribute.ValueContext в качестве единственного аргумента. Она должна вернуть тензор или тип, который можно преобразовать в тензор.
Возвращаемые значения
tf.distribute.DistributedValues, содержащий значение для каждой реплики.

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

  1. Возврат постоянного значения для каждой реплики:
strategy = tf.distribute.MirroredStrategy()
def value_fn(ctx):
  return tf.constant(1.)
distributed_values = (
     strategy.experimental_distribute_values_from_function(
       value_fn))
local_result = strategy.experimental_local_results(distributed_values)
local_result
(<tf.Tensor: shape=(), dtype=float32, numpy=1.0>,)
  1. Распределение значений в массиве на основе id реплики:
strategy = tf.distribute.MirroredStrategy()
array_value = np.array([3., 2., 1.])
def value_fn(ctx):
  return array_value[ctx.replica_id_in_sync_group]
distributed_values = (
     strategy.experimental_distribute_values_from_function(
       value_fn))
local_result = strategy.experimental_local_results(distributed_values)
local_result
(3.0,)
  1. Указание значений с использованием num_replicas_in_sync:
strategy = tf.distribute.MirroredStrategy()
def value_fn(ctx):
  return ctx.num_replicas_in_sync
distributed_values = (
     strategy.experimental_distribute_values_from_function(
       value_fn))
local_result = strategy.experimental_local_results(distributed_values)
local_result
(1,)
  1. Размещение значений на устройствах и распределение:
strategy = tf.distribute.TPUStrategy()
worker_devices = strategy.extended.worker_devices
multiple_values = []
for i in range(strategy.num_replicas_in_sync):
  with tf.device(worker_devices[i]):
    multiple_values.append(tf.constant(1.0))

def value_fn(ctx):
  return multiple_values[ctx.replica_id_in_sync_group]

distributed_values = strategy.
  experimental_distribute_values_from_function(
  value_fn)

experimental_local_results

Просмотреть исходный код

experimental_local_results(
    value
)

Возвращает список всех локальных значений на реплику, содержащихся в value.

Примечание: Это возвращает только значения на рабочем узле, инициированном этим клиентом. При использовании tf.distribute.Strategy, например tf.distribute.experimental.MultiWorkerMirroredStrategy, каждый рабочий узел будет собственным клиентом, и эта функция вернёт только значения, вычисленные на этом рабочем узле.
Аргументы
value Значение, возвращённое experimental_run(), run(), extended.call_for_each_replica(), или переменная, созданная в scope.
Возвращаемые значения
Кортеж значений, содержащихся в value. Если value представляет одно значение, это возвращает (value,).

experimental_make_numpy_dataset

Просмотреть исходный код

experimental_make_numpy_dataset(
    numpy_input
)

Создаёт tf.data.Dataset из массива NumPy. (устарело)

Предупреждение: ЭТА ФУНКЦИЯ УСТАРЕЛА. Она будет удалена после 30 сентября 2020 г. Инструкции по обновлению: Используйте tf.data.Dataset.from_tensor_slices вместо этого.

Это позволяет избежать добавления numpy_input в качестве большой константы в граф и копирует данные на машину или машины, которые будут обрабатывать ввод.

Обратите внимание, что вам, вероятно, нужно будет использовать experimental_distribute_dataset с возвращаемым набором данных, чтобы дополнительно распределить его с помощью стратегии.

Пример:

strategy = tf.distribute.MirroredStrategy()
numpy_input = np.ones([10], dtype=np.float32)
dataset = strategy.experimental_make_numpy_dataset(numpy_input)
dataset
<TensorSliceDataset shapes: (), types: tf.float32>
dataset = dataset.batch(2)
dist_dataset = strategy.experimental_distribute_dataset(dataset)
Аргументы
numpy_input Вложенный набор массивов NumPy, которые будут преобразованы в набор данных. Обратите внимание, что массивы NumPy складываются, так как это стандартное поведение tf.data.Dataset.
Возвращаемые значения
tf.data.Dataset, представляющий numpy_input.

experimental_replicate_to_logical_devices

Просмотреть исходный код

experimental_replicate_to_logical_devices(
    tensor
)

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

Примечание: Этот API поддерживается только в TPUStrategy на данный момент. Это добавляет аннотацию к тензору tensor , указывающую, что операции над tensor будут вызваны на всех логических устройствах.
# Initializing TPU system with 2 logical devices and 4 replicas.
resolver = tf.distribute.cluster_resolver.TPUClusterResolver(tpu='')
tf.config.experimental_connect_to_cluster(resolver)
topology = tf.tpu.experimental.initialize_tpu_system(resolver)
device_assignment = tf.tpu.experimental.DeviceAssignment.build(
    topology,
    computation_shape=[1, 1, 1, 2],
    num_replicas=4)
strategy = tf.distribute.TPUStrategy(
    resolver, experimental_device_assignment=device_assignment)

iterator = iter(inputs)

@tf.function()
def step_fn(inputs):
  images, labels = inputs
  images = strategy.experimental_split_to_logical_devices(
    inputs, [1, 2, 4, 1])

  # model() function will be executed on 8 logical devices with `inputs`
  # split 2 * 4  ways.
  output = model(inputs)

  # For loss calculation, all logical devices share the same logits
  # and labels.
  labels = strategy.experimental_replicate_to_logical_devices(labels)
  output = strategy.experimental_replicate_to_logical_devices(output)
  loss = loss_fn(labels, output)

  return loss

strategy.run(step_fn, args=(next(iterator),))

Аргументы: tensor: входной тензор для аннотации.

Возвращаемые значения
Аннотированный тензор с идентичным значением, как у tensor.

experimental_split_to_logical_devices

Просмотреть исходный код

experimental_split_to_logical_devices(
    tensor, partition_dimensions
)

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

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

Например, для системы с 8 логическими устройствами, если tensor это тензор изображения с формой (размер_парт, ширина, высота, канал) и partition_dimensions это [1, 2, 4, 1], то tensor будет разделен на 2 по ширине и 4 по высоте, а значения разделенных тензоров будут поданы на 8 логических устройств.

# Initializing TPU system with 8 logical devices and 1 replica.
resolver = tf.distribute.cluster_resolver.TPUClusterResolver(tpu='')
tf.config.experimental_connect_to_cluster(resolver)
topology = tf.tpu.experimental.initialize_tpu_system(resolver)
device_assignment = tf.tpu.experimental.DeviceAssignment.build(
    topology,
    computation_shape=[1, 2, 2, 2],
    num_replicas=1)
strategy = tf.distribute.TPUStrategy(
    resolver, experimental_device_assignment=device_assignment)

iterator = iter(inputs)

@tf.function()
def step_fn(inputs):
  inputs = strategy.experimental_split_to_logical_devices(
    inputs, [1, 2, 4, 1])

  # model() function will be executed on 8 logical devices with `inputs`
  # split 2 * 4  ways.
  output = model(inputs)
  return output

strategy.run(step_fn, args=(next(iterator),))

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

Возвращает исключения
ValueError

1) Если размер partition_dimensions не равен рангу tensor или 2) если произведение элементов partition_dimensions не соответствует числу логических устройств на реплику, определённому спецификацией устройства реализующей DistributionStrategy или 3) если известный размер tensor не делится на соответствующее значение в partition_dimensions.

Возвращает
Анотированный тензор со значением, идентичным значению tensor.

reduce

Просмотреть исходный код

reduce(
    reduce_op, value, axis
)

Сведение value по репликам.

Учитывая значение, возвращаемое run для каждой реплики, например, потерю на пример, пакет будет распределён по всем репликам. Эта функция позволяет агрегировать значения по репликам, а также, необязательно, по элементам пакета. Например, если у вас есть общий размер пакета 8 и 2 реплики, значения для примеров [0, 1, 2, 3] будут на реплике 0, а [4, 5, 6, 7] — на реплике 1. По умолчанию reduce просто агрегирует по репликам, возвращая [0+4, 1+5, 2+6, 3+7]. Это полезно, когда каждая реплика вычисляет скаляр или какое-то другое значение, у которого нет размерности "пакет" (например, градиент). Чаще всего вам нужно будет агрегировать по глобальному пакету, что можно сделать, указав размерность пакета как axis, обычно как axis=0. В этом случае она вернёт скаляр 0+1+2+3+4+5+6+7.

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

Аргументы
reduce_op Значение tf.distribute.ReduceOp определяющее, как должны комбинироваться значения.
value Значение "для каждой реплики", например, возвращаемое run для объединения в один тензор.
axis Указывает размерность, по которой следует уменьшить тензор каждой реплики. Как правило, её нужно установить на размерность пакета, или на None, чтобы уменьшить только по репликам (например, если тензор не имеет размерности пакета).
Возвращает
A Tensor.

run

Просмотреть исходный код

run(
    fn, args=(), kwargs=None, options=None
)

Выполнение fn на каждой реплике с заданными аргументами.

Выполняет операции, указанные в fn на каждой реплике. Если args или kwargs имеют tf.distribute.DistributedValues, такие как те, что получены от tf.distribute.DistributedDataset из tf.distribute.Strategy.experimental_distribute_dataset или tf.distribute.Strategy.experimental_distribute_datasets_from_function, когда fn выполняется на конкретной реплике, он будет выполнен с компонентом tf.distribute.DistributedValues, соответствующим этой реплике.

fn может вызвать tf.distribute.get_replica_context() для доступа к членам, таким как all_reduce.

Все аргументы в args или kwargs должны быть либо вложенными тензорами, либо tf.distribute.DistributedValues, содержащими тензоры или составные тензоры.

Ключевой момент: В зависимости от реализации tf.distribute.Strategy и от того, включена ли жадная (eager) обработка, fn может быть вызван один или несколько раз. Если fn анотирован tf.function или tf.distribute.Strategy.run вызван внутри tf.function, жадная обработка отключена, и fn вызывается один раз (или один раз на реплику, если используется MirroredStrategy) для генерации графа Tensorflow, который затем будет повторно использоваться для выполнения с новыми входными данными. В противном случае, если жадная обработка включена, fn будет вызываться на каждом шаге, как обычный Python-код.

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

  1. Входной тензор константы.
strategy = tf.distribute.MirroredStrategy()
tensor_input = tf.constant(3.0)
@tf.function
def replica_fn(input):
  return input*2.0
result = strategy.run(replica_fn, args=(tensor_input,))
result
<tf.Tensor: shape=(), dtype=float32, numpy=6.0>
  1. Входные DistributedValues.
strategy = tf.distribute.MirroredStrategy()
@tf.function
def run():
  def value_fn(value_context):
    return value_context.num_replicas_in_sync
  distributed_values = (
    strategy.experimental_distribute_values_from_function(
      value_fn))
  def replica_fn2(input):
    return input*2
  return strategy.run(replica_fn2, args=(distributed_values,))
result = run()
result
<tf.Tensor: shape=(), dtype=int32, numpy=2>
Аргументы
fn Функция для выполнения. Выход должен быть tf.nest из Tensors.
args (Необязательно) Позиционные аргументы для fn.
kwargs (Необязательно) Именные аргументы для fn.
options (Необязательно) Экземпляр tf.distribute.RunOptions, определяющий параметры выполнения fn.
Возвращает
Объединённое возвращаемое значение fn по всем репликам. Структура возвращаемого значения такая же, как у возвращаемого значения fn. Каждый элемент структуры может быть tf.distribute.DistributedValues, объектами Tensor или Tensors (например, при выполнении на одной реплике).

scope

Просмотреть исходный код

scope()

Блок кода для установки стратегии в текущий контекст и распределения переменных.

Этот метод возвращает блок кода и используется следующим образом:

strategy = tf.distribute.MirroredStrategy()
# 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>
}
# 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 для явного распределения этого вычисления. См. пример в учебнике по пользовательской циклу обучения.

Что должно быть в блоке, а что вне?

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

  • Все, что создает переменные, которые должны быть распределенными переменными, должно находиться в 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. Обратите внимание, что простой вызов 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.3/api_docs/python/tf/distribute/MirroredStrategy

Spec-Zone.ru

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