Spec-Zone.ru › TensorFlow 2.4

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, max_outstanding_requests=None
)

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

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

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

Аргументы
processing_mode Строка, определяющая политику обработки данных рабочими процессами tf.data. Может быть «parallel_epochs», чтобы каждый рабочий процесс tf.data обрабатывал копию набора данных, или «distributed_epoch», чтобы разбить одну итерацию набора данных между всеми рабочими процессами.
service Строка, указывающая, как подключиться к службе tf.data. Строка должна иметь формат «протокол://адрес», например, «grpc://localhost:5000».
job_name (Необязательно.) Имя задания. Этот аргумент позволяет нескольким наборам данных совместно использовать одно задание. По умолчанию набор данных создаёт анонимные задания, принадлежащие только ему.
max_outstanding_requests (Необязательно.) Предел запросов элементов одновременно. Этот параметр позволяет управлять объёмом используемой памяти, так как distribute не будет использовать более element_size * max_outstanding_requests памяти.
Возвращаемое значение
Dataset Итератор элементов, производимых службой данных.

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

Spec-Zone.ru

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