tf.compat.v1.distribute.experimental.TPUStrategy
Реализация стратегии распределения для TPU.
Наследуется от: Strategy
tf.compat.v1.distribute.experimental.TPUStrategy(
tpu_cluster_resolver=None, steps_per_run=None, device_assignment=None
)
| Аргументы | |
|---|---|
tpu_cluster_resolver | tf.distribute.cluster_resolver.TPUClusterResolver, предоставляющий информацию о кластере TPU. |
steps_per_run | Количество шагов для выполнения на устройстве перед возвратом на хост. Обратите внимание, что это может повлиять на производительность, хуки, метрики, сводки и т. д. Этот параметр используется только при использовании стратегии распределения с оценщиком или Keras. |
device_assignment | Необязательный tf.tpu.experimental.DeviceAssignment для указания размещения реплик в кластере TPU. В настоящее время поддерживается только случай использования одного ядра в кластере TPU. |
| Атрибуты | |
|---|---|
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 | Возвращает количество реплик, по которым агрегируются градиенты. |
steps_per_run | УСТАНОВЛЕНО: используйте .extended.steps_per_run вместо этого. |
Методы
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 экземпляра 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). В случаях с бесконечным набором данных это фрагментирование может быть выполнено путём создания реплик набора данных, которые различаются только случайным начальным значением. 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для масштабирования предсказаний. Однако вы можете вставить индекс для каждого элемента в пакете и упорядочить выходы соответственно. Обратитесь к этому фрагменту для примера того, как упорядочить выходы.
| Аргументы | |
|---|---|
dataset_fn | Функция, принимающая экземпляр tf.distribute.InputContext и возвращающая tf.data.Dataset. |
options | tf.distribute.InputOptions, используемые для управления параметрами распределения этого набора данных. |
| Возвращаемое значение | |
|---|---|
tf.distribute.DistributedDataset. |
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, 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)
| Аргументы | |
|---|---|
numpy_input | Вложенный массив входных массивов NumPy, которые будут преобразованы в набор данных. Обратите внимание, что списки массивов NumPy складываются, так как это обычное поведение tf.data.Dataset. |
session | (Только для выполнения графа TensorFlow v1.x) Сессия, используемая для инициализации. |
| Возвращаемое значение | |
|---|---|
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.
| Аргументы | |
|---|---|
fn | Функция для выполнения. Входные данные для функции должны соответствовать выходам input_iterator.get_next() . Выход должен быть tf.nest из Tensor . |
input_iterator | (Необязательно) итератор входных данных, из которого берутся входные данные. |
| Возвращаемое значение | |
|---|---|
Объединённое возвращаемое значение fn по всем репликам. Структура возвращаемого значения такая же, как возвращаемое значение fn . Каждый элемент в структуре может быть PerReplica (если значения не синхронизированы), Mirrored (если значения синхронизированы) или Tensor (если выполняется на одной реплике). |
make_dataset_iterator
make_dataset_iterator(
dataset
)
Создаёт итератор для входных данных, предоставленных через dataset.
УСТАРЕВШИЙ метод: Этот метод недоступен в TF 2.x.
Данные из заданного набора данных будут распределены равномерно по всем вычислительным репликам. Мы предположим, что входной набор данных сгруппирован по глобальному размеру пакета. С этим предположением мы сделаем всё возможное, чтобы разделить каждый пакет по всем репликам (один или несколько рабочих узлов). Если эта попытка провалится, будет выброшено исключение, и пользователь должен вместо этого использовать make_input_fn_iterator , которое предоставляет пользователю больший контроль и не пытается разделить пакет между репликами.
Пользователь также может использовать make_input_fn_iterator , если хочет настроить, какой вход подаётся на какую реплику/рабочий узел и т.д.
| Аргументы | |
|---|---|
dataset | tf.data.Dataset, который будет распределён равномерно по всем репликам. |
| Возвращаемое значение | |
|---|---|
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.
| Аргументы | |
|---|---|
input_fn | Функция, принимающая объект tf.distribute.InputContext и возвращающая tf.data.Dataset. |
replication_mode | значение перечисления tf.distribute.InputReplicationMode. В настоящее время поддерживается только PER_WORKER, что означает, что вызов input_fn будет выполнен один раз на каждый рабочий узел. Реплики будут извлекать данные из локального tf.data.Dataset на их рабочем узле. |
| Возвращаемое значение | |
|---|---|
Объект итератора, который сначала должен быть обработан .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 для уменьшения только по репликам (например, если тензор не имеет измерения пакета). |
| Возвращаемое значение | |
|---|---|
Значение Tensor. |
run
run(
fn, args=(), kwargs=None, options=None
)
Выполнить fn на каждой реплике с указанными аргументами.
Выполняет операции, указанные в fn, на каждой реплике. Если args или kwargs имеют значения «по реплике», такие как те, которые генерирует «распределённый Dataset», то при выполнении fn на конкретной реплике она будет выполняться с частью этих значений «по реплике», соответствующей этой реплике.
fn может вызывать tf.distribute.get_replica_context() для доступа к элементам, таким как all_reduce.
Все аргументы в args или kwargs должны быть либо вложенными тензорами, либо объектами «по реплике», содержащими тензоры или составные тензоры.
Пользователи могут передавать стратегические параметры в аргумент options.
Пример включения букетизации динамических форм в TPUStrategy.run:
resolver = tf.distribute.cluster_resolver.TPUClusterResolver(tpu='') tf.config.experimental_connect_to_cluster(resolver) tf.tpu.experimental.initialize_tpu_system(resolver) strategy = tf.distribute.experimental.TPUStrategy(resolver)
options = tf.distribute.RunOptions(
experimental_bucketizing_dynamic_shape=True)
dataset = tf.data.Dataset.range(
strategy.num_replicas_in_sync, output_type=dtypes.float32).batch(
strategy.num_replicas_in_sync, drop_remainder=True)
input_iterator = iter(strategy.experimental_distribute_dataset(dataset))
@tf.function() def step_fn(inputs): output = tf.reduce_sum(inputs) return output
strategy.run(step_fn, args=(next(input_iterator),), options=options)
| Аргументы | |
|---|---|
fn | Функция для выполнения. Выход должен быть tf.nest из Tensor. |
args | (Необязательно) Позиционные аргументы для fn. |
kwargs | (Необязательно) Именованные аргументы для fn. |
options | (Необязательно) Экземпляр tf.distribute.RunOptions, определяющий параметры для выполнения fn. |
| Возвращаемое значение | |
|---|---|
Объединённое возвращаемое значение fn по репликам. Структура возвращаемого значения такая же, как у возвращаемого значения из fn. Каждый элемент структуры может быть объектом «по реплике» 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, например, - Некоторые 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может иногда потребоваться внутри области, если она создаёт переменные.
strategy.run или model.fit для входа. Любые переменные, созданные вне области, не будут распределены и могут иметь последствия для производительности. Типичные вещи, создающие переменные в TF: модели, оптимизаторы, метрики. Они всегда должны создаваться внутри области. Другим источником создания переменных может быть восстановление контрольной точки — когда переменные создаются лениво. Обратите внимание, что любая переменная, созданная внутри стратегии, сохраняет информацию о стратегии. Чтение и запись этих переменных вне strategy.scope также могут работать безупречно, без необходимости ввода пользователя в область. | Возвращаемое значение | |
|---|---|
| Контекстный менеджер. |
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/TPUStrategy