Мы хотели бы поблагодарить следующих людей за их вклад, проведение бенчмарков и рецензирование этой работы: Эдвард Оукс, Менджин Ян, Картица Моди, Дхие Шах, Джош Ли, Зак Полицер, Стив Александр, Эндрю Си Ким (Google), Мао Яньцань (ByteDance).
Поскольку Ray используется во все большем количестве новых задач ИИ, мы заметили, что пользователи доводят Ray Core до предела с двух сторон:
- Пакетный вывод (batch inference), перемешивание (shuffle) и другие конвейеры данных: рабочие нагрузки, где драйвер управляет жизненным циклом миллионов объектов, выполняя цикл планирования на критическом пути, что ограничивает их устойчивую пропускную способность задач.
Пакетный вывод, перемешивание и другие конвейеры данных: рабочие нагрузки, где драйвер управляет жизненным циклом миллионов объектов, выполняя цикл планирования на критическом пути, что ограничивает их устойчивую пропускную способность задач.
- Крупномасштабное обучение с подкреплением (RL) и предварительное/последующее обучение: они выполняются на кластерах от 2 000 до 10 000 узлов, где десятки тысяч акторов планируются с учетом ограничений топологии сети. Время запуска и время перезапуска после сбоя критически важны для этих пользователей.
Крупномасштабное обучение с подкреплением (RL) и предварительное/последующее обучение: они выполняются на кластерах от 2 000 до 10 000 узлов, где десятки тысяч акторов планируются с учетом ограничений топологии сети. Время запуска и время перезапуска после сбоя критически важны для этих пользователей.
Обнаруженные нами узкие места сводились к трем типам:
- Состязание за блокировки (lock contention) между потоками.
Состязание за блокировки между потоками.
- Один перегруженный поток, обрабатывающий асинхронную работу.
Один перегруженный поток, обрабатывающий асинхронную работу.
- Планирование на основе устаревших данных о ресурсах.
Планирование на основе устаревших данных о ресурсах.
После ряда улучшений Ray Core, вот как выглядят эти рабочие нагрузки по сравнению с прошлым годом.
- На 23% быстрее при сквозном выполнении пакетного вывода на 500 узлах.
На 23% быстрее при сквозном выполнении пакетного вывода на 500 узлах.
- На 24% быстрее при сквозном выполнении перемешивания данных в Ray Data.
На 24% быстрее при сквозном выполнении перемешивания данных в Ray Data.
- Теперь Ray может запускать больше акторов на более крупных тренировочных кластерах с учетом топологических ограничений: в 62 раза быстрее подготовка групп размещения на 2 000 узлах; в 303 раза быстрее подготовка групп размещения на 10 000 узлах; в 6,5 раза быстрее сквозной запуск акторов на 2 000 узлах. Масштабируется до 40 000 акторов на кластере из 10 000 узлов — масштаб, который невозможно было поддерживать год назад.
Теперь Ray может запускать больше акторов на более крупных тренировочных кластерах с учетом топологических ограничений:
- В 62 раза быстрее подготовка групп размещения на 2 000 узлах.
В 62 раза быстрее подготовка групп размещения на 2 000 узлах.
- В 303 раза быстрее подготовка групп размещения на 10 000 узлах.
В 303 раза быстрее подготовка групп размещения на 10 000 узлах.
- В 6,5 раза быстрее сквозной запуск акторов на 2 000 узлах.
В 6,5 раза быстрее сквозной запуск акторов на 2 000 узлах.
- Масштабируется до 40 000 акторов на кластере из 10 000 узлов — масштаб, который невозможно было поддерживать год назад.
Масштабируется до 40 000 акторов на кластере из 10 000 узлов — масштаб, который невозможно было поддерживать год назад.
В остальной части этой статьи мы расскажем, как мы обнаружили каждое узкое место и что с этим сделали. Мы начнем с драйвера и конвейеров данных, затем перейдем к тренировочным кластерам на 10 000 узлов, при этом три вышеупомянутых узких места постоянно возникают в обоих случаях.
Когда драйвер становится узким местом
Все началось с загадки в Ray Data. Конвейеры пакетного вывода и перемешивания теряли время при вызовах Ray Core, а их цикл планирования замедлялся вместе с ними, что снижало устойчивую пропускную способность задач и общее время выполнения. Однако, когда мы измеряли эти вызовы отдельно, ray.wait и другие API Ray Core работали быстро. Так куда же уходило время?
Чтобы ответить на это, давайте проследим жизненный цикл одного объекта. Миллионы объектов проходят тот же путь.
Рисунок 1: Жизненный цикл одного объекта
Как показано на рисунке 1, процесс драйвера Ray имеет два потока, важных для этой истории. Поток Python выполняет ваш код, например, ray.wait и логику планирования. Поток C, фоновый поток Ray Core, обрабатывает всю асинхронную работу.
Рассматривая часть потока Python на рисунке 1: вы вызываете f.remote() для отправки задачи, получаете ссылку на объект (ref), ждете готовности объекта (реальный цикл планирования Ray Data использует тайм-аут и продолжает проверку), а затем передаете ссылку другим последующим задачам. Наконец, когда объект был использован и больше никому не нужен, вы удаляете ссылку.
Теперь за кулисами. Когда вы вызываете ray.wait, API Ray Core, он переводит ваш поток в спящий режим в ожидании переменной условия готовности объекта. Задача выполняется на узле-воркере A и сохраняет там свой выходной объект. Идеальное время ожидания должно быть в точности равно времени выполнения вашей задачи. Поток C получает уведомление после завершения задачи и, выполнив учет, включая запись местоположения объекта, пробуждает поток Python.
Теперь вы передаете ссылку последующей задаче g. Она планируется на узле-воркере B, которому нужно знать, где находится объект, поэтому он подписывается на местоположение объекта. Поток C записал это местоположение ранее, а теперь публикует его обратно. Узел B забирает объект с узла A, и задача g выполняется.
А удаление ссылки за кулисами приводит счетчик ссылок объекта к нулю, что запускает очистку, включая публикацию сообщения об ошибке для всех, кто все еще ждет местоположения этого объекта, чтобы они прекратили ожидание и отказались от попыток его получения, а также освобождение каждой копии в кластере. С учетом вышеизложенного фонового контекста, давайте рассмотрим узкие места, которые мы обнаружили в драйвере, и способы их решения.
Два потока, одна блокировка
Рисунок 2: Два потока, борющиеся за одну и ту же блокировку публикации
Как мы видим из жизненного цикла объекта выше, оба потока публикуют обновления объектов, и существует блокировка публикации, защищающая все эти операции. Как показано на рисунке 2, при удалении последней ссылки потоку Python требовалось захватить блокировку. В то же время потоку C нужна та же блокировка для публикации местоположений объектов.
Два потока, одна блокировка — это приводит к серьезному состязанию при масштабировании. Из-за характера пакетного вывода объекты запрашиваются каждую секунду, поэтому поток C постоянно захватывает блокировку публикации для публикации местоположений объектов. В то же время объекты постоянно удаляются после использования, и потоку Python также требовалось захватить блокировку публикации для публикации сообщений об ошибках. В задаче пакетного вывода на 500 узлах только одна блокировка публикации съедала 17,4% времени цикла планирования.
Решение для такого рода состязания за блокировки — перенос всей связанной работы в один поток. Мы перенесли все операции публикации, включая сообщения об ошибках, в поток C, поэтому потоку Python больше никогда не требовалось захватывать блокировку публикации, и мы увидели, что состязание за блокировки исчезло из профиля.
Один перегруженный поток
Рисунок 3: Дополнительное время блокировки, когда поток C перегружен
И здесь возникает второе «узкое место», скрывающееся в середине жизненного цикла объекта, в самом ray.wait. Как показано на Рисунке 3, время, которое вы проводите в ray.wait, определяется не только скоростью создания объекта, но и загруженностью потока C. Синяя часть вашего ожидания связана с выполнением задачи. Красная часть — это дополнительное время блокировки. Когда ваш объект готов, поток C может быть занят обработкой обновлений местоположения множества других объектов и всем остальным, и всё это может стоять в очереди операций потока C перед пробуждением. Эта задержка в очереди оказалась одним из главных факторов снижения производительности. Поэтому мы начали поиски и обнаружили две проблемы.
Первой проблемой, которую мы обнаружили, был поток крошечных обратных вызовов (callbacks) в потоке C. Каждый обратный вызов сам по себе «дешев», но если вы отправляете их миллионы, сама отправка вредит производительности. В одном прогоне shuffle количество таких обратных вызовов превысило миллион и заняло около 17% времени потока C, причем каждый из них выполнялся 0,04 мс после ожидания своей очереди в течение 5,7 секунд.
Большинство этих обратных вызовов оказались проверками аргументов. Перед отправкой задачи каждый аргумент должен быть подтвержден как созданный где-то в кластере — по одному обратному вызову на каждый аргумент, а задачи reduce в shuffle принимают огромное количество аргументов. Но reduce запускается только после завершения всех задач map, поэтому все эти аргументы уже должны были быть созданы.
Проблема заключалась в том, что драйвер отправлял обратный вызов, даже если аргумент уже был готов в момент создания задачи, чтобы избежать потенциальной взаимной блокировки (deadlock). Исправление гарантирует безопасность всех блокировок и отправляет асинхронную проверку только в том случае, если аргумент еще не готов в момент создания задачи. Одно только это ускорило shuffle на 9,2% в сквозном режиме.
Вторая проблема сразу проявилась при профилировании. Две трети времени потока C уходило просто на публикацию обновлений местоположения объектов для других узлов. Исследование показало, что каждая задача на каждом узле подписывалась на местоположение своих аргументов. Однако, если объекты уже находятся на локальном узле, задаче не обязательно подписываться на их местоположение. А поскольку Ray планирует задачи рядом с их данными, вероятность того, что объекты находятся на том же узле, довольно высока. Мы внесли улучшения, пропуская запрос для любых объектов, которые уже являются локальными, и это сократило трафик публикации драйвера вдвое.
Вместе с множеством других исправлений, не упомянутых здесь, эти изменения увеличили устойчивую пропускную способность задач до 1,6 раза. В результате конвейер пакетного вывода (batch inference) теперь работает на 23% быстрее на 500 узлах, а shuffle — на 24% быстрее.
В долгосрочной перспективе мы работаем над разделением потока C на несколько потоков, чтобы еще больше ускорить Ray Core.
Масштабирование для крупных обучающих прогонов
Помимо вышеупомянутых рабочих нагрузок Ray Data, все больше пользователей сообщают нам, что время запуска и перезапуска акторов (actor) ограничивает их крупномасштабные рабочие нагрузки RL и предварительного/последующего обучения. Простой GPU — это дорогостоящая трата, и они платят за это многократно: один раз при запуске и снова при каждом сбое, поскольку коммуникатор NCCL работает по принципу «все или ничего», и один неисправный узел означает перезапуск каждого актора.
Их также беспокоят ограничения топологии сети, поскольку скорость передачи данных между этими акторами играет большую роль в производительности обучающих прогонов, и предпочтительный способ их выражения — планирование групп размещения (placement group) с учетом топологии в Ray. Например, при использовании стоек NVIDIA GB200 или GB300 NVL72, типичная настройка может представлять собой одну группу размещения на стойку из 18 узлов, один пакет (bundle) на узел, 4 актора внутри каждого пакета, каждый из которых соответствует одному GPU в стойке, где 18 умножить на 4 — это ровно 72 в NVL72.
Таким образом, чтобы поддерживать рабочие нагрузки RL и предварительного/последующего обучения в их текущем и будущем масштабе, Ray сначала должен масштабироваться до уровня 10 000 узлов, а затем сделать планирование групп размещения и планирование акторов достаточно быстрыми.
Поддержание стабильности GCS на 10 000 узлах
GCS, глобальная служба управления Ray, находится в центре плоскости управления, выполняя критически важную работу, которая включает:
- Управление узлами
Управление узлами
- Управление жизненным циклом акторов
Управление жизненным циклом акторов
- Обслуживание внутреннего хранилища «ключ-значение»
Обслуживание внутреннего хранилища «ключ-значение»
- Сбор данных о спросе для автоскейлера
Сбор данных о спросе для автоскейлера
- Планирование групп размещения
Планирование групп размещения
Все это раньше делило один основной поток, и большая часть этой работы растет вместе с количеством узлов. Это снова второе «узкое место» — один перегруженный поток. И основной поток GCS был занят задолго до того, как поступала какая-либо пользовательская рабочая нагрузка. Поэтому, когда приходила рабочая нагрузка актора или группы размещения, GCS вскоре становился «узким местом» и легко перегружался.
Мы сделали GCS многопоточным, чтобы он мог масштабироваться вместе с кластером и справляться с нагрузкой акторов. Широковещательная рассылка обновлений состояния узлов, обслуживание внутреннего хранилища «ключ-значение» и сбор ожидающего спроса для автоскейлера были перенесены в отдельные потоки.
Мы также облегчили то, что остается в GCS. Подписки на изменения узлов перешли от O(workers) к O(nodes), информация о пакетах групп размещения теперь кэшируется, а не запрашивается из GCS во время отправки актора, и мы включили возможность вывода событий задач полностью за пределы GCS, направляя их прямо на панель мониторинга.
Вот в каком состоянии сейчас находится нагрузка в режиме ожидания. В прошлогодней версии Ray (2.51) кластер из 10 000 узлов без отправленных заданий уже загружал основной поток GCS примерно на 61%; в сегодняшней ночной сборке тот же простаивающий кластер загружен примерно на 38%.
Результатом является GCS, который остается стабильным на 10 000 узлах. Но это не означает автоматически, что планирование становится быстрее, а это уже другая проблема.
Распределенное планирование Ray
Планирование в Ray распределено. Каждый узел способен планировать, а это означает, что каждый узел должен знать ресурсы всего кластера. Как показано на Рисунке 4, каждый узел поддерживает свое собственное локальное представление ресурсов, и это локальное представление является источником истины. При любом изменении узел делает снимок своего представления ресурсов и отправляет его в GCS. GCS собирает всю информацию и рассылает каждое изменение на каждый узел, поэтому каждый узел в конечном итоге владеет представлениями ресурсов всех остальных. Внутри Ray эта служба доставки называется syncer.
Рисунок 4: Как перемещаются представления ресурсов
Теперь давайте пройдемся по простой отправке задачи, чтобы увидеть, как планирует Ray. Отправитель сначала отправляет запрос на узел. В принципе, это может быть любой узел. На практике это обычно тот узел, на котором находится отправитель, или узел, который содержит большую часть данных, необходимых задаче. Допустим, этот узел — узел A. Узел A сначала проверяет, может ли он разместить задачу локально. Если нет, он выбирает узел на основе своего представления ресурсов кластера, в данном случае узел B. Затем узел A оптимистично вычитает ресурсы задачи из своего представления и возвращает решение отправителю. Отправитель затем отправляет запрос на узел B, и тот же процесс продолжается. На этот раз B принимает запрос, выделяет ресурсы по-настоящему, и задача выполняется на узле B.
Рисунок 5: Как Ray планирует задачу
Обратите внимание, на чем основывается этот дизайн. Узел A принимает решение на основе своего представления о ресурсах, и это представление может быть устаревшим, поэтому решение может оказаться неверным. В таком случае отправитель направляет запрос на узел B, узел B отклоняет его, и мы перепланируем задачу. Мы полагаемся на синхронизатор (syncer) в доставке актуальных представлений о ресурсах и на то, что он исправит неверное оптимистичное предположение узла A. Именно в этой зависимости кроется третье «узкое место». Когда синхронизатор отстает, все вышеперечисленное планируется на основе устаревшего представления.
Группы размещения: от часов до секунд
Планирование групп размещения в масштабе было медленным от начала до конца. Одна из причин была заложена прямо в дизайне, описанном выше. Вспомните об оптимистичном предположении и синхронизаторе, который исправляет его в случае ошибки. Существует одна утечка, которую синхронизатор не может покрыть, и на рисунке 6 она показана от начала до конца. С того момента, как узел A возвращает запрос отправителю, узел A больше не отслеживает его, но он уже оптимистично вычел ресурсы узла B в своем собственном представлении. Единственный способ исправить неверное вычитание — это обновление синхронизатора. Поэтому, если запрос отбрасывается отправителем или отправитель завершает работу, то на узле B ничего не меняется, нет ничего, что могло бы активировать синхронизатор, и вычитание узла A остается «утечкой» навсегда.
Рисунок 6: Когда запрос никогда не доходит, возникает утечка вычитания
Чтобы решить эту проблему в Ray, каждые три секунды каждый узел сбрасывает свое локально измененное представление до последнего полученного снимка состояния, и GCS, который размещает группы размещения на основе своего собственного представления ресурсов, делает то же самое. Обычно это безвредно, потому что синхронизатор работает намного быстрее, чем три секунды, и свежие снимки состояния продолжают поступать. Однако при больших масштабах синхронизатор работает медленно, и сама страховка стала проблемой. Группы размещения обычно создаются сразу при запуске кластера, поэтому снимок состояния, о котором сообщал каждый узел, гласил, что при инициализации кластер был пуст. Каждые три секунды сброс заставлял занятые ресурсы снова выглядеть свободными, и планировщик на самом деле не планировал, а просто отправлял группы размещения на уже заполненные узлы, пока синхронизатор не догонял процесс.
Но действительно ли нам нужен сброс для групп размещения? Оказывается, планирование групп размещения отличается от остального планирования в Ray. Как показано на рисунке 7, группы размещения проходят через двухфазную фиксацию и могут планироваться только GCS. GCS сначала выполняет подготовку, пытаясь зарезервировать все ресурсы на узлах, и фиксирует изменения только после того, как все резервирования прошли успешно. В отличие от запроса актора или задачи, запрос группы размещения никогда не пересылается через кого-либо еще, и GCS всегда знает, был ли он успешным, благодаря двухфазной фиксации. Поэтому GCS вообще не нуждался в этой страховке. Простое удаление кода сброса для GCS сократило время планирования групп размещения на 10 000 узлах до секунд.
Рисунок 7: Группы размещения фиксируются в две фазы
Другой частью медленной работы были накладные расходы, возникающие при проверке готовности группы размещения. pg.ready() — это API для пользователя, позволяющее проверить, готова ли группа размещения. Это звучит как простая проверка. «Под капотом» оно раньше отправляло фиктивную задачу в группу размещения. Планирование задачи в группе означает, что группа должна быть готова, и поскольку задача естественным образом возвращает ссылку, pg.ready() становится асинхронным API «бесплатно». Но у этого ярлыка были реальные издержки. Фиктивная задача наследовала среду выполнения вашего задания, поэтому простая проверка ресурсов могла ожидать завершения установки torch. И поскольку API возвращает одну ссылку на группу размещения, это естественным образом провоцирует последовательный цикл. Мы обнаружили, что многие пользователи пишут это так:
Но каждый ray.get блокировался на своей собственной фиктивной задаче, поэтому проверки выполнялись одна за другой, отправляя следующую фиктивную задачу только после завершения предыдущей. Тысяча групп означала тысячу последовательных обращений. Теперь pg.ready() — это прямой асинхронный запрос к GCS. Он по-прежнему возвращает тот же ObjectRef, но за ним больше нет задачи, и накладные расходы исчезли.
Вот чего удалось достичь: то же создание групп размещения, от 200 до 10 000 узлов, одна группа на стойку из 18 узлов, один бандл на узел, на прошлогодней версии Ray и на сегодняшней ночной сборке.
Рисунок 8: Время готовности группы размещения в зависимости от размера кластера
Планирование акторов: минуты позади истины
Несмотря на то, что время работы групп размещения сократилось до секунд, запуск акторов внутри групп размещения с учетом топологии все еще был медленным на уровне 10 000 узлов. Мы развиваем два пути параллельно, по одному для каждой из двух проблем ниже.
Первая проблема заключается в том, что сам синхронизатор работает медленно. Ray предполагает, что каждый узел может принимать решения о планировании, поэтому, как показано на рисунке 9, изменения стекаются в GCS с каждого узла, и каждое из них рассылается на все 10 000 узлов. Это делает синхронизатор медленным, поэтому узлы часто имеют устаревшие представления о ресурсах и принимают неверные решения.
Рисунок 9: Синхронизатор не справляется с нагрузкой на 10 000 узлах
Устаревшее представление вредит акторам двумя способами, особенно при использовании групп размещения. Планирование акторов обычно выполняется через головной узел. После того как группы размещения с учетом топологии были размещены, головному узлу нужна информация о ресурсах группы размещения, прежде чем он сможет планировать акторов в них. Ресурсы группы размещения изначально живут только на узлах, которые их содержат, и полагаются на синхронизатор для распространения на головной узел. Синхронизатор работает медленно. На уровне 10 000 узлов мы видели, что для распространения в пиковые моменты требуется около 200 секунд. Поэтому головной узел вообще не может планировать, пока не поступит информация о группе размещения.
И как только она поступает, головной узел все еще планирует на основе устаревшего представления, поскольку синхронизатор часто не может доставить данные быстро, а трехсекундный сброс делает их еще более устаревшими. И то, и другое подталкивает головной узел к неверным решениям, отправляя акторов на узлы, которые уже заполнены. Это приводит к отклонению и перепланированию, и ситуация становится еще хуже по мере заполнения кластера, с все большим количеством перепланирований в конце. Мы наблюдали этот «длинный хвост» напрямую.
Чтобы решить эту проблему, мы в настоящее время работаем над более централизованным механизмом планирования. Большинству рабочих нагрузок Ray нужно лишь несколько узлов для принятия решений о планировании, поэтому мы отказываемся от предположения, с которого начался весь этот дизайн: что каждый узел может планировать. Как только решения будут сосредоточены, лишь немногим узлам вообще потребуется представление ресурсов всего кластера, поэтому синхронизатору придется передавать меньше данных, представление будет более точным, и, следовательно, большинство отклонений, вызванных устаревшими представлениями сегодня, исчезнут.
Вторая проблема заключается в том, что у актора сегодня «два мозга». GCS решает, когда он создается, перезапускается и уничтожается, в то время как владелец, обычно драйвер, хранит ссылки и знает, кто все еще его использует. Каждый запуск актора требует затрат на взаимодействие между ними, и на рисунке 10 показан этот процесс: как минимум четыре сообщения между владельцем и GCS, от регистрации актора и постоянной подписки на счетчик ссылок актора до получения его адреса обратно. Перезапуск также затрагивается этим процессом. Помните, что для рабочих нагрузок обучения один «мертвый» актор означает, что все акторы умирают и перезапускаются. Десятки тысяч постоянных подписок одновременно отправляются в GCS, а перезапуск повторяет весь процесс с самого начала, по одному разу для каждого актора, и снова все одновременно.
Рисунок 10: Взаимодействие между драйвером и GCS
Чтобы решить вышеуказанную проблему, в настоящее время мы работаем над переносом жизненного цикла актора из GCS к владельцу. Причина, по которой управление жизненным циклом актора изначально находится в GCS, заключается в отсоединенных (detached) акторах. Их жизненный цикл не зависит от драйвера, поэтому они должны продолжать существовать после завершения работы драйвера. Для неотсоединенных акторов, как только управление жизненным циклом переходит к владельцу, вся задержка, связанная с обменом сообщениями, исчезает.
Приведенные ниже результаты достигнуты благодаря вышеупомянутой работе над GCS и группами размещения, а также ранним наработкам по двум указанным выше направлениям, включая вклад сообщества, который позволяет синхронизатору отправлять сообщения синхронизации пакетами по мере возможности. Вот к чему это приводит: создание групп размещения с учетом топологии и последующее планирование акторов в них — сравнение Ray прошлого года и текущей ночной сборки.
Рисунок 11: Время запуска актора от начала до конца в зависимости от размера кластера
LinkК чему мы пришли
Каждое исправление, описанное в этой статье, уже присутствует в ночной сборке Ray, и мы продолжаем инвестировать в масштабируемость акторов и драйверов, включая два упомянутых выше улучшения. Если вы используете Ray на пределе возможностей, мы будем рады узнать, где возникают трудности! Вы можете найти нас на GitHub или в Ray Slack.









