Spec-Zone.ru › TensorFlow 2.4

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 Keras compile/fit разрабатывается.

  • tf.distribute.experimental.ParameterServerStrategy должен использоваться с tf.distribute.experimental.coordinator.ClusterCoordinator.

Аргументы
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 будет вызываться на процессоре 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, чтобы увидеть пример.

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

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

  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. Распределяет значения в массиве на основе 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)
  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(), 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.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, 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, соответствующий его реплике.

Ключевой момент: В зависимости от реализации tf.distribute.Strategy и от того, включено ли выполнение с немедленным выполнением, 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. Его элемент может быть 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» вводится на каждом работнике.
Примечание: Вход в область не автоматически распределяет вычисление, за исключением случая высокоуровневой обучающей среды, такой как keras model.fit. Если вы не используете model.fit, вам нужно использовать API strategy.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

Spec-Zone.ru

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