tf.distribute.experimental.ParameterServerStrategy
Стратегия tf.distribute для нескольких рабочих узлов с параметрическими серверами.
Наследуется от: Strategy
tf.distribute.experimental.ParameterServerStrategy(
cluster_resolver: tf.distribute.cluster_resolver.ClusterResolver,
variable_partitioner: tf.distribute.experimental.partitioners.Partitioner = None
)
Используется в блокнотах
| Используется в руководстве | Используется в учебниках |
|---|---|
Обучение с параметрическим сервером — это распространенный метод параллельного обучения данных для масштабирования машинной обучающей модели на нескольких машинах. Кластер обучения с параметрическим сервером состоит из рабочих узлов и параметрических серверов. Переменные создаются на параметрических серверах, и они читаются и обновляются рабочими узлами на каждом шаге. По умолчанию рабочие узлы читают и обновляют эти переменные независимо, не синхронизируясь друг с другом. В этом случае это известно как асинхронное обучение.
В TensorFlow 2 мы рекомендуем архитектуру, основанную на централизованном координировании для обучения с параметрическим сервером. Каждый рабочий узел и параметрический сервер запускают tf.distribute.Server, а на его основе задача координатора отвечает за создание ресурсов на рабочих узлах и параметрических серверах, рассылку функций и координацию обучения. Координатор использует tf.distribute.experimental.coordinator.ClusterCoordinator для координации кластера и tf.distribute.experimental.ParameterServerStrategy для определения переменных на параметрических серверах и вычислений на рабочих узлах.
Для работы обучения координатор отправляет tf.function для выполнения на удаленных рабочих узлах. Получив запросы от координатора, рабочий узел выполняет tf.function путем чтения переменных с параметрических серверов, выполнения операций и обновления переменных на параметрических серверах. Каждый рабочий узел обрабатывает только запросы от координатора и взаимодействует с параметрическими серверами без прямого взаимодействия с другими рабочими узлами в кластере.
В результате сбои некоторых рабочих узлов не препятствуют продолжению работы кластера, что позволяет кластеру обучаться с экземплярами, которые могут быть периодически недоступны (например, предварительно выделенные или временные экземпляры). Однако координатор и параметрические серверы должны быть доступны в любое время для достижения прогресса кластером.
Обратите внимание, что координатор не является одним из рабочих узлов обучения. Вместо этого он создаёт ресурсы, такие как переменные и наборы данных, рассылает tf.function, сохраняет контрольные точки и т. д. Помимо рабочих узлов, параметрических серверов и координатора, можно запустить дополнительный оценочный сервер, который периодически читает контрольные точки, сохранённые координатором, и выполняет оценку для каждой контрольной точки.
ParameterServerStrategy поддерживается с двумя API обучения: Пользовательский цикл обучения (CTL) и API обучения Keras, также известный как Model.fit. CTL рекомендуется, когда пользователи предпочитают определять детали своего цикла обучения, а Model.fit рекомендуется, когда пользователи предпочитают высокоуровневую абстракцию и обработку обучения.
При использовании CTL, ParameterServerStrategy должен работать совместно с объектом tf.distribute.experimental.coordinator.ClusterCoordinator.
При использовании Model.fit, в настоящее время поддерживается только тип входных данных tf.keras.utils.experimental.DatasetCreator.
Пример кода для координатора
Этот раздел предоставляет фрагменты кода, которые предназначены для выполнения на единственной задаче, назначенной в качестве координатора. Обратите внимание, что аргументы cluster_resolver, variable_partitioner и dataset_fn объяснены в следующих разделах «Настройка кластера», «Разбиение переменных» и «Подготовка наборов данных».
При использовании CTL,
# Prepare a strategy to use with the cluster and variable partitioning info.
strategy = tf.distribute.experimental.ParameterServerStrategy(
cluster_resolver=...,
variable_partitioner=...)
coordinator = tf.distribute.experimental.coordinator.ClusterCoordinator(
strategy=strategy)
# Prepare a distribute dataset that will place datasets on the workers.
distributed_dataset = coordinator.create_per_worker_dataset(dataset_fn=...)
with strategy.scope():
model = ...
optimizer, metrics = ... # Keras optimizer/metrics are great choices
checkpoint = tf.train.Checkpoint(model=model, optimizer=optimizer)
checkpoint_manager = tf.train.CheckpointManager(
checkpoint, checkpoint_dir, max_to_keep=2)
# `load_checkpoint` infers initial epoch from `optimizer.iterations`.
initial_epoch = load_checkpoint(checkpoint_manager) or 0
@tf.function
def worker_fn(iterator):
def replica_fn(inputs):
batch_data, labels = inputs
# calculate gradient, applying gradient, metrics update etc.
strategy.run(replica_fn, args=(next(iterator),))
for epoch in range(initial_epoch, num_epoch):
distributed_iterator = iter(distributed_dataset) # Reset iterator state.
for step in range(steps_per_epoch):
# Asynchronously schedule the `worker_fn` to be executed on an arbitrary
# worker. This call returns immediately.
coordinator.schedule(worker_fn, args=(distributed_iterator,))
# `join` blocks until all scheduled `worker_fn`s finish execution. Once it
# returns, we can read the metrics and save checkpoints as needed.
coordinator.join()
logging.info('Metric result: %r', metrics.result())
train_accuracy.reset_states()
checkpoint_manager.save()
При использовании Model.fit,
# Prepare a strategy to use with the cluster and variable partitioning info.
strategy = tf.distribute.experimental.ParameterServerStrategy(
cluster_resolver=...,
variable_partitioner=...)
# A dataset function takes a `input_context` and returns a `Dataset`
def dataset_fn(input_context):
dataset = tf.data.Dataset.from_tensors(...)
return dataset.repeat().shard(...).batch(...).prefetch(...)
# With `Model.fit`, a `DatasetCreator` needs to be used.
input = tf.keras.utils.experimental.DatasetCreator(dataset_fn=...)
with strategy.scope():
model = ... # Make sure the `Model` is created within scope.
model.compile(optimizer="rmsprop", loss="mse", steps_per_execution=..., ...)
# Optional callbacks to checkpoint the model, back up the progress, etc.
callbacks = [tf.keras.callbacks.ModelCheckpoint(...), ...]
# `steps_per_epoch` is required with `ParameterServerStrategy`.
model.fit(input, epochs=..., steps_per_epoch=..., callbacks=callbacks)
Пример кода для рабочих узлов и параметрических серверов
Помимо координатора, должны быть задачи, назначенные как «рабочий узел» или «ps». Они должны выполнить следующий код, чтобы запустить сервер TensorFlow, ожидая запросов координатора:
# Provide a `tf.distribute.cluster_resolver.ClusterResolver` that serves
# the cluster information. See below "Cluster setup" section.
cluster_resolver = ...
server = tf.distribute.Server(
cluster_resolver.cluster_spec(),
job_name=cluster_resolver.task_type,
task_index=cluster_resolver.task_id,
protocol="grpc")
# Blocking the process that starts a server from exiting.
server.join()
Настройка кластера
Для того, чтобы задачи в кластере знали адреса других задач, требуется использовать tf.distribute.cluster_resolver.ClusterResolver в координаторе, рабочем узле и ps. tf.distribute.cluster_resolver.ClusterResolver отвечает за предоставление информации о кластере, а также тип и идентификатор задачи текущей задачи. Подробности см. в tf.distribute.cluster_resolver.ClusterResolver.
Если переменная окружения TF_CONFIG установлена, необходимо также использовать tf.distribute.cluster_resolver.TFConfigClusterResolver.
Поскольку в tf.distribute.experimental.ParameterServerStrategy есть предположения относительно именования типов задач, «главный», «ps» и «рабочий узел» должны использоваться в tf.distribute.cluster_resolver.ClusterResolver для обозначения координатора, параметрических серверов и рабочих узлов соответственно.
Следующий пример демонстрирует установку TF_CONFIG для задачи, назначенной в качестве параметрического сервера (тип задачи «ps») и индекса 1 (вторая задача) в кластере с 1 главным, 2 параметрическими серверами и 3 рабочими узлами. Обратите внимание, что это должно быть установлено до использования tf.distribute.cluster_resolver.TFConfigClusterResolver.
Пример кода для настройки кластера:
os.environ['TF_CONFIG'] = '''
{
"cluster": {
"chief": ["chief.example.com:2222"],
"ps": ["ps0.example.com:2222", "ps1.example.com:2222"],
"worker": ["worker0.example.com:2222", "worker1.example.com:2222",
"worker2.example.com:2222"]
},
"task": {
"type": "ps",
"index": 1
}
}
'''
Если вы предпочитаете запускать один и тот же бинарник для всех задач, вам нужно будет разделить его на разные роли в начале программы:
# If coordinator, create a strategy and start the training program.
if cluster_resolver.task_type == 'chief':
strategy = tf.distribute.experimental.ParameterServerStrategy(
cluster_resolver)
...
# If worker/ps, create a server
elif cluster_resolver.task_type in ("worker", "ps"):
server = tf.distribute.Server(...)
...
В качестве альтернативы, вы также можете заранее запустить несколько серверов TensorFlow и подключиться к ним позже. Координатор может находиться в том же кластере или на любой машине, имеющей подключение к рабочим узлам и параметрическим серверам. Это описано в нашем руководстве и учебнике.
Создание переменных с помощью strategy.scope()
tf.distribute.experimental.ParameterServerStrategy следует контракту API tf.distribute, где ожидается создание переменной внутри контекстного менеджера, возвращаемого strategy.scope(), для правильного размещения на параметрических серверах по круговому методу:
# In this example, we're assuming having 3 ps.
strategy = tf.distribute.experimental.ParameterServerStrategy(
cluster_resolver=...)
coordinator = tf.distribute.experimental.coordinator.ClusterCoordinator(
strategy=strategy)
# Variables should be created inside scope to be placed on parameter servers.
# If created outside scope such as `v1` here, it would be placed on the
# coordinator.
v1 = tf.Variable(initial_value=0.0)
with strategy.scope():
v2 = tf.Variable(initial_value=1.0)
v3 = tf.Variable(initial_value=2.0)
v4 = tf.Variable(initial_value=3.0)
v5 = tf.Variable(initial_value=4.0)
# v2 through v5 are created in scope and are distributed on parameter servers.
# Default placement is round-robin but the order should not be relied on.
assert v2.device == "/job:ps/replica:0/task:0/device:CPU:0"
assert v3.device == "/job:ps/replica:0/task:1/device:CPU:0"
assert v4.device == "/job:ps/replica:0/task:2/device:CPU:0"
assert v5.device == "/job:ps/replica:0/task:0/device:CPU:0"
См. distribute.Strategy.scope для получения дополнительной информации.
Разбиение переменных
Использование выделенных серверов для хранения переменных позволяет разделить или «разбить» переменные по параметрическим серверам. Разбиение больших переменных между ps — это распространенный метод для повышения скорости обучения и уменьшения требований к памяти. Это позволяет параллельно выполнять вычисления и обновления на разных фрагментах переменной, что часто обеспечивает лучшее балансирование нагрузки между параметрическими серверами. Без разбиения модели с большими переменными (например, вложения), которые не помещаются в память одной машины, в противном случае не смогут обучаться.
С помощью tf.distribute.experimental.ParameterServerStrategy, если variable_partitioner предоставлен для __init__ и выполнены определенные условия, полученные переменные в области разделяются по параметрическим серверам по круговому методу. Ссылка на переменную, возвращенная из tf.Variable, становится типом, который служит контейнером разнесенных переменных. Для доступа к фактическим компонентам переменной можно использовать атрибут variables этого контейнера. Если модель создается с tf.Module или Keras, компоненты переменной собираются в атрибутах типа variables.
Рекомендуется использовать партиционеры на основе размера, такие как tf.distribute.experimental.partitioners.MinSizePartitioner, чтобы избежать разбиения малых переменных, что может отрицательно сказаться на скорости обучения модели.
# Partition the embedding layer into 2 shards.
variable_partitioner = (
tf.distribute.experimental.partitioners.MinSizePartitioner(
min_shard_bytes=(256 << 10),
max_shards = 2))
strategy = tf.distribute.experimental.ParameterServerStrategy(
cluster_resolver=...,
variable_partitioner = variable_partitioner)
with strategy.scope():
embedding = tf.keras.layers.Embedding(input_dim=1024, output_dim=1024)
assert len(embedding.variables) == 2
assert isinstance(embedding.variables[0], tf.Variable)
assert isinstance(embedding.variables[1], tf.Variable)
assert embedding.variables[0].shape == (512, 1024)
assert embedding.variables[1].shape == (512, 1024)
Контейнер разнесенной переменной можно преобразовать в Tensor с помощью tf.convert_to_tensor. Это означает, что контейнер можно напрямую использовать в большинстве операций Python, где происходит автоматическое преобразование в Tensor. Например, в приведенном выше фрагменте кода x * self.w неявно применит указанное преобразование тензора. Обратите внимание, что такое преобразование может быть дорогостоящим, так как компоненты переменных должны быть перенесены с нескольких параметрических серверов в место использования значения.
tf.nn.embedding_lookup, с другой стороны, не применяет преобразование тензора и вместо этого выполняет параллельные запросы к компонентам переменной. Это важно для масштабирования запросов к вложениям, когда переменная таблицы вложений большая.
Когда разнесенная переменная сохраняется в SavedModel, она будет сохранена как одна переменная. Это повышает эффективность обслуживания, устраняя ряд операций, которые обрабатывают аспекты разбиения.
Известные ограничения разбиения переменных:
Количество разделов не должно изменяться при сохранении/загрузке контрольных точек.
После сохранения разбиением переменных в SavedModel, SavedModel не может быть загружен с помощью
tf.saved_model.load.Переменная разбиения напрямую не работает с
tf.GradientTape, пожалуйста, используйте атрибутыvariables, чтобы получить фактические компоненты переменной и использовать их в API градиента.
Подготовка набора данных
С помощью tf.distribute.experimental.ParameterServerStrategy, набор данных создается на каждом из рабочих узлов для использования в процессе обучения. Это делается путем создания dataset_fn, который не принимает аргументов и возвращает tf.data.Dataset, и передачи dataset_fn в tf.distribute.experimental.coordinator. ClusterCoordinator.create_per_worker_dataset. Рекомендуется, чтобы набор данных был перемешан и повторялся, чтобы примеры проходили обучение как можно равномернее.
def dataset_fn():
filenames = ...
dataset = tf.data.Dataset.from_tensor_slices(filenames)
# Dataset is recommended to be shuffled, and repeated.
return dataset.shuffle(buffer_size=...).repeat().batch(batch_size=...)
coordinator =
tf.distribute.experimental.coordinator.ClusterCoordinator(strategy=...)
distributed_dataset = coordinator.create_per_worker_dataset(dataset_fn)
Ограничения
tf.distribute.experimental.ParameterServerStrategyв TF2 является экспериментальным, и API может подвергаться дальнейшим изменениям.При использовании
Model.fit,tf.distribute.experimental.ParameterServerStrategyдолжен использоваться сtf.keras.utils.experimental.DatasetCreator, иsteps_per_epochдолжен быть указан.
| Аргументы | |
|---|---|
cluster_resolver | объект tf.distribute.cluster_resolver.ClusterResolver. |
variable_partitioner | distribute.experimental.partitioners.Partitioner, который определяет, как разбить переменные. Если None, переменные не будут разбиты.
|
| Атрибуты | |
|---|---|
cluster_resolver | Возвращает резольвер кластера, связанный с этой стратегией. В общем случае при использовании многоузловой Стратегии, которые намерены иметь связанный Одноузловые стратегии обычно не имеют The
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.distribute.DistributedDataset из tf.data.Dataset.
Возвращаемый tf.distribute.DistributedDataset можно перебирать аналогично обычным наборам данных. ВАЖНО: пользователь не может добавлять больше преобразований в tf.distribute.DistributedDataset. Вы можете только создать итератор или проверить tf.TypeSpec сгенерированных им данных. См. документацию API tf.distribute.DistributedDataset для получения дополнительной информации.
Следующий пример:
global_batch_size = 2
# Passing the devices is optional.
strategy = tf.distribute.MirroredStrategy(devices=["GPU:0", "GPU:1"])
# Create a dataset
dataset = tf.data.Dataset.range(4).batch(global_batch_size)
# Distribute that dataset
dist_dataset = strategy.experimental_distribute_dataset(dataset)
@tf.function
def replica_fn(input):
return input*2
result = []
# Iterate over the `tf.distribute.DistributedDataset`
for x in dist_dataset:
# process dataset elements
result.append(strategy.run(replica_fn, args=(x,)))
print(result)
[PerReplica:{
0: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([0])>,
1: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([2])>
}, PerReplica:{
0: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([4])>,
1: <tf.Tensor: shape=(1,), dtype=int64, numpy=array([6])>
}]Три ключевых действия, происходящие за кулисами этого метода, — это пакетная обработка, разбиение и предварительная обработка.
В приведенном выше фрагменте кода dataset сгруппирован по global_batch_size, и вызов experimental_distribute_dataset на нём перегруппировывает dataset в новую размерность пакета, равную глобальной размерности пакета, делённой на количество реплик в синхронизации. Мы итерируем по нему с помощью питоновского цикла for. x — это tf.distribute.DistributedValues, содержащий данные для всех реплик, и каждая реплика получает данные новой размерности пакета. tf.distribute.Strategy.run позаботится о подаче правильных данных для каждой реплики в x к соответствующему replica_fn, выполняемому на каждой реплике.
Разбиение данных включает автоматическое разбиение данных по нескольким рабочим процессам и внутри каждого рабочего процесса. Во-первых, при распределённом обучении с несколькими рабочими процессами (то есть когда вы используете tf.distribute.experimental.MultiWorkerMirroredStrategy или tf.distribute.TPUStrategy), автоматическое разбиение набора данных по рабочим процессам означает, что каждому рабочему процессу назначается подмножество всего набора данных (если установлен правильный tf.data.experimental.AutoShardPolicy). Это делается для того, чтобы на каждом шаге каждый рабочий процесс обрабатывал глобальный пакет элементов набора данных без пересечений. Автоматическое разбиение имеет несколько вариантов, которые можно указать с помощью tf.data.experimental.DistributeOptions. Затем, разбиение внутри каждого рабочего процесса означает, что метод разделит данные между всеми устройствами рабочего процесса (если их несколько). Это произойдёт независимо от автоматического разбиения данных по нескольким рабочим процессам.
Примечание: по умолчанию режим автоматического разбиения по нескольким рабочим процессам —tf.data.experimental.AutoShardPolicy.AUTO. Этот режим попытается разбить входной набор данных по файлам, если набор данных создаётся из наборов данных для чтения (например,tf.data.TFRecordDataset,tf.data.TextLineDatasetи т. д.), или иначе разделить набор данных по данным, где каждый из рабочих процессов прочитает весь набор данных и обработает только выделенный ему фрагмент. Однако если у вас меньше одного входного файла на рабочий процесс, мы рекомендуем отключить автоматическое разбиение набора данных по рабочим процессам, установивtf.data.experimental.DistributeOptions.auto_shard_policyнаtf.data.experimental.AutoShardPolicy.OFF.
По умолчанию этот метод добавляет преобразование предварительной выборки в конце предоставленного пользователем экземпляра tf.data.Dataset. Аргумент преобразования предварительной выборки, который является buffer_size, равен количеству реплик в синхронизации.
Если описанная выше логика разбиения пакетов и разбиения набора данных нежелательна, используйте вместо неё tf.distribute.Strategy.distribute_datasets_from_function, которая не выполняет автоматического пакетного разбиения или разбиения набора данных.
Примечание: Если вы используете 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 | tf.data.Dataset, который будет распределён по всем репликам в соответствии с указанными выше правилами. |
options | tf.distribute.InputOptions, используемый для управления параметрами распределения этого набора данных. |
| Возвращаемое значение | |
|---|---|
tf.distribute.DistributedDataset. |
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.
Примечание: Это возвращает только значения на рабочем процессе, инициированном этим клиентом. При использованииtf.distribute.Strategy, например,tf.distribute.experimental.MultiWorkerMirroredStrategy, каждый рабочий процесс будет собственным клиентом, и эта функция вернёт только вычисленные на нём значения.
| Аргументы | |
|---|---|
value | Значение, возвращённое experimental_run(), run(), or a variable created inscope`. |
| Возвращаемое значение | |
|---|---|
Кортеж значений, содержащихся в value, где i-й элемент соответствует i-й реплике. Если 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 | тензор int32 размерности 0. Размерность, по которой необходимо собрать. Должно быть в диапазоне [0, rank(value)). |
| Returns | |
|---|---|
Tensor, являющийся конкатенацией value по всем репликам по размерности axis. |
reduce
reduce(
reduce_op, value, axis
)
Свести value по всем репликам и вернуть результат на текущем устройстве.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"])
def step_fn():
i = tf.distribute.get_replica_context().replica_id_in_sync_group
return tf.identity(i)
per_replica_result = strategy.run(step_fn)
total = strategy.reduce("SUM", per_replica_result, axis=None)
total
<tf.Tensor: shape=(), dtype=int32, numpy=1>Чтобы увидеть, как это будет выглядеть с несколькими репликами, рассмотрим тот же пример с MirroredStrategy и 2 графическими процессорами:
strategy = tf.distribute.MirroredStrategy(devices=["GPU:0", "GPU:1"])
def step_fn():
i = tf.distribute.get_replica_context().replica_id_in_sync_group
return tf.identity(i)
per_replica_result = strategy.run(step_fn)
# Check devices on which per replica result is:
strategy.experimental_local_results(per_replica_result)[0].device
# /job:localhost/replica:0/task:0/device:GPU:0
strategy.experimental_local_results(per_replica_result)[1].device
# /job:localhost/replica:0/task:0/device:GPU:1
total = strategy.reduce("SUM", per_replica_result, axis=None)
# Check device on which reduced result is:
total.device
# /job:localhost/replica:0/task:0/device:CPU:0
Этот API обычно используется для агрегирования результатов, возвращаемых различными репликами, например, для отчётности. Например, потерю, вычисленную из различных реплик, можно усреднить с помощью этого API перед выводом.
Примечание: Результат копируется на «текущее» устройство — обычно это ЦП узла, на котором выполняется программа. ДляTPUStrategy— это первый хост TPU. Для многоклиентскихMultiWorkerMirroredStrategy— это ЦП каждого узла.
Существует ряд различных API tf.distribute для сведения значений по всем репликам:
-
tf.distribute.ReplicaContext.all_reduce: Это отличается отStrategy.reduceтем, что предназначено для контекста реплики и не копирует результаты на устройство хоста.all_reduceобычно используется для вычислений внутри этапа обучения, таких как градиенты. -
tf.distribute.StrategyExtended.reduce_toиtf.distribute.StrategyExtended.batch_reduce_to: Эти API являются более расширенными версиямиStrategy.reduce, так как они позволяют настраивать место назначения результата. Они также вызываются в кросс-репликационном контексте.
Каким должен быть ось?
Учитывая значение, возвращённое по каждой реплике run, например, потери на пример, пакет будет разделен на все реплики. Эта функция позволяет агрегировать по репликам и, по желанию, также по элементам пакета, указав параметр axis соответствующим образом.
Например, если у вас есть глобальный размер пакета 8 и 2 реплики, значения для примеров [0, 1, 2, 3] будут на реплике 0, а [4, 5, 6, 7] — на реплике 1. С помощью axis=None, reduce будет агрегировать только по репликам, возвращая [0+4, 1+5, 2+6, 3+7]. Это полезно, когда каждая реплика вычисляет скаляр или какое-либо другое значение без размерности «пакет» (например, градиент или потерю).
strategy.reduce("sum", per_replica_result, axis=None)
Иногда вы захотите агрегировать как по глобальному пакету, так и по всем репликам. Это поведение можно получить, указав размерность пакета в качестве axis, обычно axis=0. В этом случае он вернёт скаляр 0+1+2+3+4+5+6+7.
strategy.reduce("sum", per_replica_result, axis=0)
Если есть последний частичный пакет, вам нужно указать ось, чтобы размерность результата была согласована между репликами. Таким образом, если последний пакет имеет размер 6 и он разделен на [0, 1, 2, 3] и [4, 5], у вас возникнет несоответствие размеров, если вы не укажете axis=0. Если вы укажете tf.distribute.ReduceOp.MEAN, используя axis=0, будет использован правильный знаменатель 6. Противопоставьте это вычислению reduce_mean для получения скалярного значения на каждой реплике и этой функции для усреднения этих средних значений, что будет взвешивать некоторые значения 1/8, а другие 1/4.
| Args | |
|---|---|
reduce_op | значение tf.distribute.ReduceOp, определяющее, как должны комбинироваться значения. Разрешает использовать строковое представление перечисления, например, «SUM», «MEAN». |
value | экземпляр tf.distribute.DistributedValues, например, возвращённый Strategy.run, который необходимо объединить в один тензор. Он также может быть обычным тензором, когда используется с OneDeviceStrategy или стратегией по умолчанию. |
axis | определяет размерность, по которой нужно производить свёртку в тензоре каждой реплики. Обычно устанавливается в размерность пакета или None, чтобы выполнить свёртку только по репликам (например, если тензор не имеет размерности пакета). |
| Returns | |
|---|---|
Tensor. |
run
run(
fn, args=(), kwargs=None, options=None
)
Вызывает fn на каждой реплике с указанными аргументами.
Этот метод является основным способом распределения вычислений с помощью объекта tf.distribute. Он вызывает fn на каждой реплике. Если args или kwargs содержат tf.distribute.DistributedValues, например, те, которые создаются tf.distribute.DistributedDataset из tf.distribute.Strategy.experimental_distribute_dataset или tf.distribute.Strategy.distribute_datasets_from_function, при выполнении fn на конкретной реплике, он будет выполнен с компонентом tf.distribute.DistributedValues, соответствующим этой реплике.
fn вызывается в контексте реплики. fn может вызвать tf.distribute.get_replica_context() для доступа к элементам, таким как all_reduce. С концепцией контекста реплики ознакомьтесь в документации модуля tf.distribute.
Все аргументы в args или kwargs могут быть вложенной структурой тензоров, например, списком тензоров, в этом случае args и kwargs будут переданы вызываемому fn на каждой реплике. Или args или kwargs могут быть tf.distribute.DistributedValues, содержащими тензоры или составные тензоры, то есть tf.compat.v1.TensorInfo.CompositeTensor, в этом случае каждый вызов fn получит компонент tf.distribute.DistributedValues, соответствующий своей реплике. Обратите внимание, что произвольные значения Python, которые не являются указанных типов, не поддерживаются.
Пример использования:
-
Входной тензор константы.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"]) 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 PerReplica:{ 0: <tf.Tensor: shape=(), dtype=float32, numpy=6.0>, 1: <tf.Tensor: shape=(), dtype=float32, numpy=6.0> } -
Входной DistributedValues.
strategy = tf.distribute.MirroredStrategy(["GPU:0", "GPU:1"]) @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=4> -
Использование
tf.distribute.ReplicaContextдля allreduce значений.strategy = tf.distribute.MirroredStrategy(["gpu:0", "gpu:1"]) @tf.function def run(): def value_fn(value_context): return tf.constant(value_context.replica_id_in_sync_group) distributed_values = ( strategy.experimental_distribute_values_from_function( value_fn)) def replica_fn(input): return tf.distribute.get_replica_context().all_reduce( "sum", input) return strategy.run(replica_fn, args=(distributed_values,)) result = run() result PerReplica:{ 0: <tf.Tensor: shape=(), dtype=int32, numpy=1>, 1: <tf.Tensor: shape=(), dtype=int32, numpy=1> }
| Args | |
|---|---|
fn | Функция для выполнения на каждой реплике. |
args | Необязательные позиционные аргументы для fn. Его элементы могут быть тензором, вложенной структурой тензоров или tf.distribute.DistributedValues. |
kwargs | Необязательные именованные аргументы для fn. Его элементы могут быть тензором, вложенной структурой тензоров или tf.distribute.DistributedValues. |
options | Необязательный экземпляр tf.distribute.RunOptions, определяющий параметры выполнения fn. |
| Returns | |
|---|---|
Объединённое возвращаемое значение fn по всем репликам. Структура возвращаемого значения такая же, как и у возвращаемого значения fn. Каждый элемент структуры может быть tf.distribute.DistributedValues, объектом Tensor или Tensor (например, при выполнении на одной реплике). |
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, вам необходимо использовать APIstrategy.run, чтобы явно распределить это вычисление. См. пример в руководстве по пользовательскому циклу обучения 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/ParameterServerStrategy