Масштабирование Kubernetes до 7500 узлов

Нам удалось масштабировать кластеры Kubernetes до 7 500 узлов, создав масштабируемую инфраструктуру для таких крупных моделей, как GPT‑3, CLIP и DALL·E, а также для быстрых мелкомасштабных итеративных исследований, таких как законы масштабирования для нейронных языковых моделей.
Масштабирование отдельного кластера Kubernetes до такого размера практикуется редко и требует особой осторожности, но преимуществом является простая инфраструктура, которая позволяет нашим командам исследователей в области машинного обучения двигаться быстрее и масштабироваться без изменения своего кода.
С момента нашей последней публикации о масштабировании до 2500 узлов мы продолжили развивать нашу инфраструктуру для удовлетворения потребностей исследователей, извлекая по ходу дела множество дополнительных уроков. В этой статье суммируются эти уроки, чтобы другие участники сообщества Kubernetes могли извлечь из них пользу, а в конце описаны проблемы, с которыми мы все еще сталкиваемся и которые будем решать в дальнейшем.
Наша рабочая нагрузка
Прежде чем двигаться дальше, важно описать нашу рабочую нагрузку. Приложения и аппаратное обеспечение, которые мы запускаем с помощью Kubernetes, сильно отличаются от того, с чем вы можете столкнуться в типичной компании. Наши проблемы и соответствующие решения могут как подходить, так и не подходить для вашей конфигурации!
Крупная задача машинного обучения охватывает множество узлов и выполняется наиболее эффективно, когда у нее есть доступ ко всем аппаратным ресурсам на каждом узле. Это позволяет графическим процессорам (GPU) осуществлять прямую межсетевую связь с использованием NVLink или графическим процессорам напрямую общаться с сетевым адаптером (NIC) с помощью GPUDirect. Поэтому для многих наших рабочих нагрузок один под занимает весь узел целиком. Конденсация ресурсов NUMA, CPU или PCIE не является фактором для планирования. Упаковка или фрагментация не являются распространенной проблемой. Наши текущие кластеры имеют полную пропускную способность деления пополам, поэтому мы также не принимаем во внимание топологию стоек или сети. Все это означает, что, несмотря на большое количество узлов, нагрузка на планировщик относительно невелика.
Тем не менее, нагрузка на kube-scheduler носит пиковый характер. Новая задача может состоять из сотен подов, создаваемых одновременно, а затем возвращаться к относительно низкой скорости изменений.
Наши самые крупные задачи выполняются с использованием MPI, и все поды в рамках задачи участвуют в едином коммуникаторе MPI. Если какой-либо из участвующих подов погибает, вся задача останавливается и требует перезапуска. Задача регулярно сохраняет контрольные точки, а при перезапуске возобновляет работу с последней контрольной точки. Таким образом, мы считаем поды полусостояние зависимыми (semi-stateful): завершенные поды могут быть заменены, и работа может продолжаться, но это вызывает сбои и должно быть сведено к минимуму.
Мы не слишком полагаемся на балансировку нагрузки Kubernetes. У нас очень мало HTTPS-трафика, нет необходимости в A/B-тестировании, схемах blue/green или canary-развертываниях. Поды общаются друг с другом напрямую по их IP-адресам с помощью MPI через SSH, а не через конечные точки сервисов. «Обнаружение» сервисов ограничено; мы просто выполняем однократный поиск подов, участвующих в MPI, во время запуска задачи.
Большинство задач взаимодействует с той или иной формой блочного хранилища (blob storage). Обычно они либо передают потоком сегменты набора данных, либо создают контрольные точки непосредственно из блочного хранилища, либо кэшируют их на быстрый локальный временный диск. У нас есть несколько постоянных томов (PersistentVolumes) для случаев, когда полезны семантики POSIX, но блочное хранилище гораздо более масштабируемо и не требует медленных операций отсоединения/подсоединения.
Наконец, характер нашей работы носит фундаментально исследовательский характер, что означает, что сами рабочие нагрузки постоянно меняются. Хотя команда суперкомпьютерных вычислений стремится обеспечить уровень качества вычислительной инфраструктуры, который мы считаем «производственным», приложения, работающие в этом кластере, недолговечны, а их разработчики быстро внедряют изменения. В любой момент могут возникнуть новые шаблоны использования, которые ставят под сомнение наши предположения о тенденциях и правильных компромиссах. Нам нужна устойчивая система, которая также позволяет нам быстро реагировать при изменении условий.
Сеть
По мере увеличения количества узлов и подов в наших кластерах мы столкнулись с тем, что Flannel испытывает трудности с масштабированием требуемой пропускной способности. Мы перешли на использование встроенных сетевых технологий подов для конфигураций IP в масштабируемых наборах виртуальных машин Azure (Azure VMSS) и соответствующих плагинов CNI. Это позволило нам добиться сетевой пропускной способности уровня хоста на наших подах.
Другая причина, по которой мы перешли на IP-адресацию на основе псевдонимов (alias-based), заключается в том, что в наших крупнейших кластерах одновременно может использоваться примерно 200 000 IP-адресов. Когда мы тестировали маршрутизируемую сеть подов, мы обнаружили серьезные ограничения на количество маршрутов, которые мы можем эффективно использовать.
Отказ от инкапсуляции повышает требования к базовой SDN или механизму маршрутизации, но делает нашу сетевую конфигурацию простой. Добавление VPN или туннелирования может быть выполнено без каких-либо дополнительных адаптеров. Нам не нужно беспокоиться о фрагментации пакетов из-за того, что какая-то часть сети имеет более низкий MTU. Сетевые политики и мониторинг трафика просты; нет никакой двусмысленности относительно источника и назначения пакетов.
Мы используем тегирование iptables на хосте для отслеживания использования сетевых ресурсов по пространствам имен (Namespaces) и подам. Это позволяет исследователям визуализировать шаблоны использования сети. В частности, поскольку многие из наших экспериментов имеют четкие закономерности интернет-трафика и связи между подами, часто бывает полезно выяснить, где могут возникать узкие места.
Правила iptables mangle можно использовать для произвольной маркировки пакетов, соответствующих определенным критериям. Вот наши правила для определения того, является ли трафик внутренним или предназначенным для интернета. Правила FORWARD охватывают трафик от подов по сравнению с INPUT и OUTPUT трафиком от хоста:
Plain Text
1iptables -t mangle -A INPUT ! -s 10.0.0.0/8 -m comment --comment "iptables-exporter openai traffic=internet-in"2iptables -t mangle -A FORWARD ! -s 10.0.0.0/8 -m comment --comment "iptables-exporter openai traffic=internet-in"3iptables -t mangle -A OUTPUT ! -d 10.0.0.0/8 -m comment --comment "iptables-exporter openai traffic=internet-out"4iptables -t mangle -A FORWARD ! -d 10.0.0.0/8 -m comment --comment "iptables-exporter openai traffic=internet-out"После маркировки iptables запускает счетчики для отслеживания количества байт и пакетов, соответствующих этому правилу. Вы можете визуально оценить эти счетчики с помощью самого iptables:
Plain Text
1% iptables -t mangle -L -v2Chain FORWARD (policy ACCEPT 50M packets, 334G bytes)3 pkts bytes target prot opt in out source destination4....51253K 555M all -- any any anywhere !10.0.0.0/8 /* iptables-exporter openai traffic=internet-out */61161K 7937M all -- any any !10.0.0.0/8 anywhere /* iptables-exporter openai traffic=internet-in */Мы используем экспортер Prometheus с открытым исходным кодом под названием iptables-exporter, чтобы передавать эти данные в нашу систему мониторинга. Это простой способ отслеживания пакетов, соответствующих различным типам условий.
Один несколько уникальный аспект нашей сетевой модели заключается в том, что мы полностью предоставляем исследователям диапазоны CIDR сети узлов, подов и сервисов. У нас сетевая модель «звезда» (hub and spoke), и мы используем собственные диапазоны CIDR узлов и подов для маршрутизации этого трафика. Исследователи подключаются к хабу и оттуда получают доступ к любому из отдельных кластеров (периферийным узлам). Но сами кластеры не могут общаться друг с другом. Это гарантирует, что кластеры остаются изолированными, без межкластерных зависимостей, которые могут нарушить изоляцию сбоев.
Мы используем узел «NAT» для трансляции диапазона CIDR сети сервисов для трафика, поступающего извне кластера. Такая настройка предоставляет нашим исследователям значительную гибкость в выборе того, как и какие типы сетевых конфигураций они могут использовать для своих экспериментов.
API-серверы
API-серверы Kubernetes и etcd являются важнейшими компонентами для поддержания кластера в работоспособном состоянии, поэтому мы уделяем особое внимание нагрузке на эти системы. Мы используем панели мониторинга Grafana, предоставляемые kube-prometheus, а также собственные панели мониторинга. Мы сочли полезным настроить оповещения о частоте HTTP-статусов 429 (Too Many Requests) и 5xx (Server Error) на API-серверах в качестве общего сигнала о проблемах.
Хотя некоторые запускают API-серверы внутри kube, мы всегда запускали их вне самого кластера. И etcd, и API-серверы работают на собственных выделенных узлах. В наших крупнейших кластерах работает 5 API-серверов и 5 узлов etcd для распределения нагрузки и минимизации последствий в случае выхода из строя одного из них. У нас не было никаких серьезных проблем с etcd с тех пор, как мы выделили события Kubernetes (Events) в их собственный кластер etcd еще в прошлом блог-посте. API-серверы не имеют состояния и, как правило, легко запускаются в группе самовосстанавливающихся инстансов или scaleset. Мы еще не пытались создавать какую-либо автоматизацию самовосстановления кластеров etcd, поскольку инциденты происходили крайне редко.
API-серверы могут потреблять немало памяти, и этот показатель имеет тенденцию масштабироваться линейно с количеством узлов в кластере. Для нашего кластера с 7 500 узлами мы наблюдаем использование до 70 ГБ кучи (heap) на один API-сервер, так что, к счастью, в будущем это продолжит укладываться в аппаратные возможности.
Серьезную нагрузку на API-серверы создавали операции WATCH для Endpoints. Существует несколько сервисов, таких как kubelet и node-exporter, участником которых является каждый узел в кластере. Когда узел добавлялся в кластер или удалялся из него, срабатывал этот WATCH. И поскольку обычно каждый узел сам отслеживал сервис kubelet через kube-proxy, количество и пропускная способность, необходимые для этих ответов, составляли N2 N^2 и были огромными, порой достигая 1 ГБ/с и более. EndpointSlices, запущенные в Kubernetes 1.17, стали огромным преимуществом, снизившим эту нагрузку в 1000 раз.
В целом мы очень внимательно следим за любыми запросами к API-серверу, которые масштабируются в зависимости от размера кластера. Мы стараемся избегать взаимодействия каких-либо DaemonSet с API-сервером. В тех случаях, когда вам действительно нужно, чтобы каждый узел отслеживал изменения, внедрение промежуточного сервиса кэширования, такого как Datadog Cluster Agent, представляется хорошим шаблоном для предотвращения узких мест в масштабах всего кластера.
По мере роста наших кластеров мы реже прибегаем к фактическому автомасштабированию самих кластеров. Но время от времени мы сталкивались с проблемами при слишком масштабном автомасштабировании за один раз. При добавлении нового узла в кластер генерируется множество запросов, и добавление сотен узлов одновременно может перегрузить емкость API-сервера. Сглаживание этого процесса даже на несколько секунд помогло избежать сбоев.
Метрики временных рядов с помощью Prometheus и Grafana
Мы используем Prometheus для сбора метрик временных рядов и Grafana для построения графиков, панелей мониторинга и оповещений. Мы начинали с развертывания kube-prometheus, которое собирает широкий спектр метрик и предоставляет хорошие панели мониторинга для визуализации. Со временем мы добавили множество собственных панелей мониторинга, метрик и оповещений.
По мере добавления все новых и новых узлов мы сталкивались с трудностями из-за огромного объема метрик, собираемых Prometheus. Хотя kube-prometheus предоставляет много полезных данных, некоторые из них мы на самом деле никогда не просматривали, а некоторые были слишком детализированными для эффективного сбора, хранения и выполнения запросов. Мы используем правила Prometheus, чтобы «отбрасывать» (drop) некоторые из этих метрик при сборе.
Некоторое время мы боролись с проблемой, когда Prometheus потреблял все больше и больше памяти, пока в конце концов контейнер не завершался аварийно с ошибкой Out-Of-Memory (OOM). Это казалось происходящим даже после выделения огромных объемов памяти для приложения. Хуже того, после сбоя приложению требовалось много часов на запуск для воспроизведения файлов журнала упреждающей записи (WAL), прежде чем оно снова могло использоваться.
В конце концов мы выяснили источник этих OOM-ошибок: им оказалось взаимодействие между Grafana и Prometheus, при котором Grafana использовала API /api/v1/series в Prometheus с запросом {le!=""} (в сущности, «дай мне все метрики гистограмм»). Реализация /api/v1/series была ничем не ограничена ни по времени, ни по пространству — для запроса с большим количеством результатов это продолжало потреблять все больше памяти и времени. Этот процесс также продолжался даже после того, как запрашивающий сдался и закрыл соединение. Для нас памяти всегда было недостаточно, и Prometheus в итоге падал. Мы внесли патч в Prometheus, чтобы заключить этот API в Context для принудительного ограничения по времени ожидания (timeout), что полностью решило проблему.
Хотя Prometheus стал падать гораздо реже, в те моменты, когда нам все же приходилось его перезапускать, воспроизведение WAL оставалось проблемой. Часто требовалось много часов на воспроизведение всех логов WAL, прежде чем Prometheus возобновлял сбор новых метрик и обслуживание запросов. С помощью Robust Perception мы выяснили, что применение параметра GOMAXPROCS=24 приводит к значительному улучшению. Prometheus пытается использовать все ядра во время воспроизведения WAL, и для серверов с большим количеством ядер борьба за ресурсы сводит на нет всю производительность.
Мы изучаем новые возможности для увеличения нашей емкости мониторинга, которые описаны в разделе »Нерешенные проблемы» ниже.
Проверки работоспособности (Healthchecks)
При кластере такого размера мы, разумеется, полагаемся на автоматизацию для обнаружения и удаления некорректно работающих узлов. Со временем мы создали целый ряд систем проверки работоспособности.
Пассивные проверки работоспособности
Некоторые проверки работоспособности являются пассивными и постоянно выполняются на всех узлах. Они отслеживают базовые системные ресурсы, такие как доступность сети, поврежденные или заполненные диски, а также ошибки графических процессоров (GPU). GPU могут давать сбои самыми разными способами, но одна из самых распространенных и простых проблем — это «неисправимая ошибка ECC» (Uncorrectable ECC error). Инструменты Nvidia Data Center GPU Manager (DCGM) позволяют легко отправлять запросы для выявления этой и ряда других ошибок «Xid». Один из способов отслеживания этих ошибок — использование dcgm-exporter для передачи метрик в Prometheus — нашу систему мониторинга. Это будет отображаться в виде метрики DCGM_FI_DEV_XID_ERRORS и соответствовать коду ошибки, возникшей в последний раз. Кроме того, API запросов устройств NVML предоставляет более подробную информацию о состоянии и работе GPU.
Как только мы обнаруживаем ошибку, ее часто удается устранить путем перезагрузки GPU или системы, хотя в некоторых случаях это приводит к необходимости физической замены самого графического процессора.
Другой тип проверки работоспособности отслеживает события обслуживания от облачного провайдера. Каждый из основных облачных провайдеров предоставляет способ узнать, запланировано ли для текущей виртуальной машины (ВМ) предстоящее обслуживание, которое в конечном итоге приведет к сбою. Возможно, потребуется перезагрузить ВМ для применения патча гипервизора или заменить физический узел на другое оборудование.
Эти пассивные проверки работоспособности постоянно выполняются в фоновом режиме на всех узлах. Если проверка работоспособности начинает выдавать сбои, узел автоматически изолируется (cordoned), поэтому на нем больше не планируются новые поды. В случае более серьезных сбоев мы также пытаемся выполнить эвакуацию подов, чтобы все запущенные в данный момент поды немедленно завершили работу. Решение о том, разрешать ли такую эвакуацию, по-прежнему остается за самим подом и настраивается с помощью бюджета прерываний подов (Pod Disruption Budget). В конечном итоге, либо после завершения работы всех подов, либо по истечении 7 дней (согласно нашему SLA), мы принудительно завершаем работу ВМ.
Активные тесты GPU
К сожалению, не все проблемы с GPU проявляются в виде кодов ошибок, видимых через DCGM. Мы создали собственную библиотеку тестов, которые нагружают GPU для выявления дополнительных проблем и проверки того, что аппаратное обеспечение и драйвер работают должным образом. Эти тесты не могут выполняться в фоновом режиме — для их запуска требуется эксклюзивное использование GPU в течение нескольких секунд или минут.
Сначала мы запускаем эти тесты на узлах при загрузке в рамках системы, которую мы называем «предолетной» (preflight). Все узлы присоединяются к кластеру с примененным «предолетным» осквернением (taint) и меткой (label). Это осквернение предотвращает планирование обычных подов на узле. DaemonSet настроен на запуск тестовых подов preflight на всех узлах с этой меткой. После успешного завершения теста сам тест удаляет осквернение и метку, после чего узел становится доступным для общего использования.
Затем мы также периодически запускаем эти тесты в течение жизненного цикла узла. Мы запускаем их в виде CronJob, что позволяет им выполняться на любом доступном узле в кластере. Признаем, что выбор узлов для тестирования здесь носит несколько случайный и неконтролируемый характер, но мы обнаружили, что со временем это обеспечивает достаточный охват при минимальных координации и сбоях.
Квоты и использование ресурсов
По мере масштабирования наших кластеров исследователи начали сталкиваться с трудностями при получении всей выделенной им емкости. Традиционные системы планирования заданий имеют множество различных функций для справедливого распределения работы между конкурирующими командами, которых нет в Kubernetes. Со временем мы вдохновились этими системами планирования заданий и реализовали несколько возможностей в стиле Kubernetes.
Командные осквернения (Team taints)
В каждом кластере у нас есть служба «team-resource-manager», которая выполняет множество функций. Ее источником данных является ConfigMap, определяющий кортежи (селектор узла, применяемая метка команды, объем выделенных ресурсов) для всех исследовательских команд, имеющих квоту в данном кластере. Она сопоставляет эти данные с текущими узлами в кластере, помечая соответствующее количество узлов с помощью ###NON_TRAN_LABEL### openai.com/team=teamname:NoSchedule.
«team-resource-manager» также имеет службу веб-перехватчика допуска (admission webhook), благодаря которой при отправке каждого задания применяется соответствующее правило толерантности (toleration) на основе принадлежности автора к команде. Использование осквернений позволяет нам гибко ограничивать планировщик подов Kubernetes, например, разрешая толерантность «любая» («any») для подов с более низким приоритетом, что позволяет командам заимствовать ресурсы друг у друга без необходимости сложной координации.
CPU- и GPU-баллоны
Помимо использования cluster-autoscaler для динамического масштабирования наших кластеров на базе ВМ, мы применяем его для исправления (удаления и повторного добавления) неисправных участников кластера. Для этого мы устанавливаем «минимальный размер» кластера равным нулю, а «максимальный размер» — доступной емкости. Однако, обнаружив простаивающие узлы, cluster-autoscaler попытается уменьшить масштаб до уровня необходимой емкости. По ряду причин (задержка запуска ВМ, предварительно выделенные затраты, упомянутое выше влияние на API-сервер) такое масштабирование при простаивании не является идеальным.
Поэтому мы внедрили развертывание баллонов (balloon Deployment) как для хостов только с CPU, так и с GPU. Это развертывание содержит ReplicaSet с низкоприоритетными подами в количестве, равном «максимальному размеру». Эти поды занимают ресурсы на узле, поэтому автоскейлер не считает их простаивающими. Однако, поскольку они имеют низкий приоритет, планировщик может немедленно вытеснить их, чтобы освободить место для реальной работы. (Мы решили использовать Deployment вместо DaemonSet, чтобы избежать ситуации, когда DaemonSet считается простаивающей рабочей нагрузкой на узле.)
Стоит отметить, что мы используем антиаффинити подов (pod anti-affinity), чтобы гарантировать равномерное распределение подов по узлам. В более ранних версиях планировщика Kubernetes была проблема с производительностью O (N2) O (N^2) при использовании антиаффинити подов. Эта проблема была исправлена начиная с Kubernetes 1.18.
Групповое планирование (Gang scheduling)
Наши эксперименты часто включают один или несколько объектов StatefulSet, каждый из которых управляет определенной частью процесса обучения. Для оптимизаторов исследователям необходимо, чтобы все участники StatefulSet были запланированы до того, как можно будет начать какое-либо обучение (поскольку мы часто используем MPI для координации между участниками-оптимизаторами, а MPI чувствителен к изменениям состава группы).
Тем не менее, Kubernetes по умолчанию не обязательно отдает приоритет выполнению всех запросов от одного StatefulSet над другим. Например, если два эксперимента запрашивают по 100% емкости кластера, вместо планирования одного эксперимента целиком или другого Kubernetes может запланировать лишь половину подов каждого эксперимента, что приведет к тупиковой ситуации (дедлоку), при которой ни один эксперимент не сможет продвинуться вперед.
Мы пробовали несколько вариантов, требующих написания пользовательского планировщика, но столкнулись с граничными случаями, которые вызывали конфликты с тем, как планировались обычные поды. В Kubernetes 1.18 появилась архитектура плагинов для основного планировщика Kubernetes, что значительно упростило добавление подобных функций на уровне ядра. Недавно мы остановились на плагине совместного планирования (Coscheduling) как на хорошем способе решения этой проблемы.
Нерешенные проблемы
По мере масштабирования наших кластеров Kubernetes нам предстоит решить еще множество проблем. Вот некоторые из них:
Метрики
В масштабах нашего кластера мы столкнулись со множеством трудностей: встроенный механизм хранения TSDB в Prometheus медленно выполняет сжатие (compaction) и требует много времени на воспроизведение WAL (журнала упреждающей записи) при каждой перезагрузке. Запросы также часто приводят к ошибкам «query processing would load too many samples» (обработка запроса загрузила бы слишком много сэмплов). В настоящее время мы переходим на другое хранилище и движок запросов, совместимые с Prometheus. Ждите нашу будущую запись в блоге о том, как продвигается этот процесс!
Шейпинг сетевого трафика подов
По мере масштабирования наших кластеров для каждого пода рассчитывается определенный объем доступной интернет-пропускной способности. Совокупные требования к интернет-трафику на одного пользователя стали существенными, и теперь наши исследователи могут непреднамеренно создавать значительную ресурсную нагрузку на другие ресурсы в интернете, такие как репозитории датасетов для скачивания и устанавливаемые пакетами ПО.
Заключение
Мы обнаружили, что Kubernetes — исключительно гибкая платформа для наших исследовательских потребностей. Она обладает способностью масштабироваться в соответствии с самыми требовательными рабочими нагрузками, которые мы на нее возлагаем. Тем не менее, остается много областей, требующих улучшений, и команда суперкомпьютерных вычислений OpenAI продолжит исследовать возможности масштабирования Kubernetes. Если вам интересна подобная работа, вам стоит задуматься о том, чтобы подать заявку в OpenAI!
Авторы
Полный текст статьи читайте на OpenAI
