tf.distribute.experimental.CentralStorageStrategy
Стратегия для одной машины, которая помещает все переменные на один узел.
Наследуется от: Strategy
tf.distribute.experimental.CentralStorageStrategy(
compute_devices=None, parameter_device=None
)
Использование в ноутбуках
| Используется в руководстве |
|---|
Переменные назначаются локальному процессору CPU или единственному графическому процессору GPU. Если доступно более одного графического процессора GPU, операции вычислений (кроме операций обновления переменных) будут дублироваться по всем графическим процессорам GPU.
Например:
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 будет вызываться на устройстве CPU каждого из узлов, и каждый из них создаст набор данных, где каждая реплика на этом узле будет извлекать одну партию ввода (т. е. если на узле две реплики, две партии будут извлекаться из Dataset на каждом шаге).
Этот метод можно использовать для нескольких целей. Во-первых, он позволяет указать собственную логику группировки и разделения. (В отличие от tf.distribute.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. См. 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],...
Args: 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 или другие методы, принимающие распределённые значения, когда не используются наборы данных.
| Args | |
|---|---|
value_fn | Функция, которая выполняется для генерации значений. Она вызывается для каждой реплики с tf.distribute.ValueContext в качестве единственного аргумента. Она должна возвращать тензор или тип, который может быть преобразован в тензор. |
| Returns | |
|---|---|
A 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>) -
Распределение значений в массиве на основе id реплики:
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-й размерности. Результат копируется на "текущее" устройство, которое обычно является процессором узла, на котором выполняется программа. Для tf.distribute.TPUStrategy это первый хост TPU. Для многоклиентской tf.distribute.MultiWorkerMirroredStrategy это процессор каждого рабочего узла.
Этот 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 | |
|---|---|
A 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 | |
|---|---|
A Tensor. |
run
run(
fn, args=(), kwargs=None, options=None
)
Выполнить fn на каждой реплике с заданными аргументами.
В CentralStorageStrategy, fn вызывается на каждой вычислительной реплике с предоставленными аргументами "для каждой реплики", специфичными для этого устройства.
| Args | |
|---|---|
fn | Функция, которую нужно выполнить. Вывод должен быть tf.nest тензоров. |
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()теперь вернёт эту стратегию. За пределами этого контекста, она возвращает стратегию по умолчанию без действий. - Вход в контекст также вводит «меж-репличный контекст». Обратитесь к
tf.distribute.StrategyExtendedдля объяснения меж-репличных и репличных контекстов. - Создание переменных внутри
scopeперехватывается стратегией. Каждая стратегия определяет, как она хочет повлиять на создание переменных. Синхронные стратегии, такие какMirroredStrategy,TPUStrategyиMultiWorkerMiroredStrategy, создают переменные, дублированные на каждой реплике, в то время какParameterServerStrategyсоздаёт переменные на серверах параметров. Это делается с помощью пользовательскогоtf.variable_creator_scope. - В некоторых стратегиях также может быть введён контекст устройства по умолчанию: в
MultiWorkerMiroredStrategyна каждом рабочем узле вводится контекст устройства по умолчанию «/CPU:0».
Примечание: Вход в контекст не автоматически распределяет вычисления, за исключением случаев высокоуровневых обучающих фреймворков, таких как kerasmodel.fit. Если вы не используетеmodel.fit, вам нужно использоватьstrategy.runAPI для явного распределения вычислений. См. пример в руководстве по пользовательским циклам обучения custom training loop tutorial.
Что должно быть в контексте, а что вне его?
Существует ряд требований к тому, что должно происходить внутри контекста. Однако в тех местах, где у нас есть информация о используемой стратегии, мы часто входим в контекст за пользователя, чтобы они не должны были делать это явно (т.е. вызов внутри или вне контекста допустим).
- Любой код, создающий переменные, которые должны быть распределёнными переменными, должен вызываться в контексте
strategy.scope. Этого можно добиться, либо напрямую вызвав функцию создания переменной в контексте контекста, либо используя другой API, такой какstrategy.runилиkeras.Model.fit, чтобы автоматически войти в него за вас. Любая переменная, созданная вне контекста, не будет распределена и может повлиять на производительность. Некоторые распространённые объекты, создающие переменные в TF, — это модели, оптимизаторы, метрики. Такие объекты всегда должны быть инициализированы в контексте, а любые функции, которые могут лениво создавать переменные (например,Model.call(), отслеживаниеtf.functionи т.д.), аналогичным образом должны быть вызваны в контексте. Ещё одним источником создания переменных может быть восстановление контрольной точки — когда переменные создаются лениво. Обратите внимание, что любая переменная, созданная внутри стратегии, фиксирует информацию о стратегии. Поэтому чтение и запись в эти переменные за пределами контекстаstrategy.scopeтакже могут работать беспрепятственно, без необходимости входа пользователя в контекст. - Некоторые API стратегий (например,
strategy.runиstrategy.reduce), которые требуют нахождения в контексте стратегии, входят в контекст автоматически, что означает, что при использовании этих API вам не нужно явно входить в контекст. - Когда
tf.keras.Modelсоздаётся внутриstrategy.scope, объект модели фиксирует информацию о контексте. Когда затем вызываются методы высокоуровневых обучающих фреймворков, такие какmodel.compile,model.fitи т.д., захваченный контекст будет автоматически введён, и соответствующая стратегия будет использоваться для распределения обучения и т.д. Подробный пример см. в distributed keras tutorial. ВНИМАНИЕ: Простое вызовmodel(..)не автоматически входит в захваченный контекст — только высокоуровневые API обучающих фреймворков поддерживают это поведение:model.compile,model.fit,model.evaluate,model.predictиmodel.saveмогут быть вызваны как внутри, так и вне контекста. - Следующее может быть как внутри, так и вне контекста:
- Создание наборов данных для ввода
- Определение
tf.function, которые представляют ваш шаг обучения - Сохранение API, такие как
tf.saved_model.save. Загрузка создаёт переменные, поэтому это должно выполняться внутри контекста, если вы хотите обучить модель в распределённом режиме. - Сохранение контрольных точек. Как указано выше —
checkpoint.restoreиногда может потребоваться внутри контекста, если он создаёт переменные.
| Возвращаемое значение | |
|---|---|
| Менеджер контекста. |
© 2022 The TensorFlow Authors. All rights reserved.
Licensed under the Creative Commons Attribution License 4.0.
Code samples licensed under the Apache 2.0 License.
https://www.tensorflow.org/api_docs/python/tf/distribute/experimental/CentralStorageStrategy