Spec-Zone.ru › TensorFlow 2.9

tf.data.experimental.service.distribute

Преобразование, которое перемещает обработку набора данных в службу tf.data.

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

Псевдонимы совместимости для миграции

См. Руководство по миграции для получения дополнительной информации.

tf.compat.v1.data.experimental.service.distribute

tf.data.experimental.service.distribute(
    processing_mode,
    service,
    job_name=None,
    consumer_index=None,
    num_consumers=None,
    max_outstanding_requests=None,
    data_transfer_protocol=None,
    compression='AUTO',
    target_workers='AUTO'
)

При итерации по набору данных, содержащему преобразование distribute, служба tf.data создает «задачу», которая генерирует данные для итерации по набору данных.

Служба tf.data использует кластер рабочих узлов для подготовки данных для обучения вашей модели. Аргумент processing_mode для tf.data.experimental.service.distribute описывает, как использовать несколько рабочих узлов для обработки входного набора данных. В настоящее время доступны два режима обработки: «distributed_epoch» и «parallel_epochs».

«distributed_epoch» означает, что набор данных будет разделен между всеми рабочими узлами службы tf.data. Диспечер создает «разделы» для набора данных и отправляет их рабочим узлам для дальнейшей обработки. Например, если набор данных начинается со списка имен файлов, диспечер будет перебирать имена файлов и отправлять имена файлов рабочим узлам tf.data, которые выполнят остальные преобразования набора данных над этими файлами. «distributed_epoch» полезен, когда вашей модели необходимо увидеть каждый элемент набора данных ровно один раз или если ей необходимо увидеть данные в обобщённо последовательном порядке. «distributed_epoch» работает только для наборов данных со сплитовыми источниками, такими как Dataset.from_tensor_slices, Dataset.list_files или Dataset.range.

«parallel_epochs» означает, что весь входной набор данных будет независимо обрабатываться каждым из рабочих узлов службы tf.data. По этой причине важно перемешивать данные (например, имена файлов) непредсказуемым образом, чтобы каждый рабочий узел обрабатывал элементы набора данных в другом порядке. «parallel_epochs» можно использовать для распределения наборов данных, которые не являются сплитовыми.

При использовании двух рабочих узлов «parallel_epochs» каждый элемент набора данных будет создан дважды:

dispatcher = tf.data.experimental.service.DispatchServer()
dispatcher_address = dispatcher.target.split("://")[1]
# Start two workers
workers = [
    tf.data.experimental.service.WorkerServer(
        tf.data.experimental.service.WorkerConfig(
            dispatcher_address=dispatcher_address)) for _ in range(2)
]
dataset = tf.data.Dataset.range(10)
dataset = dataset.apply(tf.data.experimental.service.distribute(
    processing_mode="parallel_epochs", service=dispatcher.target))
print(sorted(list(dataset.as_numpy_iterator())))
[0, 0, 1, 1, 2, 2, 3, 3, 4, 4, 5, 5, 6, 6, 7, 7, 8, 8, 9, 9]

«distributed_epoch», с другой стороны, по-прежнему будет генерировать каждый элемент один раз:

dispatcher = tf.data.experimental.service.DispatchServer()
dispatcher_address = dispatcher.target.split("://")[1]
workers = [
    tf.data.experimental.service.WorkerServer(
        tf.data.experimental.service.WorkerConfig(
            dispatcher_address=dispatcher_address)) for _ in range(2)
]
dataset = tf.data.Dataset.range(10)
dataset = dataset.apply(tf.data.experimental.service.distribute(
    processing_mode="distributed_epoch", service=dispatcher.target))
print(sorted(list(dataset.as_numpy_iterator())))
[0, 1, 2, 3, 4, 5, 6, 7, 8, 9]

При использовании apply(tf.data.experimental.service.distribute(...)), набор данных перед преобразованием apply выполняется в службе tf.data, в то время как операции после apply выполняются в локальном процессе.

dispatcher = tf.data.experimental.service.DispatchServer()
dispatcher_address = dispatcher.target.split("://")[1]
workers = [
    tf.data.experimental.service.WorkerServer(
        tf.data.experimental.service.WorkerConfig(
            dispatcher_address=dispatcher_address)) for _ in range(2)
]
dataset = tf.data.Dataset.range(5)
dataset = dataset.map(lambda x: x*x)
dataset = dataset.apply(
   tf.data.experimental.service.distribute("parallel_epochs",
                                           dispatcher.target))
dataset = dataset.map(lambda x: x+1)
print(sorted(list(dataset.as_numpy_iterator())))
[1, 1, 2, 2, 5, 5, 10, 10, 17, 17]

В приведённом выше примере операции над набором данных (перед применением функции distribute к элементам) будут выполняться на рабочих узлах tf.data, а элементы будут предоставляться через RPC. Остальные преобразования (после вызова distribute) будут выполняться локально. Диспечер и рабочие узлы будут связываться с неиспользуемыми свободными портами (которые выбираются случайным образом), чтобы взаимодействовать друг с другом. Однако для привязки их к определённым портам можно передать параметр port.

Аргумент job_name позволяет совместно использовать задачи между несколькими наборами данных. Вместо того, чтобы каждый набор данных создавал свою задачу, все наборы данных с одинаковым значением job_name будут потреблять одну и ту же задачу. Новая задача будет создаваться для каждой итерации набора данных (при каждом повторении Dataset.repeat считается новая итерация). Предположим, что DispatchServer обслуживает localhost:5000 и два рабочих узла обучения (в настройке с одним клиентом или несколькими клиентами) итерируются по набору данных ниже, и есть один рабочий узел tf.data:

range5_dataset = tf.data.Dataset.range(5)
dataset = range5_dataset.apply(tf.data.experimental.service.distribute(
    "parallel_epochs", "localhost:5000", job_name="my_job_name"))
for iteration in range(3):
  print(list(dataset))

Элементы каждой задачи будут разделены между двумя процессами, причём элементы будут потребляться процессами по принципу «первый пришёл — первый обслужен». Возможный результат состоит в том, что процесс 1 выведет

[0, 2, 4]
[0, 1, 3]
[1]

а процесс 2 выведет

[1, 3]
[2, 4]
[0, 2, 3, 4]

Имена задач не должны повторно использоваться для разных задач обучения в течение жизненного цикла службы tf.data. В общем случае ожидается, что служба tf.data будет существовать в течение всего времени выполнения одной задачи обучения. Чтобы использовать службу tf.data с несколькими задачами обучения, убедитесь, что вы используете разные имена задач, чтобы избежать конфликтов. Например, предположим, что задача обучения вызывает distribute со значением job_name="job" и читает до конца входных данных. Если другая независимая задача подключается к той же службе tf.data и пытается прочитать из job_name="job", она сразу же получит конец входных данных, не получив никаких данных.

Координированный считывание данных

По умолчанию, когда несколько потребителей читают из одной и той же задачи, они получают данные по принципу «первый пришёл — первый обслужен». В некоторых случаях целесообразно координировать потребителей. На каждом шаге потребители читают данные с одного и того же рабочего узла.

Например, служба tf.data может использоваться для координации размеров примеров в кластере во время синхронного обучения, таким образом, чтобы на каждом шаге все реплики обучались на элементах с похожими размерами. Для этого определите набор данных, который генерирует раунды num_consumers последовательных подобных по размеру пакетов, а затем включите координированное чтение, установив consumer_index и num_consumers.

Примечание: Чтобы сохранить потребителей в синхронизации, потребление данных по круговому алгоритму требует, чтобы набор данных имел бесконечную мощность. Это можно получить, добавив .repeat() в конце определения набора данных.

Keras и стратегии распределения

Набор данных, созданный преобразованием distribute , можно передать в Model.fit Keras или в tf.distribute.Strategy.experimental_distribute_dataset стратегии распределения, как и любой другой tf.data.Dataset. Рекомендуется установить job_name при вызове distribute , чтобы, если есть несколько рабочих узлов, они получали данные из одной и той же задачи. Обратите внимание, что автоматическое фрагментирование, обычно выполняемое experimental_distribute_dataset , будет отключено при установке job_name, так как совместное использование задачи уже приводит к разделению данных между рабочими узлами. При использовании общей задачи данные будут динамически сбалансированы между рабочими узлами, так что они примерно в одно и то же время достигнут конца входных данных. Это приводит к лучшему использованию рабочих узлов, чем при автоматическом фрагментировании, где каждый рабочий узел обрабатывает независимый набор файлов, и некоторые рабочие узлы могут раньше других закончить обработку данных.

Аргументы
processing_mode tf.data.experimental.service.ShardingPolicy , определяющий, как разбить набор данных между рабочими узлами tf.data. Подробности см. в tf.data.experimental.service.ShardingPolicy. Для обеспечения обратной совместимости processing_mode также может быть установлен на строки "parallel_epochs" или "distributed_epoch", которые соответственно эквивалентны ShardingPolicy.OFF и ShardingPolicy.DYNAMIC.
service Строка или кортеж, указывающие, как подключиться к службе tf.data. Если это строка, она должна иметь формат [<protocol>://]<address>, где <address> идентифицирует адрес диспечера, а <protocol> можно использовать для переопределения используемого протокола по умолчанию. Если это кортеж, он должен иметь вид (протокол, адрес).
job_name (Необязательно.) Название задачи. Если указано, оно должно быть непустой строкой. Этот аргумент позволяет нескольким наборам данных совместно использовать одну и ту же задачу. По умолчанию набор данных создает анонимные задачи, которые принадлежат только ему.
consumer_index (Необязательно.) Индекс потребителя в диапазоне от 0 до num_consumers. Должен быть указан вместе с num_consumers. При указании потребители будут читать из задачи в строгом циклическом порядке, вместо стандартного порядка «первый пришёл — первый обслужен».
num_consumers (Необязательно.) Количество потребителей, которые будут потреблять из задачи. Должен быть указан вместе с consumer_index. При указании потребители будут читать из задачи в строгом циклическом порядке, вместо стандартного порядка «первый пришёл — первый обслужен». Когда num_consumers указано, набор данных должен иметь бесконечную мощность, чтобы предотвратить преждевременное завершение поставщика и привести к несовпадению потребителей.
max_outstanding_requests (Необязательно.) Ограничение на количество элементов, которые могут быть запрошены одновременно. Этот параметр можно использовать для управления объёмом используемой памяти, так как distribute не будет использовать больше, чем element_size * max_outstanding_requests памяти.
data_transfer_protocol (Необязательно.) Протокол для передачи данных со службой tf.data. По умолчанию данные передаются с помощью gRPC.
compression Способ сжатия элементов набора данных перед передачей их по сети. "AUTO" оставляет решение о том, как сжимать, за службой tf.data runtime. None указывает на то, что сжатие не требуется.
target_workers (Необязательно.) Рабочие узлы для чтения. Если "AUTO", служба tf.data решает, с каких рабочих узлов читать. Если "ANY", читает с любых рабочих узлов службы tf.data. Если "LOCAL", читает только с локальных рабочих узлов службы tf.data в процессе. "AUTO" хорошо подходит для большинства случаев, но пользователи могут указать другие цели. Например, "LOCAL" помогает избежать RPC и копирования данных, если каждый рабочий узел TF находится вместе с рабочим узлом службы tf.data. Потребители общей задачи должны использовать те же target_workers. По умолчанию "AUTO".
Возвращаемое значение
Dataset Итератор элементов, генерируемых службой данных.

© 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/data/experimental/service/distribute

Spec-Zone.ru

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