tf.data.experimental.service.distribute
Преобразование, перемещающее обработку набора данных в службу tf.data.
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',
cross_trainer_cache=None,
target_workers='AUTO'
) -> Callable[tf.data.Dataset, tf.data.Dataset]
При итерации по набору данных, содержащему преобразование 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 набор данных может передаваться в Keras' Model.fit или стратегию распределения 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 указывает, что не следует сжимать. |
cross_trainer_cache | (Необязательно.) Если предоставлен объект CrossTrainerCache, итерация по набору данных будет совместно использоваться между одновременно работающими трейнерами. Подробности см. в https://www.tensorflow.org/api_docs/python/tf/data/experimental/service#sharing_tfdata_service_with_concurrent_trainers. |
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/api_docs/python/tf/data/experimental/service/distribute