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 | Возвращает резолвер кластера, связанный с данной стратегией. В общем случае, при использовании многоузловой Стратегии, которые намерены иметь связанный Стратегии с одним узлом обычно не имеют
os.environ['TF_CONFIG'] = json.dumps({
'cluster': {
'worker': ["localhost:12345", "localhost:23456"],
'ps': ["localhost:34567"]
},
'task': {'type': 'worker', 'index': 0}
})
# This implicitly uses TF_CONFIG for the cluster and current task info.
strategy = tf.distribute.experimental.MultiWorkerMirroredStrategy()
...
if strategy.cluster_resolver.task_type == 'worker':
# Perform something that's only applicable on workers. Since we set this
# as a worker above, this block will run on this particular instance.
elif strategy.cluster_resolver.task_type == 'ps':
# Perform something that's only applicable on parameter servers. Since we
# set this as a worker above, this block will not run on this particular
# instance.
Для получения дополнительной информации, пожалуйста, обратитесь к документации API |
extended | tf.distribute.StrategyExtended с дополнительными методами. |
num_replicas_in_sync | Возвращает количество реплик, по которым агрегируются градиенты. |
Методы
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,))
| Аргументы | |
|---|---|
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, содержащий значение для каждой реплики. |
Пример использования:
- Возврат постоянного значения для каждой реплики:
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>,)
- Распределение значений в массиве на основе идентификатора реплики:
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,)
- Указание значений с использованием 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,)
- Размещение значений на устройствах и распределение:
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. (устаревшее)
Это позволяет избежать добавления 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 не равен рангу |
| 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" вводится на каждом работнике.
Примечание: Вход в область не автоматически распределяет вычисление, за исключением случаев использования высокоуровневых обучающих фреймворков, таких как kerasmodel.fit. Если вы не используетеmodel.fit, вам необходимо использовать APIstrategy.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