Spec-Zone.ru › TensorFlow 2.9

tf.distribute.experimental.ParameterServerStrategy

Просмотреть исходный код на GitHub

Многоузловая стратегия tf.distribute с параметрическими серверами.

Наследуется от: Strategy

Просмотр псевдонимов

Основные псевдонимы

tf.distribute.ParameterServerStrategy

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, сохраняет контрольные точки и т. д. Помимо рабочих узлов, параметрических серверов и координатора, можно запустить необязательный оценщик, который периодически считывает контрольные точки, сохранённые координатором, и выполняет оценки по каждой контрольной точке.

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. Разбиение больших переменных между 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, переменные не будут разбиваться.
  • Для этого аргумента можно использовать предопределенные разбиерители из tf.distribute.experimental.partitioners. Часто используемый разбиеритель — MinSizePartitioner(min_shard_bytes = 256 << 10, max_shards = num_ps), который выделяет не менее 256 КБ на фрагмент, а каждый ps получает не более одного фрагмента.

  • variable_partitioner будет вызываться для каждой созданной переменной в рамках стратегии scope для указания способа разбиения переменной. Переменные, имеющие только один фрагмент по оси разбиения (т. е. не требующие разбиения), будут созданы как обычная tf.Variable.

  • Поддерживается только разбиение по первой/внешней оси.

  • Для разбиения переменных используется стратегия разбиения по делениям. Предполагая, что мы назначаем последовательные целочисленные идентификаторы вдоль первой оси переменной, идентификаторы назначаются фрагментам непрерывно, при этом пытаясь сохранить размер каждого фрагмента одинаковым. Если идентификаторы не делятся на количество фрагментов без остатка, каждому из первых нескольких фрагментов будет назначен на один идентификатор больше. Например, переменная, первая размерность которой равна 13, имеет 13 идентификаторов, и они разбиваются на 5 фрагментов следующим образом: [[0, 1, 2], [3, 4, 5], [6, 7, 8], [9, 10], [11, 12]].

  • Переменные, созданные в рамках strategy.extended.colocate_vars_with не будут разбиваться.

Атрибуты
cluster_resolver Возвращает связанный с этой стратегией кластер-разрешитель.

В общем случае при использовании стратегии многоуровневой обработки tf.distribute, такой как tf.distribute.experimental.MultiWorkerMirroredStrategy или tf.distribute.TPUStrategy(), с используемой стратегией связан tf.distribute.cluster_resolver.ClusterResolver, и такой экземпляр возвращается этим свойством.

Стратегии, которые планируют иметь связанный tf.distribute.cluster_resolver.ClusterResolver, должны установить соответствующий атрибут или переопределить это свойство; в противном случае, по умолчанию возвращается None. Эти стратегии также должны предоставить информацию о том, что возвращается этим свойством.

Стратегии с одним рабочим процессом обычно не имеют tf.distribute.cluster_resolver.ClusterResolver, и в таких случаях это свойство вернет None.

Разрешитель кластера tf.distribute.cluster_resolver.ClusterResolver может быть полезен, когда пользователю нужно получить доступ к информации, такой как описание кластера, тип задачи или идентификатор задачи. Например,

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 tf.distribute.cluster_resolver.ClusterResolver.

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, чтобы увидеть пример.

Важно: Возвращаемый dataset_fn tf.data.Dataset должен иметь размер пакета на реплику, в отличие от experimental_distribute_dataset, который использует глобальный размер пакета. Это может быть вычислено с помощью input_context.get_per_replica_batch_size.
Примечание: Если вы используете 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, содержащий значение для каждой реплики.

Пример использования:

  1. Возвращает постоянное значение для каждой реплики:
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>)
  1. Распределяет значения в массиве на основе идентификатора реплики:
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)
  1. Указывает значения с помощью 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)
  1. Размещает значения на устройствах и распределяет:
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.DistributedValues value, компоненты тензоров должны иметь ненулевой ранг. В противном случае рассмотрите использование 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, ранг(значение)).
Возвращаемое значение
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 перед выводом.

Примечание: Результат копируется на устройство "текущее" — обычно это процессор (CPU) рабочей машины, на которой выполняется программа. Для TPUStrategy, это первый хост TPU. Для многоклиентских MultiWorkerMirroredStrategy, это CPU каждого рабочего узла.

Существует ряд различных 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, чтобы сократить только по репликам (например, если тензор не имеет измерения пакета).
Возвращаемое значение
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, которые не относятся к указанным выше типам, не поддерживаются.

Важно: В зависимости от реализации tf.distribute.Strategy и от того, включена ли жадная (eager) обработка, fn может вызываться один или несколько раз. Если fn помечена аннотацией tf.function или tf.distribute.Strategy.run вызывается внутри tf.function (жадная обработка отключена внутри tf.function по умолчанию), fn вызывается один раз на каждую реплику для создания графа TensorFlow, который затем будет повторно использоваться для выполнения с новыми входами. В противном случае, если жадная обработка включена, fn будет вызываться один раз на каждую реплику на каждом шаге, как и обычный Python-код.

Пример использования:

  1. Входной тензор с постоянными значениями.
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>
}
  1. Входные данные 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>
  1. Использование 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. Его элементы могут быть тензорами, вложенными структурами тензоров или tf.distribute.DistributedValues.
kwargs Необязательные ключевые аргументы для fn. Его элементы могут быть тензорами, вложенными структурами тензоров или 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" вводится на каждом рабочем узле.
Примечание: Вход в область не автоматически распределяет вычисления, за исключением случаев использования высокоуровневых обучающих фреймворков, таких как Keras model.fit. Если вы не используете model.fit, вам нужно использовать API strategy.run для явного распределения вычислений. См. пример в учебном пособии по настройке цикла обучения https://www.tensorflow.org/tutorials/distribute/custom_training.

Что должно быть в области и что вне её?

Существует ряд требований к тому, что должно происходить внутри области. Однако в тех местах, где у нас есть информация о используемой стратегии, мы часто входим в область для пользователя, чтобы он не должен был делать это явно (т.е. вызов внутри или вне области допустим).

  • Все, что создаёт переменные, которые должны быть распределёнными, должно вызываться в 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 сохраняет информацию о области видимости. Когда затем вызываются методы высокоуровневого фреймворка обучения, такие как 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 иногда может потребоваться внутри области видимости, если оно создаёт переменные.
Возвращаемое значение
Менеджер контекста.

© 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/versions/r2.9/api_docs/python/tf/distribute/experimental/ParameterServerStrategy

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API