Spec-Zone.ru › TensorFlow 2.3

tf.distribute.experimental.CentralStorageStrategy

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

Стратегия для одной машины, которая помещает все переменные на один узел.

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

tf.distribute.experimental.CentralStorageStrategy(
    compute_devices=None, parameter_device=None
)

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

Например:

strategy = tf.distribute.experimental.CentralStorageStrategy()
# Create a dataset
ds = tf.data.Dataset.range(5).batch(2)
# Distribute that dataset
dist_dataset = strategy.experimental_distribute_dataset(ds)

with strategy.scope():
  @tf.function
  def train_step(val):
    return val + 1

  # Iterate over the distributed dataset
  for x in dist_dataset:
    # process dataset elements
    strategy.run(train_step, args=(x,))
Атрибуты
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 может быть полезен, когда пользователю требуется получить доступ к информации, такой как спецификация кластера, тип задачи или идентификатор задачи. Например,

os.environ['TF_CONFIG'] = json.dumps({
'cluster': {
'worker': ["localhost:12345", "localhost:23456"],
'ps': ["localhost:34567"]
},
'task': {'type': 'worker', 'index': 0}
})

# This implicitly uses TF_CONFIG for the cluster and current task info.
strategy = tf.distribute.experimental.MultiWorkerMirroredStrategy()

...

if strategy.cluster_resolver.task_type == 'worker':
# Perform something that's only applicable on workers. Since we set this
# as a worker above, this block will run on this particular instance.
elif strategy.cluster_resolver.task_type == 'ps':
# Perform something that's only applicable on parameter servers. Since we
# set this as a worker above, this block will not run on this particular
# instance.

Для получения дополнительной информации, пожалуйста, обратитесь к документации API 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
)

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

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

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

Возвращаемый dataset — это обернутый dataset стратегии, который создаёт многоузельный итератор под капотом. Он предварительно загружает входные данные на указанные узлы на рабочем узле. На возвращаемый распределённый dataset можно итерироваться так же, как и на обычных dataset'ах.

Примечание: В настоящее время пользователь не может добавить больше преобразований в распределённый dataset.

Например:

strategy = tf.distribute.CentralStorageStrategy()  # with 1 CPU and 1 GPU
dataset = tf.data.Dataset.range(10).batch(2)
dist_dataset = strategy.experimental_distribute_dataset(dataset)
for x in dist_dataset:
  print(x)  # Prints PerReplica values [0, 1], [2, 3],...

Аргументы: dataset: tf.data.Dataset для предварительной загрузки на устройство.

Возвращаемое значение
"Распределённый Dataset" для итерирования.

experimental_distribute_datasets_from_function

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

experimental_distribute_datasets_from_function(
    dataset_fn
)

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

dataset_fn будет вызван один раз для каждого рабочего узла в стратегии. В этом случае у нас есть только один рабочий узел, поэтому dataset_fn вызывается один раз. Каждая реплика на этом рабочем узле затем извлекает пакет элементов из этого локального dataset.

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

Например:

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)

inputs = strategy.experimental_distribute_datasets_from_function(dataset_fn)

for batch in inputs:
  replica_results = strategy.run(replica_fn, args=(batch,))
Ключевой момент: tf.data.Dataset, возвращаемый dataset_fn должен иметь размер пакета на реплику, в отличие от experimental_distribute_dataset, который использует глобальный размер пакета. Это может быть вычислено с помощью input_context.get_per_replica_batch_size.
Аргументы
dataset_fn Функция, принимающая экземпляр tf.distribute.InputContext и возвращающая tf.data.Dataset.
Возвращаемое значение
Распределённый Dataset, на котором можно итерироваться, как на обычных dataset'ах.

experimental_distribute_values_from_function

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

experimental_distribute_values_from_function(
    value_fn
)

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

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

Аргументы
value_fn Функция для генерации значений. Она вызывается для каждой реплики с tf.distribute.ValueContext в качестве единственного аргумента. Она должна возвращать Tensor или тип, который может быть преобразован в Tensor.
Возвращаемое значение
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. Распределение значений в массиве на основе идентификатора реплики:
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.

В CentralStorageStrategy существует один рабочий узел, поэтому возвращаемое значение будет содержать все значения на этом рабочем узле.

Аргументы
value Значение, возвращённое run(), extended.call_for_each_replica(), или переменной, созданной в scope
Возвращаемое значение
Кортеж значений, содержащихся в value . Если value представляет собой единственное значение, возвращается (value,).

experimental_make_numpy_dataset

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

experimental_make_numpy_dataset(
    numpy_input
)

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

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

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

Обратите внимание, что вам, вероятно, потребуется использовать experimental_distribute_dataset с возвращаемым 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)
Args
numpy_input вложенный набор массивов NumPy, которые будут преобразованы в набор данных. Обратите внимание, что массивы NumPy расположены стопкой, так как это стандартное поведение tf.data.Dataset.
Returns
Объект 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),))

Args: tensor: Входной тензор для аннотации.

Returns
Тензор с аннотацией, имеющий идентичное значение, как у 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),))

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

Raises
ValueError

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

Returns
Тензор с аннотацией, имеющий идентичное значение, как у 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.

Пример:

strategy = tf.distribute.experimental.CentralStorageStrategy(
    compute_devices=['CPU:0', 'GPU:0'], parameter_device='CPU:0')
ds = tf.data.Dataset.range(10)
# Distribute that dataset
dist_dataset = strategy.experimental_distribute_dataset(ds)

with strategy.scope():
  @tf.function
  def train_step(val):
    # pass through
    return val

  # Iterate over the distributed dataset
  for x in dist_dataset:
    result = strategy.run(train_step, args=(x,))

result = strategy.reduce(tf.distribute.ReduceOp.SUM, result,
                         axis=None).numpy()
# result: array([ 4,  6,  8, 10])

result = strategy.reduce(tf.distribute.ReduceOp.SUM, result, axis=0).numpy()
# result: 28
Args
reduce_op Значение tf.distribute.ReduceOp, определяющее, как следует объединять значения.
value «Значение на реплику», например, возвращаемое run, которое следует объединить в один тензор.
axis Указывает измерение для сведения вдоль тензора каждой реплики. Обычно должно быть установлено на размер пакета или None для сведения только по репликам (например, если тензор не имеет размера пакета).
Returns
Tensor.

run

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

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

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

В CentralStorageStrategy, fn вызывается на каждой вычислительной реплике с предоставленными аргументами «на реплику», специфичными для этого устройства.

Args
fn Функция для выполнения. Результат должен быть tf.nest из Tensor .
args (Необязательно) Позиционные аргументы для fn.
kwargs (Необязательно) Именованные аргументы для fn.
options (Необязательно) Экземпляр tf.distribute.RunOptions, задающий параметры для запуска fn.
Returns
Возвращаемое значение при запуске fn.

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/experimental/CentralStorageStrategy

Spec-Zone.ru

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