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 | Возвращает количество реплик, по которым агрегируются градиенты. |
Методы
distribute_datasets_from_function
distribute_datasets_from_function(
dataset_fn, options=None
)
Распределяет экземпляры tf.data.Dataset, созданные вызовами dataset_fn.
Аргумент dataset_fn, который пользователи передают, представляет собой функцию ввода, имеющую аргумент tf.distribute.InputContext и возвращающую экземпляр tf.data.Dataset. Ожидается, что возвращаемый набор данных из dataset_fn уже сгруппирован по размеру пакета на реплику (то есть, глобальный размер пакета, деленный на количество реплик в синхронизации) и поделен. tf.distribute.Strategy.distribute_datasets_from_function не группирует и не разделяет экземпляр tf.data.Dataset, возвращаемый функцией ввода. dataset_fn будет вызываться на процессорном устройстве каждого из узлов, и каждый из них сгенерирует набор данных, где каждая реплика на этом узле будет извлекать одну партию ввода (то есть, если узел имеет две реплики, две партии будут извлечены из Dataset на каждом шаге).
Этот метод может использоваться для нескольких целей. Во-первых, он позволяет указать собственную логику группировки и разделения. (В отличие от tf.distribute.experimental_distribute_dataset, которая выполняет группировку и разделение за вас.) Например, где experimental_distribute_dataset не может разделить входные файлы, этот метод может использоваться для ручного разделения набора данных (избегая медленного поведения по умолчанию в experimental_distribute_dataset). В случаях, когда набор данных бесконечен, это разделение может быть выполнено путем создания реплик набора данных, которые различаются только своим случайным начальным значением.
Функция dataset_fn должна принимать экземпляр tf.distribute.InputContext, где можно получить информацию о группировании и репликации ввода.
Вы можете использовать свойство element_spec возвращенного этим API tf.distribute.DistributedDataset для запроса tf.TypeSpec элементов, возвращаемых итератором. Это может использоваться для установки свойства input_signature tf.function. Следуйте tf.distribute.DistributedDataset.element_spec, чтобы увидеть пример.
Примечание: Если вы используете TPUStrategy, порядок обработки данных узлами при использованииtf.distribute.Strategy.experimental_distribute_datasetилиtf.distribute.Strategy.distribute_datasets_from_functionне гарантируется. Это обычно необходимо, если вы используетеtf.distributeдля масштабирования прогнозирования. Вы можете, однако, вставить индекс для каждого элемента в пачке и упорядочить результаты соответственно. Обратитесь к этому фрагменту для примера того, как упорядочить результаты.
Примечание: Состоятельные преобразования наборов данных в настоящее время не поддерживаются сtf.distribute.experimental_distribute_datasetилиtf.distribute.distribute_datasets_from_function. Любые состоятельные операции, которые может иметь набор данных, в настоящее время игнорируются. Например, если ваш набор данных имеетmap_fn, который используетtf.random.uniformдля поворота изображения, то у вас есть граф набора данных, который зависит от состояния (то есть, случайного начального значения) на локальной машине, где выполняется процесс python.
Для получения дополнительной информации об использовании и свойствах этого метода ознакомьтесь с руководством по распределённому вводу. Если вы заинтересованы в обработке последней частичной партии, прочтите эту секцию.
| Аргументы | |
|---|---|
dataset_fn | Функция, принимающая экземпляр tf.distribute.InputContext и возвращающая tf.data.Dataset. |
options | tf.distribute.InputOptions, используемый для управления параметрами распределения этого набора данных. |
| Возвращаемое значение | |
|---|---|
tf.distribute.DistributedDataset. |
experimental_distribute_dataset
experimental_distribute_dataset(
dataset, options=None
)
Распределяет экземпляр tf.data.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 для предварительной загрузки на устройство. options: tf.distribute.InputOptions, используемый для управления параметрами распределения этого набора данных.
| Возвращаемое значение | |
|---|---|
"Распределённый Dataset", по которому можно выполнить итерацию. |
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, содержащий значение для каждой реплики. |
Пример использования:
- Возврат постоянного значения для каждой реплики:
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>)
- Распределение значений в массиве в зависимости от идентификатора реплики:
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)
- Указание значений с помощью 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)
- Размещение значений на устройствах и распределение:
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 имеется один узел, поэтому возвращаемое значение будет содержать все значения на этом узле.
| Args | |
|---|---|
value | Значение, возвращённое run(), extended.call_for_each_replica(), или переменная, созданная в scope. |
| Returns | |
|---|---|
Кортеж значений, содержащихся в value. Если value представляет одиночное значение, возвращается (value,). |
gather
gather(
value, axis
)
Сборка value по репликам вдоль axis на текущее устройство.
Принимая во внимание объект типа tf.distribute.DistributedValues или tf.Tensor value, этот API собирает и конкатенирует value по репликам вдоль axis измерения. Результат копируется на «текущее» устройство
- которое, как правило, является процессором (CPU) рабочего узла, на котором выполняется программа. Для
tf.distribute.TPUStrategy, это первый хост TPU. Для многоклиентскихMultiWorkerMirroredStrategy, это процессор (CPU) каждого рабочего узла.
Этот 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.DistributedValuesvalue, тензоры компонентов должны иметь ранг не равный нулю. В противном случае рассмотрите возможность использованияtf.expand_dimsперед их сборкой.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
# A DistributedValues with component tensor of shape (2, 1) on each replica
distributed_values = strategy.experimental_distribute_values_from_function(lambda _: tf.identity(tf.constant([[1], [2]])))
@tf.function
def run():
return strategy.gather(distributed_values, axis=0)
run()
<tf.Tensor: shape=(4, 1), dtype=int32, numpy=
array([[1],
[2],
[1],
[2]], dtype=int32)>
Рассмотрим следующий пример для получения более полных комбинаций:
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1", "GPU:2", "GPU:3"])
single_tensor = tf.reshape(tf.range(6), shape=(1,2,3))
distributed_values = strategy.experimental_distribute_values_from_function(lambda _: tf.identity(single_tensor))
@tf.function
def run(axis):
return strategy.gather(distributed_values, axis=axis)
axis=0
run(axis)
<tf.Tensor: shape=(4, 2, 3), dtype=int32, numpy=
array([[[0, 1, 2],
[3, 4, 5]],
[[0, 1, 2],
[3, 4, 5]],
[[0, 1, 2],
[3, 4, 5]],
[[0, 1, 2],
[3, 4, 5]]], dtype=int32)>
axis=1
run(axis)
<tf.Tensor: shape=(1, 8, 3), dtype=int32, numpy=
array([[[0, 1, 2],
[3, 4, 5],
[0, 1, 2],
[3, 4, 5],
[0, 1, 2],
[3, 4, 5],
[0, 1, 2],
[3, 4, 5]]], dtype=int32)>
axis=2
run(axis)
<tf.Tensor: shape=(1, 2, 12), dtype=int32, numpy=
array([[[0, 1, 2, 0, 1, 2, 0, 1, 2, 0, 1, 2],
[3, 4, 5, 3, 4, 5, 3, 4, 5, 3, 4, 5]]], dtype=int32)>
| Args | |
|---|---|
value | экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который нужно объединить в один тензор. Он также может быть обычным тензором при использовании с tf.distribute.OneDeviceStrategy или по умолчанию стратегией. Тензоры, составляющие DistributedValues, могут быть только плотными тензорами с ненулевым рангом, НЕ tf.IndexedSlices. |
axis | 0-мерный int32 тензор. Измерение, вдоль которого выполняется сборка. Должно быть в диапазоне [0, ранг(значение)). |
| Returns | |
|---|---|
Tensor, который является конкатенацией value по репликам вдоль axis измерения. |
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 из Tensors. |
args | (Необязательно) Позиционные аргументы для fn. |
kwargs | (Необязательно) Аргументы ключевых слов для fn. |
options | (Необязательно) Экземпляр tf.distribute.RunOptions, определяющий параметры запуска fn. |
| Returns | |
|---|---|
Возвращаемое значение запуска fn. |
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()теперь будет возвращать эту стратегию. За пределами этой области она возвращает стратегию по умолчанию - no-op. - Вход в область также означает вход в «межреплицируемый контекст». См.
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.4/api_docs/python/tf/distribute/experimental/CentralStorageStrategy