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

Мы используемKubernetes для исследований в области глубокого обучения уже более двух лет. Хотя наши самые масштабные рабочие нагрузки напрямую управляют виртуальными машинами в облаке, Kubernetes обеспечивает быстрый цикл итераций, приемлемую масштабируемость и избавляет от лишнего шаблонного кода, что делает его идеальным выбором для большинства наших экспериментов. Сейчас мы эксплуатируем несколько кластеров Kubernetes (некоторые в облаке, а некоторые на физическом оборудовании), самый крупный из которых мы довели до более чем 2500 узлов. Этот кластер работает в Azure на комбинации виртуальных машин D15v2 и NC24.
На пути к такому масштабу множество системных компонентов вызывали сбои, включая etcd, мастер-узлы Kube, скачивание образов Docker, сеть, KubeDNS и даже ARP-кэши наших машин. Мы решили, что будет полезно поделиться конкретными проблемами, с которыми мы столкнулись, и тем, как мы их решили.
etcd
После того как наш кластер перешагнул отметку в 500 узлов, исследователи начали жаловаться на регулярные тайм-ауты в инструменте командной строки kubectl. Мы попытались добавить больше мастер-узлов Kube (виртуальных машин, запущенных под управлением kube-apiserver). Поначалу это казалось временным решением, но когда количество реплик превысило 10, мы поняли, что лечим симптомы, а не причину (для сравнения, GKE использует одну виртуальную машину с 32 ядрами на 500 узлов).
Это заставило нас серьезно заподозрить наш кластер etcd, который является центральным хранилищем состояний для мастер-узлов Kube. Заглянув в Datadog, мы увидели, что задержка записи подскочила до сотен миллисекунд на машинах DS15v2, где запускались наши реплики etcd, несмотря на то, что каждая машина использовала SSD-накопитель P30, способный выдавать 5000 IOPS.

Проведя бенчмаркинг производительности с помощью fio, мы обнаружили, что etcd может использовать лишь около 10% доступных IOPS, поскольку задержка записи составляла 2 мс, а etcd выполняет последовательный ввод-вывод, что делает его зависимым от задержек (latency-bound).
Затем мы перенесли каталог etcd для каждого узла на локальный временный диск — SSD, подключенный напрямую к инстансу, а не к сети. Переключение на локальный диск снизило задержку записи до 200 мкс, и состояние etcd нормализовалось!
Наш кластер работал хорошо до тех пор, пока мы не перешли отметку примерно в 1000 узлов, после чего мы снова столкнулись с высокой задержкой фиксации (commit latency) от etcd. На этот раз мы заметили, что kube-apiserver считывают из etcd более 500 МБ/с. Мы настроили Prometheus для мониторинга серверов API, а также установили флаги --audit-log-path и --audit-log-maxbackup, чтобы включить более подробное логирование на apiserver. Это выявило ряд медленных запросов и чрезмерных вызовов LIST API для событий (Events).
Первопричина: настройки по умолчанию для процессов мониторинга Fluentd и Datadog предполагали отправку запросов к apiserver с каждого узла в кластере (например, эта проблема теперь уже исправлена). Мы просто изменили конфигурацию этих процессов, сделав их опросы менее агрессивными, и нагрузка на apiserver снова стабилизировалась:

Еще одной полезной настройкой стало хранение Kubernetes Events в отдельном кластере etcd, чтобы всплески создания событий не влияли на производительность основных инстансов etcd. Для этого мы просто установили флаг --etcd-servers-overrides примерно следующим образом: --etcd-servers-overrides=/events#https://0.example.com:2381;https://1.example.com:2381;https://2.example.com:2381
Еще одна проблема после преодоления порога в 1000 узлов заключалась в достижении жесткого лимита хранения etcd (по умолчанию 2 ГБ), из-за чего он переставал принимать операции записи. Это спровоцировало каскадный сбой: все наши узлы Kube не прошли проверку работоспособности (health checks), и наш автомасштабировщик (autoscaler) решил, что необходимо завершить работу всех воркеров. Мы увеличили максимальный размер etcd с помощью флага --quota-backend-bytes, и теперь автомасштабировщик имеет защитную проверку, предотвращающую любые действия, если они приведут к завершению работы более 50% кластера.
Мастер-узлы Kube
Мы размещаем процессы kube-apiserver, kube-controller-manager и kube-scheduler на одних и тех же машинах. Для обеспечения высокой доступности (high availability) у нас всегда есть как минимум 2 мастер-узла, а флаг --apiserver-count установлен равным количеству запущенных apiserver (в противном случае мониторинг Prometheus может запутаться между инстансами).
Мы используем Kubernetes в основном как систему пакетного планирования (batch scheduling) и полагаемся на наш автомасштабировщик, чтобы динамически масштабировать кластер вверх и вниз — это позволяет нам существенно сократить расходы на простаивающие узлы, сохраняя при этом низкую задержку при быстрых итерациях. Политика kube-scheduler по умолчанию заключается в равномерном распределении нагрузки по узлам, но нам нужно противоположное поведение, чтобы неиспользуемые узлы можно было выключать, а крупные поды (pods) могли планироваться быстро. Поэтому мы переключились на следующую политику:
Plain Text
1{2"kind" : "Policy",3"apiVersion" : "v1",4"predicates" : [5 {"name" : "GeneralPredicates"},6 {"name" : "MatchInterPodAffinity"},7 {"name" : "NoDiskConflict"},8 {"name" : "NoVolumeZoneConflict"},9 {"name" : "PodToleratesNodeTaints"}10 ],11"priorities" : [12 {"name" : "MostRequestedPriority", "weight" : 1},13 {"name" : "InterPodAffinityPriority", "weight" : 2}14 ]15}Мы активно используем KubeDNS для обнаружения сервисов (service discovery), но вскоре после развертывания новой политики планирования в нем начались проблемы с надежностью. Мы выяснили, что сбои происходят только на определенных подах KubeDNS. Из-за новой политики планирования на некоторых машинах оказалось запущено более 10 копий KubeDNS, что создало «горячие точки» (hotspots), и мы превысили предел примерно в 200 QPS, разрешенный для каждой виртуальной машины Azure при запросах к внешним доменам.
Мы исправили это, добавив правило антиаффинити (anti-affinity rule) к нашим подам KubeDNS:
Plain Text
1affinity:2 podAntiAffinity:3 requiredDuringSchedulingIgnoredDuringExecution:4 - weight: 1005 labelSelector:6 matchExpressions:7 - key: k8s-app8 operator: In9 values:10 - kube-dns11 topologyKey: kubernetes.io/hostnameЗагрузка образов Docker
Наш проект Dota начинался на Kubernetes, и по мере его масштабирования мы заметили, что на новых узлах Kubernetes поды часто подолгу находятся в статусе Pending. Игровой образ весит около 17 ГБ, и его скачивание на свежий узел кластера часто занимало 30 минут, поэтому нам было понятно, почему контейнер Dota какое-то время остается в статусе Pending — однако это справедливо и для других контейнеров. Разобравшись в вопросе, мы обнаружили, что у kubelet есть флаг --serialize-image-pulls, который по умолчанию равен true, что означает: загрузка образа Dota блокировала все остальные образы. Переключение на false требовало перевода Docker на использование драйвера overlay2 вместо AUFS. Чтобы еще больше ускорить загрузку, мы также перенесли корневой каталог Docker на SSD, подключенный к инстансу, как сделали это для машин etcd.
Даже после оптимизации скорости скачивания мы сталкивались с тем, что поды не удавалось запустить из-за загадочного сообщения об ошибке: rpc error: code = 2 desc = net/http: request canceled. Логи kubelet и Docker также содержали сообщения о том, что загрузка образа была отменена из-за отсутствия прогресса. Мы отследили первопричину: слишком большие образы скачивались и распаковывались слишком долго, либо накапливалась большая очередь образов для скачивания. Чтобы решить эту проблему, мы установили флаг kubelet --image-pull-progress-deadline на 30 минут, а параметр демона Docker max-concurrent-downloads — на 10. (Второй параметр не ускорил распаковку больших образов, но позволил загружать очередь образов параллельно.)
Наша последняя проблема с загрузкой в Docker была связана с Google Container Registry. По умолчанию kubelet скачивает специальный образ из gcr.io (управляется флагом --pod-infra-container-image), который используется при запуске любого нового контейнера. Если эта загрузка по какой-либо причине не удается, например из-за исчерпания вашей квоты, этот узел не сможет запустить ни одного контейнера. Поскольку наши узлы выходят в интернет через NAT для доступа к gcr.io, а не имеют собственных публичных IP-адресов, мы весьма вероятной рискуем упереться в этот лимит квот на один IP-адрес. Чтобы исправить это, мы просто предварительно загружаем данный образ Docker в образ машины для наших воркеров Kubernetes, используя docker image save -o /opt/preloaded_docker_images.tar и docker image load -i /opt/preloaded_docker_images.tar. Для повышения производительности мы делаем то же самое для белого списка распространенных внутренних образов OpenAI, таких как образ Dota.
Сеть
По мере роста наших экспериментов они превращаются во все более сложные распределенные системы, которые сильно зависят от сети в своей работе. Когда мы только начали запускать распределенные эксперименты, сразу стало очевидно, что наша сеть настроена плохо. Напрямую между машинами мы получали пропускную способность 10–15 Гбит/с, но наши поды Kube с использованием Flannel выдавали максимум около 2 Гбит/с. Публичные бенчмарки Machine Zone показывают схожие цифры, а значит, дело было вряд ли просто в плохой конфигурации — скорее, это нечто присущее нашей среде. (Для сравнения, Flannel не добавляет такие накладные расходы на наших физических машинах.)
Чтобы обойти это, пользователи могут использовать две разные настройки для отключения Flannel для своего пода: hostNetwork: true и dnsPolicy: ClusterFirstWithHostNet. (Впрочем, перед этим обязательно ознакомьтесь с предупреждениями в документации Kubernetes.)
ARP-кэш
Несмотря на нашу настройку DNS, мы все равно сталкивались с периодическими проблемами разрешения имен. Однажды инженер сообщил, что nc -v к их серверу Redis выводит сообщение об установлении соединения более 30 секунд. Мы отследили проблему до стека ARP ядре. Первоначальное расследование хоста пода Redis показало, что с сетью что-то серьезно не так: связь по любому порту зависала на несколько секунд, а любые DNS-имена не могли быть разрешены через локальный демон dnsmasq, при этом dig выдавал лишь загадочное сообщение об ошибке: socket.c:1915: internal_send: 127.0.0.1#53: Invalid argument. Лог dmesg был более информативен: neighbor table overflow!, что означало переполнение ARP-кэша. ARP используется для сопоставления сетевого адреса, такого как адрес IPv4, с физическим адресом, например MAC-адресом. К счастью, это легко исправить, задав несколько параметров в /etc/sysctl.conf:
Plain Text
1net.ipv4.neigh.default.gc_thresh1 = 800002net.ipv4.neigh.default.gc_thresh2 = 900003net.ipv4.neigh.default.gc_thresh3 = 100000Тюнинг этого параметра является обычной практикой в кластерах высокопроизводительных вычислений (HPC) и особенно актуален в кластерах Kubernetes, поскольку каждый под имеет собственный IP-адрес, который расходует место в ARP-кэше.
Наши кластеры Kubernetes работают без инцидентов уже около 3 месяцев, и в 2018 году мы планируем масштабировать их до еще более крупных размеров. Недавно мы обновились до версии 1.8.4 и рады видеть, что теперь в ней официально поддерживается 5000 узлов. Если вам интересно создавать вычислительные кластеры масштабного уровня, мы нанимаем сотрудников!

Автор
Полный текст статьи читайте на OpenAI
