Модуль: tf.compat.v1.data.experimental.service
API для использования сервиса tf.data.
Этот модуль содержит:
- Реализации серверов tf.data для запуска сервиса tf.data.
- API для регистрации наборов данных в сервисе tf.data и чтения из зарегистрированных наборов данных.
Сервис tf.data предоставляет следующие преимущества:
- Горизонтальное масштабирование обработки входных конвейеров tf.data для решения проблем с узкими местами ввода.
- Координация данных для распределенного обучения. Координированное чтение позволяет всем репликам обучаться на примерах одинаковой длины на каждом глобальном шаге обучения, что улучшает время шагов при синхронном обучении.
- Динамическое балансирование данных между репликами обучения.
dispatcher = tf.data.experimental.service.DispatchServer()
dispatcher_address = dispatcher.target.split("://")[1]
worker = tf.data.experimental.service.WorkerServer(
tf.data.experimental.service.WorkerConfig(
dispatcher_address=dispatcher_address))
dataset = tf.data.Dataset.range(10)
dataset = dataset.apply(tf.data.experimental.service.distribute(
processing_mode=tf.data.experimental.service.ShardingPolicy.OFF,
service=dispatcher.target))
print(list(dataset.as_numpy_iterator()))
[0, 1, 2, 3, 4, 5, 6, 7, 8, 9]
Настройка
В этом разделе описано, как настроить сервис tf.data.
Запуск серверов tf.data
Сервис tf.data состоит из одного сервера диспетчера и n серверов рабочих процессов. Сервера tf.data должны быть запущены вместе с вашими задачами обучения, а затем остановлены после завершения задач. Используйте tf.data.experimental.service.DispatchServer для запуска сервера диспетчера и tf.data.experimental.service.WorkerServer для запуска серверов рабочих процессов. Для целей тестирования серверы можно запустить в одном процессе или масштабировать на отдельных машинах.
См. https://github.com/tensorflow/ecosystem/tree/master/data_service для примера использования Google Kubernetes Engine (GKE) для управления сервисом tf.data. Обратите внимание, что реализация сервера в tf_std_data_server.py не специфична для GKE и может быть использована для запуска сервиса tf.data в других контекстах.
Пользовательские операции
Если ваш набор данных использует пользовательские операции, эти операции необходимо сделать доступными для серверов tf.data, вызвав load_op_library из процессов диспетчера и рабочих процессов при запуске.
Использование
Пользователи взаимодействуют с сервисом tf.data, программно регистрируя свои наборы данных в сервисе tf.data, а затем создавая наборы данных, которые считывают из зарегистрированных наборов данных. Функция register_dataset регистрирует набор данных, а затем функция from_dataset_id создает новый набор данных, который считывает из зарегистрированного набора данных. Функция distribute оборачивает register_dataset и from_dataset_id в одну удобную трансформацию, которая регистрирует входной набор данных, а затем считывает из него. distribute позволяет использовать сервис tf.data с одной строкой кода. Однако он предполагает, что набор данных создается и потребляется одним и тем же субъектом, и это предположение не всегда может быть верным или желательным. В частности, в определенных сценариях, таких как распределенное обучение, может быть желательно отделить создание и потребление набора данных (через register_dataset и from_dataset_id соответственно), чтобы избежать необходимости создания набора данных на каждом из рабочих процессов обучения.
Пример
distribute
Для использования трансформации distribute, примените её после префикса вашего конвейера ввода, который вы хотите выполнить с помощью сервиса tf.data (обычно в конце).
dataset = ... # Define your dataset here.
# Move dataset processing from the local machine to the tf.data service
dataset = dataset.apply(
tf.data.experimental.service.distribute(
processing_mode=tf.data.experimental.service.ShardingPolicy.OFF,
service=FLAGS.tf_data_service_address,
job_name="shared_job"))
# Any transformations added after `distribute` will be run on the local machine.
dataset = dataset.prefetch(1)
Приведённый выше код создаст задачу сервиса tf.data, которая итерируется по набору данных для генерации данных. Чтобы поделиться данными из задачи между несколькими клиентами (например, при использовании TPUStrategy или MultiWorkerMirroredStrategy), установите общий job_name для всех клиентов.
register_dataset и from_dataset_id
register_dataset регистрирует набор данных в сервисе tf.data, возвращая идентификатор набора данных для зарегистрированного набора данных. from_dataset_id создаёт набор данных, который считывает из зарегистрированного набора данных. Эти API могут быть использованы для сокращения времени построения набора данных при распределённом обучении. Вместо того, чтобы строить набор данных на всех рабочих процессах обучения, мы можем построить его один раз и затем зарегистрировать набор данных с помощью register_dataset. Затем все рабочие процессы могут вызвать from_dataset_id без необходимости самостоятельного построения набора данных.
dataset = ... # Define your dataset here.
dataset_id = tf.data.experimental.service.register_dataset(
service=FLAGS.tf_data_service_address,
dataset=dataset)
# Use `from_dataset_id` to create per-worker datasets.
per_worker_datasets = {}
for worker in workers:
per_worker_datasets[worker] = tf.data.experimental.service.from_dataset_id(
processing_mode=tf.data.experimental.service.ShardingPolicy.OFF,
service=FLAGS.tf_data_service_address,
dataset_id=dataset_id,
job_name="shared_job")
Режимы обработки
processing_mode определяет, как фрагментировать набор данных среди рабочих процессов сервиса tf.data. Сервис tf.data поддерживает OFF, DYNAMIC, FILE, DATA, FILE_OR_DATA, HINT политики фрагментации.
OFF: Фрагментация не будет выполнена. Весь входной набор данных будет обрабатываться независимо каждым из рабочих процессов сервиса tf.data. По этой причине важно перемешивать данные (например, имена файлов) не детерминированным образом, чтобы каждый рабочий процесс обрабатывал элементы набора данных в разном порядке. Этот режим может быть использован для распределения наборов данных, которые не могут быть разделены.
Если рабочий процесс добавляется или перезапускается во время обработки в режиме ShardingPolicy.OFF, он создаст новую копию набора данных и начнёт генерировать данные с начала.
Динамическая фрагментация
DYNAMIC: В этом режиме сервис tf.data делит набор данных на две составляющие: компонент источника, который генерирует «разделы», такие как имена файлов, и компонент обработки, который принимает разделы и выводит элементы набора данных. Компонент источника выполняется централизованно диспетчером сервиса tf.data, который генерирует различные разделы входных данных. Компонент обработки выполняется параллельно рабочими процессами сервиса tf.data, каждый из которых работает с различными наборами разделов входных данных.
Например, рассмотрим следующий набор данных:
dataset = tf.data.Dataset.from_tensor_slices(filenames)
dataset = dataset.interleave(TFRecordDataset)
dataset = dataset.map(preprocess_fn)
dataset = dataset.batch(batch_size)
dataset = dataset.apply(
tf.data.experimental.service.distribute(
processing_mode=tf.data.experimental.service.ShardingPolicy.DYNAMIC,
...))
from_tensor_slices будет выполняться на диспетчере, а interleave, map, и batch будут выполняться на рабочих процессах сервиса tf.data. Рабочие процессы будут получать имена файлов от диспетчера для обработки. Для обработки набора данных с динамической фрагментацией набор данных должен иметь разделимый источник, и все его преобразования должны быть совместимы с разделением. Хотя большинство источников и преобразований поддерживают разделение, есть исключения, такие как пользовательские наборы данных, которые могут не реализовывать API разделения. Пожалуйста, создайте вопрос на Github, если вы хотите использовать обработку распределённых эпох для текущего недопустимого источника или преобразования набора данных.
Если во время обучения рабочие процессы не перезапускаются, режим динамической фрагментации посетит каждый пример ровно один раз. Если рабочие процессы перезапускаются во время обучения, разделы, которые они обрабатывали, не будут полностью посещены. Диспетчер поддерживает курсор через разделы набора данных. Предполагая, что включена отказоустойчивость (см. «Отказоустойчивость» ниже), диспетчер будет хранить состояние курсора в журналах предварительной записи, чтобы курсор можно было восстановить в случае перезапуска диспетчера в середине обучения. Это обеспечивает гарантию посещения элементов не более одного раза в присутствии перезапусков сервера.
Статическая фрагментация
Ниже приведены политики статической фрагментации. Семантика аналогична tf.data.experimental.AutoShardPolicy. Эти политики требуют:
- Кластер сервиса tf.data настроен с фиксированным списком рабочих процессов в DispatcherConfig.
- Каждый клиент читает только с локального рабочего процесса сервиса tf.data.
Если рабочий процесс перезапускается во время выполнения статической фрагментации, он начнёт обрабатывать свой фрагмент снова с начала.
FILE: Фрагментация по входным файлам (т.е. каждый рабочий процесс получит фиксированный набор файлов для обработки). При выборе этого варианта убедитесь, что файлов не меньше, чем рабочих процессов. Если файлов меньше, чем рабочих процессов, произойдёт ошибка во время выполнения.
DATA: Фрагментация по элементам, производимым набором данных. Каждый рабочий процесс обработает весь набор данных и отбросит часть, которая не предназначена для него. Обратите внимание, что для правильного разбиения элементов набора данных набор данных должен генерировать элементы в детерминированном порядке.
FILE_OR_DATA: Попытка фрагментации по файлам, переходя к фрагментации по данным при ошибке.
HINT: Ищет наличие shard(SHARD_HINT, ...), которое обрабатывается как заполнитель, который нужно заменить на shard(num_workers, worker_index).
Для обратной совместимости processing_mode также может быть установлен на строки "parallel_epochs" или "distributed_epoch", которые соответственно эквивалентны ShardingPolicy.OFF и ShardingPolicy.DYNAMIC.
Координированное чтение данных
По умолчанию, когда несколько потребителей читают из одной задачи, они получают данные в порядке их поступления. В некоторых случаях полезно координировать потребителей. На каждом шаге потребители считывают данные с одного и того же рабочего процесса.
Например, сервис tf.data можно использовать для координации размеров примеров в кластере во время синхронного обучения, чтобы на каждом шаге все реплики обучались на примерах с одинаковым размером. Для этого определите набор данных, который генерирует раунды num_consumers последовательных партий одинакового размера, а затем включите координированное чтение, установив consumer_index и num_consumers.
Примечание: Для синхронизации потребителей координированное чтение требует, чтобы набор данных имел бесконечную мощность. Вы можете получить это, добавив .repeat() в конце определения набора данных.
Задачи
Задача сервиса tf.data относится к процессу чтения из набора данных, управляемого сервисом tf.data, с использованием одного или нескольких потребителей данных. Задачи создаются при итерации по наборам данных, которые считывают данные из сервиса tf.data. Данные, производимые задачей, определяются (1) набором данных, связанным с задачей, и (2) режимом обработки задачи. Например, если задача создана для набора данных Dataset.range(5), и режим обработки — ShardingPolicy.OFF, каждый рабочий процесс tf.data будет производить элементы {0, 1, 2, 3, 4} для задачи, что приведет к тому, что задача произведёт 5 * num_workers элементов. Если режим обработки — ShardingPolicy.DYNAMIC, задача будет производить только 5 элементов.
Один или несколько потребителей могут потреблять данные из задания. По умолчанию задания являются «анонимными», что означает, что только потребитель, создавший задание, может читать из него. Чтобы поделиться результатами задания между несколькими потребителями, вы можете установить общий job_name.
Устойчивость к отказам
По умолчанию сервер диспетчеризации tf.data хранит своё состояние в оперативной памяти, что делает его единственной точкой отказа во время обучения. Чтобы избежать этого, передайте fault_tolerant_mode=True при создании вашего DispatchServer. Устойчивость к отказам диспетчера требует, чтобы work_dir была настроена и доступна для диспетчера как до, так и после перезапуска (например, путь к GCS). При включенном режиме устойчивости к отказам диспетчер будет регистрировать своё состояние в каталоге работы, чтобы не потерять состояние при перезапуске диспетчера.
WorkerServers могут быть свободно перезапущены, добавлены или удалены во время обучения. При запуске рабочие узлы будут регистрироваться у диспетчера и начнут обработку всех ожидающих заданий с самого начала.
Использование с tf.distribute
tf.distribute — это API TensorFlow для распределённого обучения. Существует несколько способов использования tf.data с tf.distribute: strategy.experimental_distribute_dataset, strategy.distribute_datasets_from_function, и (для PSStrategy) coordinator.create_per_worker_dataset. В следующих разделах приведены примеры кода для каждого случая.
В общем случае мы рекомендуем использовать tf.data.experimental.service.{register_dataset,from_dataset_id} вместо tf.data.experimental.service.distribute по двум причинам:
- Набор данных необходимо создать и оптимизировать только один раз, а не на каждом рабочем узле. Это может значительно сократить время запуска, потому что текущие реализации
experimental_distribute_datasetиdistribute_datasets_from_functionсоздают и оптимизируют наборы данных рабочих узлов последовательно. - Если набор данных зависит от таблиц поиска или переменных, которые присутствуют только на одном узле, набор данных должен быть зарегистрирован с этого узла. Как правило, это происходит только тогда, когда ресурсы размещаются на главном узле или рабочем узле 0. Регистрация набора данных с главного узла позволит избежать проблем, связанных с зависимостью от удалённых ресурсов.
strategy.experimental_distribute_dataset
При использовании strategy.experimental_distribute_dataset ничего особенного не требуется, просто примените register_dataset и from_dataset_id как и выше, убедившись, что указано job_name, чтобы все рабочие узлы потребляли данные из одного задания службы tf.data.
dataset = ... # Define your dataset here.
dataset_id = tf.data.experimental.service.register_dataset(
service=FLAGS.tf_data_service_address,
dataset=dataset)
dataset = tf.data.experimental.service.from_dataset_id(
processing_mode=tf.data.experimental.service.ShardingPolicy.OFF,
service=FLAGS.tf_data_service_address,
dataset_id=dataset_id,
job_name="shared_job")
dataset = strategy.experimental_distribute_dataset(dataset)
strategy.distribute_datasets_from_function
Сначала убедитесь, что набор данных, созданный функцией dataset_fn, не зависит от input_context для рабочего узла обучения, на котором он выполняется. Вместо того чтобы каждый рабочий узел создавал свой собственный (разделенный) набор данных, один рабочий узел должен зарегистрировать неразделенный набор данных, а остальные рабочие узлы должны получать данные из этого набора данных.
dataset = dataset_fn()
dataset_id = tf.data.experimental.service.register_dataset(
service=FLAGS.tf_data_service_address,
dataset=dataset)
def new_dataset_fn(input_context):
del input_context
return tf.data.experimental.service.from_dataset_id(
processing_mode=tf.data.experimental.service.ShardingPolicy.OFF,
service=FLAGS.tf_data_service_address,
dataset_id=dataset_id,
job_name="shared_job")
dataset = strategy.distribute_datasets_from_function(new_dataset_fn)
coordinator.create_per_worker_dataset
create_per_worker_dataset работает так же, как и distribute_datasets_from_function.
dataset = dataset_fn()
dataset_id = tf.data.experimental.service.register_dataset(
service=FLAGS.tf_data_service_address,
dataset=dataset)
def new_dataset_fn(input_context):
del input_context
return tf.data.experimental.service.from_dataset_id(
processing_mode=tf.data.experimental.service.ShardingPolicy.OFF,
service=FLAGS.tf_data_service_address,
dataset_id=dataset_id,
job_name="shared_job")
dataset = coordinator.create_per_worker_dataset(new_dataset_fn)
Ограничения
- Обработка данных на основе Python: Наборы данных, использующие обработку данных на основе Python (например,
tf.py_function,tf.numpy_functionилиtf.data.Dataset.from_generator) в настоящее время не поддерживаются. - Несериализуемые ресурсы: Наборы данных могут зависеть только от ресурсов TF, которые поддерживают сериализацию. Сериализация в настоящее время поддерживается для таблиц поиска и переменных. Если ваш набор данных зависит от ресурса TF, который нельзя сериализовать, пожалуйста, создайте вопрос на Github.
- Удаленные ресурсы: Если набор данных зависит от ресурса, набор данных должен быть зарегистрирован из того же процесса, который создал ресурс (например, «главный» узел ParameterServerStrategy).
Классы
class DispatcherConfig: Класс конфигурации для диспетчеров tf.data service.
class ShardingPolicy: Указывает, как разделить данные между рабочими узлами tf.data service.
class WorkerConfig: Класс конфигурации для диспетчеров tf.data service.
Функции
distribute(...): Преобразование, которое перемещает обработку наборов данных в tf.data service.
from_dataset_id(...): Создаёт набор данных, который считывает данные из tf.data service.
register_dataset(...): Регистрирует набор данных в tf.data service.
© 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/compat/v1/data/experimental/service