Spec-Zone.ru › TensorFlow 2.9

Модуль: tf.data.experimental.service

API для использования сервиса tf.data.

Этот модуль содержит:

  1. Реализации серверов tf.data для запуска сервиса tf.data.
  2. 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 service, которая итерирует по набору данных для генерации данных. Чтобы предоставить данные задачи для нескольких клиентов (например, при использовании 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 service относится к процессу чтения из набора данных, управляемого сервисом 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) могут свободно перезапускаться, добавляться или удаляться во время обучения. При запуске рабочие (workers) зарегистрируются у диспетчера и начнут обработку всех незавершенных заданий с самого начала.

Использование с 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 DispatchServer: Сервер диспетчеризации tf.data в процессе.

class DispatcherConfig: Класс конфигурации для диспетчеров tf.data service.

class ShardingPolicy: Определяет, как разделить данные между рабочими узлами tf.data service.

class WorkerConfig: Класс конфигурации для диспетчеров tf.data service.

class WorkerServer: Сервер рабочего узла tf.data в процессе.

Функции

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

from_dataset_id(...): Создаёт набор данных, который считывает данные из службы tf.data.

register_dataset(...): Регистрирует набор данных в службе tf.data.

© 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

Spec-Zone.ru

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