tf.data.experimental.service.distribute
Преобразование, перемещающее обработку набора данных в службу tf.data.
tf.data.experimental.service.distribute(
processing_mode, service, job_name=None, max_outstanding_requests=None
)
При итерации по набору данных, содержащему преобразование distribute, служба tf.data создает «задачу», которая генерирует данные для итерации по набору данных.
Аргумент processing_mode определяет, какие данные генерирует задача службы tf.data. В настоящее время поддерживается только режим "parallel_epochs".
processing_mode="parallel_epochs" означает, что несколько рабочих процессов tf.data будут итерировать по набору данных параллельно, каждый из которых генерирует все элементы набора данных. Например, если набор данных содержит {0, 1, 2}, каждый рабочий процесс tf.data, используемый для выполнения, будет генерировать {0, 1, 2}. Если используется 3 рабочих процесса, задача сгенерирует элементы {0, 0, 0, 1, 1, 1, 2, 2, 2} (хотя, возможно, не в таком порядке). Для учета этого рекомендуется случайным образом перемешать ваш набор данных, чтобы разные рабочие процессы tf.data итерировались по нему в разных порядках.
В будущем появятся дополнительные режимы обработки. Например, режим "one_epoch", который распределяет набор данных по рабочим процессам tf.data, так что потребители видят каждый элемент набора данных только один раз.
dataset = tf.data.Dataset.range(5)
dataset = dataset.map(lambda x: x*x)
dataset = dataset.apply(
tf.data.experimental.service.distribute("parallel_epochs",
"grpc://dataservice:5000"))
dataset = dataset.map(lambda x: x+1)
for element in dataset:
print(element) # prints { 1, 2, 5, 10, 17 }
В примере первые две строки (перед вызовом distribute) будут выполнены в рабочих процессах tf.data, а предоставленные элементы будут передаваться по RPC. Остальные преобразования (после вызова distribute) будут выполнены локально.
Аргумент job_name позволяет обмениваться задачами между несколькими наборами данных. Вместо того, чтобы каждый набор данных создавал свою задачу, все наборы данных с одинаковым значением job_name будут получать данные из одной задачи. Новая задача будет создана для каждой итерации набора данных (при каждом повторении Dataset.repeat будет считаться новой итерацией). Предположим, что два рабочих процесса обучения (в одном клиенте или в настройке с несколькими клиентами) итерируются по набору данных ниже, и есть один рабочий процесс tf.data:
range5_dataset = tf.data.Dataset.range(5)
dataset = range5_dataset.apply(tf.data.experimental.service.distribute(
"parallel_epochs", "grpc://dataservice: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". |
service | Строка, указывающая, как подключиться к службе tf.data. Строка должна иметь формат <протокол>://<адрес>, например grpc://localhost:5000. |
job_name | (Необязательно.) Имя задачи. Этот аргумент позволяет нескольким наборам данных совместно использовать одну задачу. По умолчанию набор данных создает анонимные задачи, которые принадлежат только ему. |
max_outstanding_requests | (Необязательно.) Ограничение на количество элементов, которые могут быть запрошены одновременно. Вы можете использовать этот параметр для управления объемом используемой памяти, так как distribute не будет использовать более element_size * max_outstanding_requests памяти. |
| Возвращаемое значение | |
|---|---|
Dataset | Объект 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.3/api_docs/python/tf/data/experimental/service/distribute