tf.compat.v1.distribute.experimental.CentralStorageStrategy
Стратегия для одной машины, помещающая все переменные на один узел.
Наследуется от: Strategy
tf.compat.v1.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_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 путём определения, создаётся ли 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 экземпляра tf.distribute.DistributedDataset, возвращаемого этим API, для запроса 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
experimental_distribute_datasets_from_function(
dataset_fn, options=None
)
Распределяет экземпляры tf.data.Dataset, созданные вызовами dataset_fn.
dataset_fn будет вызван один раз для каждого узла в стратегии. Каждая реплика на этом узле будет извлекать один пакет входных данных из локального Dataset (т. е. если у узла две реплики, две партии будут извлекаться из Dataset на каждом шаге).
Этот метод может быть использован для нескольких целей. Например, в тех случаях, когда experimental_distribute_dataset не может разбить входные файлы, этот метод может быть использован для ручного разбиения набора данных (избегая медленного поведения по умолчанию в experimental_distribute_dataset). В случаях, когда набор данных бесконечен, это разбиение можно выполнить, создав реплики наборов данных, которые различаются только своим seed для генерации случайных чисел. 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.distribute.Strategy.experimental_distribute_datasetилиtf.distribute.Strategy.experimental_distribute_datasets_from_functionне гарантируется. Это обычно требуется, если вы используетеtf.distributeдля масштабирования прогнозирования. Однако вы можете вставить индекс для каждого элемента в пакете и упорядочить результаты соответственно. См. этот фрагмент здесь для примера упорядочивания результатов.
| Args | |
|---|---|
dataset_fn | Функция, принимающая экземпляр tf.distribute.InputContext и возвращающая tf.data.Dataset. |
options | tf.distribute.InputOptions, используемый для управления параметрами распределения этого набора данных. |
| Returns | |
|---|---|
tf.distribute.DistributedDataset. |
experimental_local_results
experimental_local_results(
value
)
Возвращает список всех локальных значений на реплику, содержащихся в value.
Примечание: Это возвращает только значения на рабочем узле, инициированном этим клиентом. При использованииtf.distribute.Strategy, например,tf.distribute.experimental.MultiWorkerMirroredStrategy, каждый рабочий узел будет своим клиентом, и эта функция вернёт только значения, вычисленные на этом рабочем узле.
| Args | |
|---|---|
value | Значение, возвращённое функциями experimental_run(), run(), extended.call_for_each_replica(), или переменной, созданной в scope |
| Returns | |
|---|---|
Кортеж значений, содержащихся в value . Если value представляет одно значение, то возвращается (value,). |
experimental_make_numpy_dataset
experimental_make_numpy_dataset(
numpy_input, session=None
)
Создаёт tf.data.Dataset для входных данных, предоставленных через массив NumPy.
Это позволяет избежать добавления numpy_input в качестве большой константы в граф и копирует данные на машину или машины, которые будут обрабатывать входные данные.
Обратите внимание, что вам, вероятно, потребуется использовать tf.distribute.Strategy.experimental_distribute_dataset с возвращаемым набором данных, чтобы далее распределить его с помощью стратегии.
Пример:
numpy_input = np.ones([10], dtype=np.float32) dataset = strategy.experimental_make_numpy_dataset(numpy_input) dist_dataset = strategy.experimental_distribute_dataset(dataset)
| Args | |
|---|---|
numpy_input | Вложенный набор массивов NumPy, которые будут преобразованы в набор данных. Обратите внимание, что списки массивов NumPy складываются, так как это нормальное поведение tf.data.Dataset. |
session | (Только для выполнения графов TensorFlow v1.x) Сессия, используемая для инициализации. |
| Returns | |
|---|---|
tf.data.Dataset, представляющий numpy_input. |
experimental_run
experimental_run(
fn, input_iterator=None
)
Выполняет операции в fn на каждой реплике с входными данными из input_iterator.
УСТАНОВЛЕНО КАК УСТАРЕВШЕЕ: Этот метод недоступен в TF 2.x. Используйте run вместо этого.
При включённом режиме выполнения Eager выполняет операции, указанные в fn на каждой реплике. В противном случае создаёт граф для выполнения операций на каждой реплике.
Каждая реплика примет один, отличный вход из входных данных, предоставленных одним вызовом get_next для итератора входных данных.
fn может вызвать tf.distribute.get_replica_context() для доступа к членам, таким как replica_id_in_sync_group.
| Args | |
|---|---|
fn | Функция для выполнения. Входы в функцию должны соответствовать выходам input_iterator.get_next(). Выход должен быть tf.nest Tensor |
input_iterator | (Необязательно) итератор входных данных, из которого берутся входные данные. |
| Returns | |
|---|---|
Объединённое значение возврата fn по всем репликам. Структура возвращаемого значения такая же, как у возвращаемого значения fn Каждый элемент структуры может быть PerReplica (если значения не синхронизированы), Mirrored (если значения синхронизированы) или Tensor (если выполняется на одной реплике). |
make_dataset_iterator
make_dataset_iterator(
dataset
)
Создаёт итератор для входных данных, предоставленных через dataset.
УСТАНОВЛЕНО КАК УСТАРЕВШЕЕ: Этот метод недоступен в TF 2.x.
Данные из заданного набора данных будут равномерно распределяться по всем вычислительным репликам. Мы будем предполагать, что входной набор данных сгруппирован по размеру глобальной пакетной обработки. При этом предположении мы будем стараться разделить каждый пакет по всем репликам (один или несколько рабочих узлов). Если эта попытка не удастся, будет выброшено исключение, и пользователь должен вместо этого использовать make_input_fn_iterator, которое предоставляет больше контроля пользователю и не пытается разделить пакет по репликам.
Пользователь также может использовать make_input_fn_iterator , если хочет настроить, какой вход подаётся на какую реплику/рабочий узел и т. д.
| Args | |
|---|---|
dataset | tf.data.Dataset, который будет равномерно распределяться по всем репликам. |
| Returns | |
|---|---|
tf.distribute.InputIterator, который возвращает входные данные для каждого шага вычисления. Пользователь должен вызвать initialize на возвращённом итераторе. |
make_input_fn_iterator
make_input_fn_iterator(
input_fn, replication_mode=tf.distribute.InputReplicationMode.PER_WORKER
)
Возвращает итератор, распределённый по репликам, созданный из функции входных данных.
УСТАНОВЛЕНО КАК УСТАРЕВШЕЕ: Этот метод недоступен в TF 2.x.
Функция input_fn должна принимать объект tf.distribute.InputContext, где можно получить информацию о разделении на пакеты и ввода:
def input_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)
with strategy.scope():
iterator = strategy.make_input_fn_iterator(input_fn)
replica_results = strategy.experimental_run(replica_fn, iterator)
Возвращаемый tf.data.Dataset из input_fn должен иметь размер пакета на реплику, который можно вычислить с помощью input_context.get_per_replica_batch_size.
| Args | |
|---|---|
input_fn | Функция, принимающая объект tf.distribute.InputContext и возвращающая tf.data.Dataset. |
replication_mode | значение перечисления tf.distribute.InputReplicationMode. В настоящее время поддерживается только PER_WORKER, что означает, что будет один вызов input_fn на один рабочий узел. Реплики будут извлекать элементы из локального tf.data.Dataset на своих рабочих узлах. |
| Returns | |
|---|---|
Объект итератора, который сначала нужно вызвать .initialize(). Затем его можно передать в strategy.experimental_run() или вызвать iterator.get_next() для получения следующего значения для передачи в strategy.extended.call_for_each_replica(). |
reduce
reduce(
reduce_op, value, axis=None
)
Сведение 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, содержащими тензоры или составные тензоры.
Пример использования:
- Входной тензор-константа.
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>
- Входные данные 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, или Tensor (например, при выполнении на одной реплике). |
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.functions, представляющих ваш шаг обучения ** Сохранение API, такие какtf.saved_model.save. Загрузка создаёт переменные, поэтому это должно быть внутри области действия, если вы хотите обучить модель в распределённом режиме. ** Сохранение контрольных точек. Как уже упоминалось ранее —checkpoint.restoreиногда может потребоваться находиться внутри области действия, если оно создаёт переменные.
| Возвращаемое значение | |
|---|---|
| Менеджер контекста. |
update_config_proto
update_config_proto(
config_proto
)
Возвращает копию config_proto, изменённую для использования с данной стратегией.
УСТАРЕЛО: Этот метод недоступен в TF 2.x.
Обновлённая конфигурация содержит необходимые данные для работы с стратегией, например, конфигурацию для запуска коллективных операций или фильтры устройств для улучшения производительности распределённого обучения.
| Аргументы | |
|---|---|
config_proto | Объект tf.ConfigProto. |
| Возвращаемое значение | |
|---|---|
Обновлённая копия config_proto. |
© 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/compat/v1/distribute/experimental/CentralStorageStrategy