tf.distribute.experimental.ParameterServerStrategy
| Просмотреть исходный код на GitHub |
Многоузловая стратегия tf.distribute с параметрическими серверами.
Наследуется от: Strategy
tf.distribute.experimental.ParameterServerStrategy(
cluster_resolver, variable_partitioner=None
)
Обучение с параметрическими серверами — это распространённый метод параллельного обучения на данных для масштабирования машинной обучающей модели на нескольких машинах. Кластер обучения с параметрическими серверами состоит из рабочих узлов и параметрических серверов. Переменные создаются на параметрических серверах, и рабочие узлы считывают и обновляют их на каждом шаге. По умолчанию рабочие узлы считывают и обновляют эти переменные независимо, не синхронизируясь друг с другом. В этом случае это известно как асинхронное обучение.
В TensorFlow 2 мы рекомендуем централизованную архитектуру на основе координации для обучения с параметрическими серверами, где рабочие узлы и параметрические серверы выполняют tf.distribute.Server, и есть ещё одна задача, которая создаёт ресурсы на рабочих узлах и параметрических серверах, отправляет функции и координирует обучение. Мы называем эту задачу «координатор». Координатор использует tf.distribute.experimental.coordinator.ClusterCoordinator для координации кластера и tf.distribute.experimental.ParameterServerStrategy для определения переменных на параметрических серверах и вычислений на рабочих узлах.
Для работы обучения координатор отправляет tf.function для выполнения на удалённых рабочих узлах. После получения запросов от координатора рабочий узел выполняет tf.function путём чтения переменных с параметрических серверов, выполнения операций и обновления переменных на параметрических серверах. Каждый из рабочих узлов обрабатывает только запросы от координатора и взаимодействует с параметрическими серверами без прямого взаимодействия с другими рабочими узлами в кластере.
В результате сбои некоторых рабочих узлов не препятствуют продолжению работы кластера, что позволяет кластеру обучаться с экземплярами, которые могут быть время от времени недоступны (например, предварительно выделенные или временные экземпляры). Однако координатор и параметрические серверы должны быть доступны в любое время для прогресса кластера.
Обратите внимание, что координатор не является одним из обучающих рабочих узлов. Вместо этого он создаёт ресурсы, такие как переменные и наборы данных, отправляет tf.function, сохраняет контрольные точки и т. д. Помимо рабочих узлов, параметрических серверов и координатора, можно запустить необязательный оценщик на стороне, который периодически считывает контрольные точки, сохранённые координатором, и выполняет оценки для каждой контрольной точки.
tf.distribute.experimental.ParameterServerStrategy должен работать в связке с объектом tf.distribute.experimental.coordinator.ClusterCoordinator. Автономное использование tf.distribute.experimental.ParameterServerStrategy без центральной координации в настоящее время не поддерживается.
Пример кода для координатора
Вот пример использования API с пользовательским циклом обучения для обучения модели. Этот фрагмент кода предназначен для выполнения на (единственной) задаче, назначенной координатором. Обратите внимание, что cluster_resolver, variable_partitioner, и dataset_fn аргументы объяснены в следующих разделах «Настройка кластера», «Разбиение переменных» и «Подготовка набора данных».
# Set the environment variable to allow reporting worker and ps failure to the
# coordinator. This a short-term workaround.
os.environ["GRPC_FAIL_FAST"] = "use_caller"
# 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()
Пример кода для рабочих узлов и параметрических серверов
Помимо координатора, должны быть задачи, назначенные как «рабочий узел» или «ps». Они должны выполнить следующий код, чтобы запустить сервер TensorFlow, ожидая запросов координатора:
# Set the environment variable to allow reporting worker and ps failure to the
# coordinator.
os.environ["GRPC_FAIL_FAST"] = "use_caller"
# 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. Обратите внимание, что по соображениям совместимости на некоторых платформах для типа задачи координатора используется «chief», как показано в следующем примере. Здесь мы устанавливаем 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
}
}
'''
Если вы предпочитаете запускать один и тот же двоичный файл для всех задач, вам нужно будет разделить двоичный файл на разные роли в начале программы:
os.environ["GRPC_FAIL_FAST"] = "use_caller"
cluster_resolver = tf.distribute.cluster_resolver.TFConfigClusterResolver()
# 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. Разбиение больших переменных между ps — это распространённый приём для повышения производительности обучения и устранения ограничений памяти. Это позволяет выполнять параллельные вычисления и обновления разных фрагментов переменной, а часто обеспечивает лучшую балансировку нагрузки между параметрическими серверами. Без фрагментации модели с большими переменными (например, вложениями), которые не помещаются в память одной машины, не могут обучаться.
С tf.distribute.experimental.ParameterServerStrategy, если variable_partitioner предоставляется __init__ и выполняются определённые условия, создаваемые в области переменные распределяются по параметрическим серверам в порядке круговой очереди. Ссылка на переменную, возвращённую из tf.Variable, становится типом, используемым как контейнер фрагментированных переменных. Можно получить доступ к атрибуту variables этого контейнера для фактических компонентов переменной. При построении модели с tf.Module или Keras компоненты переменной собираются в атрибутах, подобных variables.
class Dense(tf.Module):
def __init__(self, name=None):
super().__init__(name=name)
self.w = tf.Variable(tf.random.normal([100, 10]), name='w')
def __call__(self, x):
return x * self.w
# Partition the dense layer into 2 shards.
variable_partitioiner = (
tf.distribute.experimental.partitioners.FixedShardsPartitioner(
num_shards = 2))
strategy = ParameterServerStrategy(cluster_resolver=...,
variable_partitioner = variable_partitioner)
with strategy.scope():
dense = Dense()
assert len(dense.variables) == 2
assert isinstance(dense.variables[0], tf.Variable)
assert isinstance(dense.variables[1], tf.Variable)
assert dense.variables[0].name == "w/part_0"
assert dense.variables[1].name == "w/part_1"
Контейнер фрагментированной переменной может быть преобразован в 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 может быть изменён.tf.distribute.experimental.ParameterServerStrategyпока не поддерживает обучение с GPU. Это запрос на функцию, который разрабатывается.tf.distribute.experimental.ParameterServerStrategyв настоящее время поддерживает только API пользовательского цикла обучения в TF2. Использование с API Kerascompile/fitразрабатывается.tf.distribute.experimental.ParameterServerStrategyдолжен использоваться сtf.distribute.experimental.coordinator.ClusterCoordinator.
| Аргументы | |
|---|---|
cluster_resolver | объект tf.distribute.cluster_resolver.ClusterResolver. |
variable_partitioner | distribute.experimental.partitioners.Partitioner, определяющий способ разбиения переменных. Если None, переменные не будут разбиты.
|
| Атрибуты | |
|---|---|
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 возвращаемого этим 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 до нового размера пакета, равного глобальному размеру пакета, делённому на количество реплик в синхронизации. Мы проходим по нему с помощью цикла Python. 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>)
- Распределяет значения в массиве на основе 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.
Примечание: Это возвращает только значения на рабочем процессе, инициированном этим клиентом. При использованииtf.distribute.Strategy, например,tf.distribute.experimental.MultiWorkerMirroredStrategy, каждый рабочий процесс будет своим клиентом, и эта функция будет возвращать только значения, вычисленные на этом рабочем процессе.
| Аргументы | |
|---|---|
value | Значение, возвращенное experimental_run(), run(), extended.call_for_each_replica(), или переменной, созданной в scope . |
| Возвращаемое значение | |
|---|---|
Кортеж значений, содержащихся в value. Если value представляет единственное значение, возвращается (value,). |
gather
gather(
value, axis
)
Собрать value по репликам вдоль axis на текущее устройство.
Учитывая tf.distribute.DistributedValues или tf.Tensor-подобный объект value, этот API собирает и конкатенирует value по репликам вдоль axis измерения. Результат копируется на «текущее» устройство
- которым обычно является ЦП рабочего процесса, на котором выполняется программа. Для
tf.distribute.TPUStrategyэто первый хост TPU. Для многоклиентского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)>
| Аргументы | |
|---|---|
value | экземпляр tf.distribute.DistributedValues, например, возвращенный Strategy.run, для объединения в один тензор. Он также может быть обычным тензором при использовании с tf.distribute.OneDeviceStrategy или по умолчанию стратегией. Тензоры, составляющие DistributedValues, могут быть только плотными тензорами с ненулевым рангом, НЕ tf.IndexedSlices. |
axis | 0-мерный тензор int32. Измерение, по которому выполняется сбор. Должно находиться в диапазоне [0, rank(value)). |
| Возвращаемое значение | |
|---|---|
Tensor">value, представляющий конкатенацию 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, скажем, потерю на пример, пакет будет разделен по всем репликам. Эта функция позволяет агрегировать по репликам и необязательно также по элементам пакета, указав параметр оси соответствующим образом.
Например, если у вас есть глобальный размер пакета 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.
| Аргументы | |
|---|---|
reduce_op | значение tf.distribute.ReduceOp, указывающее, как следует объединять значения. Позволяет использовать строковое представление перечисления, такое как «SUM», «MEAN». |
value | экземпляр tf.distribute.DistributedValues, например, возвращаемый Strategy.run, который необходимо объединить в один тензор. Он также может быть обычным тензором, когда используется с OneDeviceStrategy или стратегией по умолчанию. |
axis | указывает размерность для сокращения вдоль тензора каждой реплики. Обычно следует устанавливать в размерность пакета или None для сокращения только по репликам (например, если тензор не имеет размерности пакета). |
| Возвращаемое значение | |
|---|---|
A 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 должны быть либо Python-значениями вложенной структуры тензоров, например, список тензоров, в этом случае args и kwargs будут переданы в fn , вызываемую на каждой реплике. Или args или kwargs могут быть tf.distribute.DistributedValues, содержащими тензоры или составные тензоры, т.е. tf.compat.v1.TensorInfo.CompositeTensor, в этом случае каждый fn вызов получит компонент tf.distribute.DistributedValues, соответствующий его реплике.
Пример использования:
- Ввод тензора константы.
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>
}
| Аргументы | |
|---|---|
fn | Функция, выполняемая на каждой реплике. |
args | Необязательные позиционные аргументы для fn. Его элемент может быть Python-значением, тензором или tf.distribute.DistributedValues. |
kwargs | Необязательные именованные аргументы для fn. Его элемент может быть Python-значением, тензором или tf.distribute.DistributedValues. |
options | Необязательный экземпляр tf.distribute.RunOptions, задающий параметры запуска fn. |
| Возвращаемое значение | |
|---|---|
Объединённое возвращаемое значение 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.
Что должно быть в области, а что вне?
Существует ряд требований к тому, что должно происходить внутри области. Однако в местах, где у нас есть информация о используемой стратегии, мы часто входим в область для пользователя, чтобы он не должен был делать это явно (т.е. вызов внутри или вне области допустим).
END_OF_DOCUMENT_MARKER- Все, что создает переменные, которые должны быть распределёнными, должно находиться в
strategy.scope. Это можно сделать, либо напрямую поместив их в область видимости, либо, используя другой API, напримерstrategy.runилиmodel.fit, для их помещения в область видимости. Любые переменные, созданные вне области видимости, не будут распределены и могут привести к снижению производительности. Общие вещи, которые создают переменные в TF: модели, оптимизаторы, метрики. Они всегда должны создаваться внутри области видимости. Другим источником создания переменных может быть восстановление контрольной точки — когда переменные создаются лениво. Обратите внимание, что любая переменная, созданная внутри стратегии, захватывает информацию о стратегии. Поэтому чтение и запись в эти переменные внеstrategy.scopeтакже могут работать бесперебойно, без необходимости пользователю входить в область видимости. - Некоторые API стратегий (например,
strategy.runиstrategy.reduce) требующие нахождения в области видимости стратегии, автоматически входят в область видимости, что означает, что при использовании этих API вам не нужно входить в область видимости самостоятельно. - Когда
tf.keras.Modelсоздается внутриstrategy.scope, эта информация фиксируется. Когда методы высокоуровневых обучающих фреймворков, такие какmodel.compile,model.fitи т. д., вызываются на этой модели, мы автоматически входим в область видимости и используем эту стратегию для распределения обучения и т. д. Подробный пример см. в учебнике по распределенному Keras. Обратите внимание, что простое вызовmodel(..)не затронуто — только высокоуровневые API фреймворков обучения.model.compile,model.fit,model.evaluate,model.predictиmodel.saveмогут вызываться внутри или вне области видимости. - Следующие действия могут выполняться как внутри, так и вне области видимости:
- Создание входных наборов данных
- Определение
tf.functionфункций, представляющих ваш шаг обучения - API сохранения, такие как
tf.saved_model.save. Загрузка создает переменные, поэтому это должно выполняться внутри области видимости, если вы хотите обучать модель в распределённом режиме. - Сохранение контрольных точек. Как упоминалось выше —
checkpoint.restoreиногда может потребоваться внутри области видимости, если это создаёт переменные.
| Возвращаемое значение | |
|---|---|
| Объект контекстного менеджера. |
© 2020 The TensorFlow Authors. All rights reserved.
Licensed under the Creative Commons Attribution License 3.0.
Code samples licensed under the Apache 2.0 License.
https://www.tensorflow.org/versions/r2.4/api_docs/python/tf/distribute/experimental/ParameterServerStrategy