Это перевод статьи Бена Чесса, в которой автор рассказывает, как спланировал и осуществил эксперимент по созданию Kubernetes-кластера на миллион узлов.
Зачем?
Несколько лет назад, работая в OpenAI, я помогал писать статью «Scaling Kubernetes to 7500 Nodes», которая до сих пор остаётся одним из самых популярных постов в блоге CNCF. Затем Alibaba рассказала о работе с кластерами Kubernetes на 10 тысяч узлов, а Google совместно с Bayer Crop Science — о 15 тысячах узлов. Сегодня GKE поддерживает кластеры размером до 65 тысяч узлов, а AWS недавно анонсировала поддержку до 100 тысяч узлов.
На форумах и в разговорах с коллегами я часто сталкиваюсь со спорами о том, насколько большим может быть кластер Kubernetes. Но им обычно не достаёт конкретики и экспериментальных данных. Мне доводилось работать с инженерами, которые боялись выходить за рамки уже известного им масштаба из-за страха перед проблемами, которые потенциально могут возникнуть. Или если что-то шло не так, то их первой реакцией было уменьшить размер кластера, а не разобраться в причинах и устранить узкое место.
Суть проекта k8s-1m — нащупать жёсткие препятствия для роста кластера и попытаться найти ответы на следующие вопросы:
Каковы реальные пределы возможностей конкретной системы и чем они обусловлены?
Что нужно сделать, чтобы обойти эти ограничения и отмасштабировать систему ещё сильнее?
Какие типы рабочих нагрузок поддаются масштабированию лучше остальных?
Многие ограничения создаются самим софтом. И наша задача как инженеров — преодолевать эти программные лимиты. Только когда мы сделаем всё возможное на стороне ПО, а проблемы будут упираться исключительно в железо, можно будет сложить руки и ждать оборудования нового поколения.
У меня нет планов монетизировать этот проект — по крайней мере, не в его текущем виде. Я с удовольствием делюсь кодом. Понятно, что все заявления и результаты в этом проекте реальны и проверяемы, но у каждого крупного кластера Kubernetes своя уникальная специфика рабочих нагрузок, которая может оказаться несовместимой с предлагаемой архитектурой. Поэтому я не позиционирую своё решение как готовое к production и не рекомендую использовать его для критически важных задач.
Тем не менее, надеюсь, что предложенные здесь архитектурные идеи покажут сообществу, на что способен Kubernetes, и повлияют на будущую дорожную карту проекта.
Компоненты
В самом начале работы над проектом я спросил у ChatGPT: «Хочу масштабировать Kubernetes до 1 миллиона узлов. Какие проблемы мне придётся преодолеть?»
Вот часть того, что я получил в ответ:
«Масштабирование Kubernetes до 1 миллиона узлов — сложнейшая задача, требующая преодоления множества технических препятствий. Вот ключевые направления, над которыми придется поработать:
Масштабируемость и производительность etcd: etcd является основой хранилища Kubernetes и обрабатывает данные всех объектов API. В случае 1 миллиона узлов объём управляемых etcd данных существенно возрастет. Крайне важно оптимизировать её производительность, включая эффективное партиционирование данных, настройку для работы под высокими нагрузками, а также, возможно, потребуется поиск альтернатив или дополнений к etcd для определенных типов данных.
Производительность kube-apiserver: в кластере такого масштаба kube-apiserver будет подвергаться колоссальным нагрузкам. Оптимизация его производительности включает горизонтальное масштабирование (запуск нескольких экземпляров kube-apiserver), ограничение частоты запросов (rate limiting) и создание эффективных механизмов кеширования. Кроме того, важно минимизировать watch-события и сократить общее количество запросов к kube-apiserver.
Сеть: управление сетевым трафиком для 1 миллиона узлов представляет собой нетривиальную задачу. Она включает управление IP-адресами, обеспечение эффективной маршрутизации и, возможно, сегментацию сети, чтобы она не превратилась в узкое место. Может потребоваться внедрение решений вроде Cilium, которые используют технологию eBPF для сетевого взаимодействия и обеспечения безопасности.
Планирование и управление ресурсами: время отклика планировщика возрастёт, поскольку ему придётся оптимально размещать поды на огромном количестве узлов. Потребуется модернизация планировщика для эффективной обработки операций в таком масштабе, в том числе, возможно, путём внедрения кастомных планировщиков для решения конкретных задач».
Далее в списке упоминаются и другие проблемы: логирование, мониторинг, безопасность, а также управление хранилищем.
Но первые пункты я считаю самыми важными. Вариантов настроить логи или подключить PVC полно, но если у вас не работает etcd, не планируются поды или они не в состоянии общаться между собой, то далеко вы не уедете.
Ниже расскажу о трёх главных направлениях, которыми занимался.
Сетевое взаимодействие
Классически сложности с сетью в Kubernetes связаны с пулом IP-адресов, обнаружением сервисов и сетевыми политиками (то есть с файрволами). По сравнению со следующими этапами настроить сеть для 1 миллиона узлов оказалось довольно просто.
IP-адреса подов
Планирование пула IP-адресов в больших кластерах — непростая задача. Для эффективной маршрутизации каждому узлу обычно выделяют непрерывный диапазон CIDR. Это ограничивает число подов, которое на нём можно запустить. В результате некоторые узлы с мелкими нагрузками исчерпывают доступные IP-адреса раньше, чем у них заканчиваются аппаратные ресурсы. Диапазон 10.0.0.0/8 дает 16 миллионов IP-адресов. Для кластера из миллиона узлов это означает всего по 15 IP-адресов на узел, чего явно недостаточно.
Выход один — полностью и безоговорочно переходить на IPv6. Огромный пул IPv6-адресов позволяет без проблем раздать каждому поду по глобально доступному IP-адресу.
С поддержкой IPv6 у Kubernetes все отлично, и для создания полностью функционального кластера, работающего исключительно на IPv6, не требуется вносить какие-либо изменения в исходники.
Моя цель — сделать так, чтобы каждый узел получал префикс IPv6-адресов с диапазоном, достаточным для выделения индивидуального IP-адреса каждому поду на этом узле. Плюс хотя бы один адрес должен остаться для самого хоста.
Разумеется, для IPv6 нужна поддержка со стороны облачного провайдера. Все крупные облачные провайдеры (и даже многие игроки поменьше) в той или иной мере поддерживают IPv6. С публичным IPv6 легко поднять единый кластер, охватывающий разные облака.
Я преимущественно ориентировался на AWS, GCP и Vultr (можете иронизировать над Vultr, но у них дешёвые виртуалки, а я плачу за этот проект из своего кармана). При этом у каждого из них свои нюансы в том, как реализована поддержка адресации IPv6 внутри виртуальных машин. Чтобы показать разницу в подходах, кратко опишу их ниже.
Vultr: каждый узел получает подсеть /64. Основным IP-адресом узла является ::1 из этого диапазона. При этом сервер автоматически принимает весь трафик, направленный на любой IP-адрес из подсети.
GCP: каждый узел получает подсеть /96. Основной IP-адрес узла выбирается случайным образом из этого диапазона. Сервер должен отправлять корректные пакеты NDP RA для тех IP-адресов, трафик для которых он хочет принимать.
AWS: каждый узел получает подсеть /128. Можно добавить префикс /80 (из другого диапазона) с помощью API-вызова к существующему сетевому интерфейсу или виртуальной машине. (Хотя API Create, по всей видимости, должен поддерживать одновременную установку IPv6-адреса и IPv6-диапазона при создании, но это приведёт к ошибке). Сервер должен отправлять корректные пакеты NDP RA для тех IP-адресов, трафик для которых он хочет принимать, а все исходящие пакеты должны использовать тот единственный MAC-адрес, который соответствует основному IP-адресу.
Чтобы удовлетворить все эти требования (особенно ограничение на MAC-адреса), я создаю общий сетевой мост для интерфейсов всех подов. Интерфейс хоста остаётся отдельным, а для передачи трафика между мостом и интерфейсом хоста включается переадресация пакетов. Локальный IPAM хоста настраивается на префикс IPv6 /96, полученный от провайдера. В итоге мы получаем 232 IP-адресов на узел — более чем достаточно для подов.
Поскольку это глобальные публичные IPv6-адреса, никакая особая маршрутизация не требуется. Инкапсуляция пакетов или NAT не используются. Трафик от каждого пода уходит с его реальным исходным IP-адресом в качестве отправителя, независимо от назначения.
Зависимости от внешних IPv4-сервисов
C IPv6-адресом можно обращаться только к другим IPv6-адресам в интернете. Всё, что работает исключительно по IPv4, напрямую недоступно.
Большинство сервисов, который я использовал в проекте, работали без проблем: репозитории Ubuntu, пакеты PyPi, реестр образов docker.io. Главным исключением стал GitHub. Он до сих пор поддерживает только IPv4.
У многих сервисов AWS есть эндпоинты с поддержкой dual-stack, но вот Elastic Container Registry (ECR), который мне нужен, её лишен. За это им минус.
Чтобы устройства с IPv6 могли связываться с IPv4-хостами, большинство облачных провайдеров предлагают шлюзы NAT64. Также можно поднять собственный шлюз на виртуальной машине Linux. Я здесь немного перемудрил и настроил кастомный сервер WireGuard. Все виртуальные машины подключаются к нему через WireGuard и используют его как IPv4-шлюз.
Сетевые политики
Если говорить в общих чертах, то я просто закрыл глаза на эту проблему и не стал настраивать сетевые политики между рабочими нагрузками.
На 1 миллион узлов пришлось бы завести 1 миллион отдельных IPv6-префиксов — ни один файрвол такое количество записей не переварит. Специалисты по безопасности могут схватиться за сердце, но я не использовал масштабное межсетевое экранирование для защиты кластера от внешнего интернета. Настроил лишь правила файрвола для ограничения доступа к нескольким конкретным портам, но в остальном для защиты серверов и подов от несанкционированного входящего трафика пришлось полагаться на другие методы.
Почти все потребности безопасности в этом проекте закрывает повсеместное использование TLS. Кроме того, из-за огромного адресного пространства IPv6 сканировать сеть практически бесполезно. Изолировать поды друг от друга могли бы Cilium, kube-proxy или другие сетевые плагины, но это создало бы огромную нагрузку на управляющий слой (control plane) из-за избыточных watch-подписок.
Если использовать одного провайдера для всех узлов, возможно, они получат IPv6-адреса из одного или нескольких крупных диапазонов. Такое небольшое количество правил уже можно без проблем прописать в файрволе.
Сетевой трафик (количество TCP-соединений)
Как kube-apiserver’ы, так и etcd поддерживают HTTP/2 и gRPC. Множество отдельных запросов и потоков мультиплексируются в рамках одного TCP-соединения. В Kubernetes по умолчанию установлено ограничение HTTP/2 в 100 одновременных запросов (точнее, потоков) на одно TCP-соединение. (HTTP/2 может поддерживать гораздо больше, но с ростом числа потоков возникают проблемы с производительностью, например блокировка начала очереди — head-of-line blocking.) Соответственно, каждому kubelet'у требуется хотя бы одно подключение к управляющему слою kube-apiserver. Ещё по одному соединению понадобится для kube-proxy или любого другого CNI-плагина вроде Cilium или Calico. Для 1 миллиона узлов это означает, что kube-apiserver должен обрабатывать не менее 2 миллионов TCP-соединений. В случае 8 инстансов kube-apiserver'а на каждый сервер будет приходиться по 250 тысяч подключений от kubelet'ов.
Linux способна справиться с таким количеством соединений после небольшой настройки. И, конечно, нужно убедиться, что системе выделено достаточное количество файловых дескрипторов. Однако такая нагрузка может превысить лимиты облачного провайдера. Например, Azure заявляет о лимите в 500 тысяч входящих и 500 тысяч исходящих соединений на одну виртуальную машину. GCP и AWS не публикуют точные лимиты, но ограничения на общее число одновременных соединений и скорость создания новых существуют в любой системе.
Управление состоянием
Под «управлением состоянием» я подразумеваю API-интерфейс, который Kubernetes предоставляет для работы с ресурсами. Тонкая настройка позволяет масштабировать kube-apiserver до весьма высокой пропускной способности. Но etcd в этой схеме остаётся узким местом. В этом разделе я объясню, почему так, и представлю альтернативное решение, которое способно справиться с нагрузками кластера из миллиона узлов.
kube-apiserver против etcd
Сначала кратко разберем, как устроена работа с состоянием в Kubernetes. Множество различных клиентов взаимодействуют с kube-apiserver'ами, которые, в свою очередь, работают с etcd.

kube-apiserver’ы не хранят состояние (stateless). Полноценным постоянным хранилищем для всех ресурсов Kubernetes служит etcd. Все CRUD-операции, направляемые в kube-apiserver, в конечном итоге сохраняются в etcd.
Для работы с состоянием kube-apiserver использует семь стандартных методов-глаголов (verbs):
create
get
list
update (оно же replace)
patch
delete
watch
У etcd их четыре:
put
range (включая get с пустым range_end)
deleteRange
watch
Операции create, update, patch и delete kube-apiserver'а соответствуют операции put в etcd (при этом удаление — это просто put с пустым значением). etcd не поддерживает частичное обновление данных, только перезапись всего значения. Поэтому любая операция, изменяющая ресурс, приводит к вызову etcd put для всего его содержимого.
Операция watch kube-apiserver'а приводит (но не всегда) к watch в etcd. Подробнее об этом ниже.
Выжимаем нужный QPS для кластера на 1 миллион узлов
kubelet'ы взаимодействуют с kube-apiserver’ом в основном через ресурсы двух типов:
Node — ресурс, представляющий сервер для запуска подов;
Lease — обновляемый kubelet’ом легковесный объект для периодической проверки работоспособности узла.
Lease играет критически важную роль: если он не обновляется вовремя, NodeController помечает узел как неработоспособный (NotReady). По умолчанию каждый kubelet обновляет свой Lease каждые 10 секунд. В масштабе 1 миллиона узлов это даёт 100 тысяч операций записи в секунду только для поддержания статуса работоспособности узлов.
Если добавить к этому постоянное обновление других ресурсов, системе придётся проводить сотни тысяч операций записи в секунду параллельно с большим объёмом операций чтения.
В случае kube-apiserver'а эта задача вполне выполнима. Поскольку он не хранит состояние, показатель QPS можно легко масштабировать, запустив дополнительные реплики. Если один инстанс не справляется с нагрузкой, добавляются новые, и трафик распределяется между ними.
С etcd ситуация обстоит иначе. etcd сохраняет состояние (stateful), что сильно усложняет масштабирование операций записи.
etcd слишком медленная
Утилита etcd-benchmark показала, что пропускная способность инстанса etcd с NVMe-накопителем составляет около 50 тысяч операций записи в секунду. Важно, что добавление реплик никак эту картину не меняет. Пропускная способность на запись на самом деле снижается при увеличении количества участников, так как каждую запись необходимо согласовывать через кворум реплик для соблюдения консистентности. Поэтому в стандартном сетапе из трёх реплик реальный QPS на запись будет даже ниже тестовых 50 тысяч в секунду. Этого крайне мало для кластера из 1 миллиона узлов.
На первый взгляд, 50 тысяч QPS кажутся на удивление скромным результатом для современного оборудования. Один NVMe-диск может выполнять более 1 миллиона записей блоками по 4 КБ в секунду, а один модуль памяти DDR5 — на порядок больше. Почему же etcd так сильно отстаёт от физического предела железа?
Ответ кроется в интерфейсе и гарантиях, которые даёт etcd. Во-первых, etcd обеспечивает гарантированную запись данных на диск. При каждом вызове put или delete etcd выполняет fsync на диск, прежде чем подтвердить успешное выполнение запроса. Это гарантирует сохранность данных при сбое хоста или отключении питания. Но погоня за такой надёжностью резко снижает реальный показатель IOPS, который может выдать современный NVMe-диск.
Во-вторых, etcd обладает весьма широким набором функций:
Работает как хранилище «ключ-значение» и поддерживает чтение, запись и удаление одиночных объектов.
Умеет возвращать выборку диапазона отсортированных ключей.
Хранит историю изменений, позволяя запрашивать старые версии определённого ключа или диапазона ключей. Устаревшие изменения со временем сжимаются (compacted) для экономии места.
Поддерживает механизм watch — стриминг изменений, происходящих с конкретным ключом или диапазоном ключей.
Имеет Lease API, где ключам задаётся время жизни (TTL), по истечении которого они автоматически удаляются, если Lease не был продлен.
Поддерживает транзакции для выполнения атомарных проверок If/Then/Else.
Реализация всех этих интерфейсов делает программу довольно сложной. Помимо базовых операций записи и удаления, etcd должна управлять транзакциями, хранить историю версий и координировать консенсус Raft между репликами.
Именно эти возможности обеспечивают согласованность и надёжность Kubernetes, но они же ограничивают производительность. Строгая согласованность требует сериализации: операции нельзя свободно распараллелить. Запись должна пройти весь путь через Raft, WAL (лог предзаписи) и сжатие, чтобы все реплики согласовали состояние, прежде чем вернуть успешный ответ.
В итоге физическое железо способно обрабатывать миллионы записей в секунду, но etcd из-за необходимости соблюдать гарантии выдаёт результат на несколько порядков ниже.
Но действительно ли нам нужны все эти возможности?
Снижаем требования к сохранности и убираем реплики
Думаю, это мое самое смелое утверждение за весь проект: большинству кластеров на самом деле не нужен тот уровень надёжности и сохранности данных, который предоставляет etcd.
Как мы увидим в следующем разделе, основная часть записей в кластере Kubernetes относится к эфемерным ресурсам.
События (Events) в Kubernetes «живут» всего несколько минут.
Объекты Lease обычно устаревают за десятки секунд.
Если в работе кластера произойдёт сбой, восстанавливать эти объекты обычно нет никакого смысла, и уж точно нет нужды восстанавливать их с точностью до миллисекунд. Даже в случае с долгоживущими объектами Kubernetes изначально спроектирован так, чтобы выполнять реконсиляцию (синхронизировать состояние) автоматически:
Узлы постоянно обновляют свой статус через kubelet.
Контроллеры приводят состояние DaemonSet'ов и Deployment'ов в соответствие с реальными подами.
Если бы мы перестали делать fsync для этих временных записей или вообще отказались от их сохранения на диск, оставив только в оперативной памяти, кластеры могли бы обрабатывать гораздо больше операций и работать значительно быстрее.
Более того, в некоторых окружениях не страшна даже полная потеря данных управляющего слоя. Многие кластеры сами по себе эфемерны, а вся их конфигурация описана в виде кода в Terraform, Helm и/или других GitOps-инструментах. В таких случаях пересоздать кластер с нуля зачастую проще, чем пытаться сохранить каждую операцию записи. Некоторые компании уже относятся к кластерам Kubernetes как к расходникам.
Если вы ещё не начали плеваться, вот вам добавка: реплики etcd, скорее всего, тоже не нужны.
За 5 лет работы с кластерами Kubernetes в OpenAI мы ни разу не сталкивались с незапланированным падением виртуальных машин, на которых работала etcd. Ресурсов ей требуется минимум: размер базы данных ограничен 8 ГБ, процессору нужно не более 2–4 ядер. Большинство облачных провайдеров умеют выполнять живую миграцию таких небольших виртуальных машин. А при наличии сетевых дисков (например, EBS) восстановление выглядит предельно просто: запускаем новую ВМ, подключаем диск и продолжаем работу без потери данных.
При падении единственного инстанса etcd отключился бы управляющий слой кластера. При этом сами поды продолжат работать, узлы останутся доступными, а кластер, скорее всего, и дальше будет обслуживать трафик. Если etcd использовала EBS, восстановление займёт ровно столько времени, сколько требуется на запуск новой виртуалки и подключение диска — без какой-либо потери данных.
Да, одиночный инстанс etcd — это единая точка отказа. Однако сбои происходят редко, а их реальные последствия часто незначительны. В то же время реплики etcd ощутимо бьют по производительности. Для многих рабочих нагрузок такой компромисс просто того не стоит.
В любом случае стоит отключить запись Event и Lease на диск. Далее есть несколько вариантов:
Живучесть данных не критична: запускаем одну реплику и храним всё состояние в оперативной памяти.
Допустима потеря нескольких миллисекунд обновлений: используем одну реплику с сетевым диском, но отключаем
fsync.Хочется избежать потери данных: запускаем несколько реплик на случай сбоя, но не пишем изменения на диск. Полагаемся на доступность остальных реплик для сохранности данных.
Максимальная защита от потери данных: запускаем одну реплику с сетевым диском и включённым
fsync.
Урезаем интерфейс
Как уже говорилось, у etcd довольно раздутый интерфейс. Но использует ли Kubernetes все эти возможности на практике?
Чтобы проверить это, я написал простенькую утилиту под названием etcd proxy. Этот прокси-сервер находится между Kubernetes и etcd, незаметно перенаправляя весь трафик и одновременно фиксируя в логах каждый запрос и ответ.
После этого запустил кластер Kubernetes и прогнал через него Sonobuoy — стандартный набор тестов на соответствие требованиям Kubernetes. Sonobuoy последовательно тестирует все возможности API Kubernetes, проверяя систему на соответствие официальным стандартам. Запустив его через прокси, я собрал полный трейс реальных запросов и нагрузок, с которыми сталкивается etcd в стандартном кластере.
Выяснилось, что Kubernetes на самом деле задействует лишь небольшую часть интерфейса etcd.
Разумеется, там присутствуют операции read, write, range и watch, но все они укладываются в несколько простых паттернов.
Txn-Put
Kubernetes действительно использует транзакционные запросы (Txn), но они всегда имеют следующий вид:
{ "method": "/etcdserverpb.KV/Txn", "request": { "compare": [ { "key": "SOMEKEY", "modRevision": "SOMEREV", "target": "MOD" } ], "success": [ { "requestPut": { "key": "SOMEKEY", "value": "..." } } ], "failure": [ { "requestRange": { "key": "SOMEKEY" } } ] } }
Другими словами, выполнить запись (put), если ревизия изменения (modRev) этого ключа совпадает с заданным значением, иначе — вернуть текущую версию. Это имеет смысл, поскольку Kubernetes часто обновляет (через update или patch) существующие ресурсы. Чтобы безопасно преобразовать это в put всего ресурса, нужно гарантировать, что исходный ресурс не изменился в процессе выполнения.
Lease'ы
Обратите внимание, что объекты Lease в Kubernetes и Lease в etcd — это не одно и то же. Lease'ы Kubernetes реализованы как обычные пары ключ-значение (K/V) в etcd. Сам Kubernetes создаёт крайне мало Lease’ов именно на уровне etcd.
Основная область, где Kubernetes использует Lease'ы etcd — это объекты Events, например:
{ "method": "/etcdserverpb.Lease/LeaseGrant", "request": { "TTL": "3660" }, "response": { "ID": "7587883212297104637", "TTL": "3660" } } { "method": "/etcdserverpb.KV/Txn", "request": { "compare": [ { "key": "/registry/events/NAMESPACE/SOMEEVENT", "modRevision": "205", "target": "MOD" } ], "failure": [ { "requestRange": { "key": "/registry/events/NAMESPACE/SOMEEVENT", } } ], "success": [ { "requestPut": { "key": "/registry/events/NAMESPACE/SOMEEVENT", "lease": "7587883212297104637", "value": "..." } } ] } }
Всё это нужно для того, чтобы задать событиям адекватный TTL (время жизни). Для модели консистентности Kubernetes эта штука вообще не критична.
Range’ы
etcd можно было бы реализовать в виде простой хеш-таблицы со скоростью вставки O(1), если бы не range-запросы. Они возвращают отсортированный список ключей в заданных пределах, из-за чего данные приходится хранить в упорядоченной структуре. Вставка в отсортированный список или B-дерево имеет сложность O(log n). На мой взгляд, поддержка range-запросов — самое жёсткое ограничение, с которым etcd приходится мириться ради совместимости с Kubernetes. Увы, без этого не обойтись.
К счастью, можно использовать предсказуемую структуру пространства ключей:
/registry/[$APIGROUP/]$APIKIND/[$NAMESPACE/]$NAME
Range-запросы обычно ограничиваются либо конкретным неймспейсом, либо всеми неймспейсами для определённого типа (Kind) ресурса. Kubernetes не выполняет range-запросы, которые охватывают несколько типов ресурсов одновременно (например, Pods и ConfigMaps вместе).
Открывается интересная возможность: вместо одного глобального B-дерева для всего пространства ключей можно поддерживать отдельные B-деревья для каждого типа ресурса. Это сокращает эффективное значение n в сложности O(log n) до количества объектов только одного типа, что ускоряет как вставку, так и выполнение запросов.
Ещё один нюанс заключается в использовании лимитов (limit) при range-запросах. Kubernetes редко запрашивает все объекты сразу; запросы обычно возвращают по 500, 1000 или 10 000 результатов за раз. Однако в ответах на range-запросы также ожидается наличие поля count, отражающего общее количество оставшихся объектов. Это сводит на нет преимущество использования limit, так как даже простой подсчёт всех оставшихся ключей может оказаться ресурсоёмким.
На практике, однако, Kubernetes не требует, чтобы значение count было точным. Ему достаточно знать, что доступно больше результатов, чем указано в limit. Такое мягкое требование позволяет использовать приближённые значения, что открывает возможности для дальнейшей оптимизации.
mem_etcd: кастомная реализация in-memory etcd
Я написал на Rust утилиту mem_etcd. Она прикидывается обычным etcd, но со всеми описанными выше упрощениями. Кроме того, она обеспечивает полностью корректную семантику API etcd, от работы которого зависит Kubernetes.
mem_etcd опирается на две основные структуры данных:
хеш-таблицу для хранения всего пространства ключей;
B-дерево, индексирующее ключи внутри каждого префикса.
Каждое значение также хранит несжатую историю ревизий для данного ключа. Такой подход позволяет выполнять запись по уже существующим ключам за O(1), тогда как запись по новым ключам и range-запросы выполняются за O(log n) (где n — количество ресурсов данного Kind). Range-запросы также требуют дополнительных линейных операций в пределах заданного лимита.
Пусть слово «mem» в названии вас не смущает — mem_etcd умеет в сохранность данных, записывая WAL на диск. Каждый префикс типа /registry/[$APIGROUP/]$APIKIND/[$NAMESPACE/] пишется в отдельный файл. По умолчанию запись буферизуется, поэтому вызовы put могут завершаться до того, как данные окончательно запишутся на диск. Это поведение можно изменить с помощью флага CLI, который включает fsync и принудительно сбрасывает все данные на диск перед успешным завершением put. Кроме того, можно сделать так, чтобы определённые префиксы вообще не сохранялись на диск.

fsync ограничивает производительность на уровне около 100 тыс. операций/сек., в то время как буферизация на диск позволяет преодолеть планку в 1 млн. etcd тяжело приходится на больших масштабах даже при записи на ramdisk, где fsync по идее вообще ничего делать не должен
fsync задержки сильно растут, потому что операции записи выстраиваются в очередь% (cd /tmpfs ; etcd-3.5.16 --snapshot-count=9999999999 --quota-backend-bytes=9999999999) & % parallel -j $X --results out_{#}.txt './benchmark put --total 10000000 --clients 1000 --conns 10 --key-space-size 10000000 --key-size=48 --val-size=1024' ::: {1..$X}
Эти тесты прогонялись на паре инстансов типа c4d-standard-192-lssd: на одной ВМ был запущен mem_etcd, а на другой — клиент для тестирования. По результатам хорошо видно, насколько пагубно включение fsync сказывается на пропускной способности и задержках. Стоит учесть, что для сравнения использовалась одиночная реплика etcd v3.5.16, запущенная на ram-диске с tmpfs. Это идеальные условия для etcd, поскольку физический диск не задействуется вообще, а вызов fsync (хоть и остаётся системным) по сути отрабатывает вхолостую. В то же время mem_etcd писал свой WAL на локальный NVMe-накопитель (в терминологии GCE — Titanium SSD). Хотя у этого типа ВМ есть 16 локальных дисков, в тесте использовался только один.

% timeout 10 parallel -j $X --results out_{#}.txt './etcd-lease-flood -num-keys 1000 -workers 100 -key-prefix {#}' ::: {1..$X}
etcd-lease-flood — это кастомный бенчмарк, призванный имитировать доминирующий тип нагрузки в крупных кластерах Kubernetes. Каждый клиент создаёт по 100 объектов Lease напрямую в etcd, используя то же кодирование protobuf, что и Kubernetes. Для каждого объекта Lease клиент непрерывно шлёт put-запросы на обновление в плотном цикле, пытаясь обновить Lease как можно быстрее.
Watch()
Существует несколько разновидностей watch-запросов, и у каждой своя производительность. Давайте рассмотрим их подробнее.
В документации Kubernetes хорошо описано, как именно kube-apiserver обрабатывает параметры watch-запроса:

resourceVersion не задан: получение состояния и запуск с самой актуальной версии.
Watch запускается с самой последней версии ресурса, которая должна быть консистентной (она запрашивается из etcd через кворумное чтение). Чтобы установить начальное состояние, watch сначала отправляет синтетические события «Added» для всех экземпляров ресурсов, существующих на момент стартовой версии. Все последующие watch-события отражают изменения, произошедшие после этой версии.
resourceVersion="0": получение состояния и запуск с любой версии.
Watch запускается с любой версии ресурса; предпочтение отдается самой последней доступной версии, но это не обязательно. Чтобы установить начальное состояние, watch сначала отправляет синтетические события «Added» для всех экземпляров ресурсов, существующих на момент стартовой версии. Все последующие watch-события отражают изменения, произошедшие после этой версии.
resourceVersion="{значение, отличное от 0}": запуск с конкретной версии.
Запускает watch со строго опредёленной версии ресурса. Watch-cобытия отражают все изменения, произошедшие после указанной версии ресурса. В отличие от режимов «получение состояния и запуск с самой актуальной версии» и «получение состояния и запуск с любой версии», здесь watch не начинается с синтетических событий «Added» для указанной версии ресурса.
Давайте представим всё это в виде таблицы:
resourceVersion не установлен |
resourceVersion=0 |
resourceVersion>0 |
|
Читается из кеша kube-apiserver (вместо etcd) |
❌ |
✅ |
✅ |
Включает начальный список объектов |
✅ |
✅ |
❌ |
Таким образом, операции watch часто предшествует вызов list. List возвращает моментальный снимок (snapshot) набора ресурсов на определённый момент времени, помеченный номером ревизии. Затем вы запускаете watch с этого номера ревизии, и он начинает транслировать все изменения, произошедшие с тех пор.
Когда задан параметр resourceVersion, watch-запрос к kube-apiserver не порождает новый watch-запрос к etcd. При запуске kube-apiserver создаёт watch-потоки к etcd для каждого из стандартных, общеизвестных ресурсов. Всякий раз, когда клиент создаёт watch, kube-apiserver обрабатывает этот поток самостоятельно, опираясь на тот единственный watch-поток к etcd, который поддерживает. Таким образом, хотя клиентские watch-запросы могут быть ресурсоёмкими для kube-apiserver'а, они не оказывают дополнительной нагрузки на etcd. Это позволяет масштабировать kube-apiserver горизонтально.
Да и в целом watch-запросы не так страшны для etcd. У watch есть начальный и конечный диапазоны, которые укладываются в те же префиксы, что и range-запросы. При каждом put нужно выполнить поиск за O(log n) в списке активных watch-запросов, чтобы найти совпадения по ключу. Но самих watch-запросов гораздо меньше, чем объектов. Значение n здесь невелико, к тому же поиск происходит асинхронно уже после фиксации записи, так что на время выполнения запроса на запись это не влияет.
Что действительно создаёт проблемы, так это мультиплексирование сетевого трафика. На каждую запись в etcd может приходиться N активных watch-запросов для этого объекта. В результате объём исходящего трафика из etcd сильно возрастает. Получателями данных выступают kube-apiserver'ы. Они группируют свои watch-запросы, но etcd всё равно отправляет копию данных на каждый kube-apiserver. Хотя добавление реплик kube-apiserver'а помогает решить многие проблемы масштабирования, каждая новая реплика увеличивает нагрузку на сетевой интерфейс etcd. Пропускная способность сети etcd — это первое аппаратное узкое место в крупных кластерах Kubernetes. Впрочем, эта проблема касается исключительно канала между etcd и kube-apiserver'ами. В рамках одного дата-центра с современным железом вполне реально обеспечить достаточную пропускную способность сетевых соединений между ними.
Количество watch-запросов на узел
Увеличивая количество узлов, я смог измерить, сколько watch-запросов создаёт каждый узел. На одну пару kubelet + kube-proxy приходится:
4 watch-запроса для configmap'ов;
по 2 watch-запроса для объектов типа
pod,secret,service,node;по 1 watch-запросу для примитивов
namespace,endpoint,csidriver,runtimeclass,endpointslice,networkpolicy.
В сумме это дает 18 watch-запросов на узел, то есть 18 млн активных watch'ей на 1 млн узлов. Все они обрабатываются исключительно на уровне kube-apiserver'а и не нагружают etcd напрямую. Если развернуть достаточное количество инстансов kube-apiserver'а, Kubernetes вполне справится с этой задачей.
Update()
Давайте вернёмся к задаче по обработке миллиона Lease'ов от kubelet. Kubelet отправляет запрос update (также известный как replace или put) для своего ресурса Lease каждые 10 секунд.
Вот пример старого объекта Lease:
apiVersion: coordination.k8s.io/v1 kind: Lease metadata: creationTimestamp: "2025-06-26T18:27:28Z" name: my-node namespace: kube-node-lease ownerReferences: - apiVersion: v1 kind: Node name: my-node uid: ef4d9943-841b-49cc-9fc2-a5faab77e63f resourceVersion: "1556549" uid: 7e2ec4e2-263f-4350-9397-76f37ceb83cd spec: holderIdentity: my-node leaseDurationSeconds: 40 renewTime: "2025-07-01T21:41:50.646654Z"
А так выглядит тело вызова update() при продлении этой Lease:
apiVersion: coordination.k8s.io/v1 kind: Lease metadata: creationTimestamp: "2025-06-26T18:27:28Z" name: my-node namespace: kube-node-lease ownerReferences: - apiVersion: v1 kind: Node name: my-node uid: ef4d9943-841b-49cc-9fc2-a5faab77e63f resourceVersion: "1556549" uid: 7e2ec4e2-263f-4350-9397-76f37ceb83cd spec: holderIdentity: my-node leaseDurationSeconds: 40 renewTime: "2025-07-01T21:51:50.650000Z"
Обратите внимание, что значение renewTime обновилось на время, которое на 10 секунд больше предыдущего. (На самом деле renewTime всегда устанавливается на 40 секунд вперёд на случай неудачных или затянувшихся обновлений Lease.)
Ещё одно важное поле — resourceVersion. Когда клиент отправляет запрос update() в kube-apiserver, он прикрепляет resourceVersion, которая была у ресурса до обновления. Это нужно для подстраховки — чтобы убедиться, что никакой другой клиент не вклинился и не обновил ресурс параллельно. Каждый раз, когда ресурс обновляется на сервере, тот присваивает ему новую монотонно растущую resourceVersion. Операция update должна содержать resourceVersion, которая показывает, какую старую версию ресурса планируется заменить. Это защищает от случайного перезаписывания изменений, сделанных другими клиентами.
Логично предположить, что kube-apiserver мог бы просто превратить update в транзакцию Txn-Put в etcd, пробросив эту команду дальше stateless-способом. Однако на практике реализация update в kube-apiserver всегда требует получения полной старой версии ресурса. На это есть пара причин:
Поля на стороне сервера: у некоторых ресурсов есть поля вроде
statusиmanagedFields, которые обновляются только сервером.Проверки допуска: интерфейс проверки admission-контроллеров принимает как старую, так и новую версию ресурса.
Для ускорения update-вызовов kube-apiserver хранит кеши watch-событий для самых популярных ресурсов. При update-вызове старая версия извлекается из локального watch-кеша. Если по какой-то причине старой версии там нет, kube-apiserver сначала отправляет запрос range в etcd, чтобы получить старый ресурс, и только после этого вызывает Txn-Put.
Необходимость делать два синхронных вызова к etcd на каждый update удвоила бы наши потребности в QPS и увеличила задержки, поэтому важно иметь актуальный watch-кеш.
Однако для 1 миллиона узлов это создаёт новое требование и ограничение: kube-apiserver’ы должны обрабатывать поток watch-событий со скоростью не менее 100 тысяч событий в секунду.
В моих тестах именно на этом шаге ситуация становилась критической.
Кеширование и блокировки
kube-apiserver каждую секунду десериализует 100 тысяч вложенных словарей (и, что ещё хуже, выделяет память под них). Всё это он складывает в кеш на базе B-дерева, который прикрыт мьютексом RWMutex. На этом мьютексе начинается жуткая «давка»:
Вызовы
Update()пытаются читать из кеша старые объекты.Успешные запросы
Update()(через финализаторGuaranteedUpdate()) пишут в кеш новые значения.Поток watch-событий из etcd также пишет новые значения в этот же кеш.
Увеличение числа kube-apiserver'ов помогает снизить конкуренцию за блокировки, вызванную операциями update, но не снижает нагрузку от получения watch-событий: каждой реплике kube-apiserver'а по-прежнему приходится обрабатывать полный поток всех изменений. Кроме того, добавление новых реплик kube-apiserver'а создаёт дополнительную нагрузку на etcd — ей приходится рассылать копии этого watch-потока на каждый сервер.
Перевод кеша kube-apiserver'а на структуру B-дерева произошёл сравнительно недавно. Раньше он работал на базе хеш-карты (hash map). Фича включается переключателем BtreeWatchCache, который в Kubernetes 1.32 включён по умолчанию. Насколько могу судить, мотивацией для перехода на B-дерево стало желание ускорить ответы на вызовы List(). Метод List() должен возвращать элементы в отсортированном порядке, а хранение данных в B-дереве значительно повышает скорость его работы. С другой стороны, операции Get() и Update() для существующих элементов теперь выполняются за O(n log n) вместо O(1).
В моих тестах максимум, чего мне удалось добиться от кеша на B-дереве, — 40 тысяч обновлений в секунду на инстансе GCP c4a-standard-72. Кеш начинает отставать от потока watch-событий, и слишком много времени теряется на ожидании блокировки.
При использовании старого кеша на хеш-таблице и 11 реплик kube-apiserver'а ресурсов хватает, чтобы справиться с нагрузкой в 100 тысяч обновлений Lease в секунду.
Сборка мусора
kube-apiserver'ы парсят и декодируют все ресурсы на отдельные поля. Из-за этого ресурсы с большим количеством полей создают кучу мелких объектов в Go, что сильно нагружает сборщик мусора. Добавление новых реплик kube-apiserver'а здесь не поможет, если все они следят за одними и теми же потоками событий. Идеального решения проблемы не существует, но настройка параметров GOMEMLIMIT и GOGC помогает облегчить ситуацию.
Я установил значение GOMEMLIMIT на 10–20 % меньше объёма доступной памяти, а GOGC увеличил до нескольких сотен.
Планировщик
Какой толк от кластера на миллион узлов, если на них невозможно запланировать поды? Планировщик Kubernetes довольно часто выступает узким местом для крупных задач. Попытка запланировать 50 тысяч подов на 50 тысяч узлов заняла около 4,5 минут. Некомфортно долго, да?
Предупреждение
Если вы создаёте поды через контроллеры вроде Deployment, DaemonSet или StatefulSet, имейте в виду: они могут стать узким местом ещё до того, как подключится планировщик. DaemonSet, к примеру, создаёт поды порциями по 500 штук за раз, а затем ждёт подтверждения их создания из watch-потока, прежде чем продолжить (скорость зависит от многих факторов, но вряд ли превысит 5 тысяч в секунду). У планировщика даже нет шанса приступить к работе, пока эти поды не будут созданы.
Для эксперимента с миллионом узлов я поставил себе амбициозную задачу — распланировать 1 миллион подов за 1 минуту. Понятно, что цифра взята скорее с потолка, но уж очень красиво смотрелась магия из букв «М» (миллион, минута).
Плюс мне хотелось сохранить полную совместимость со стандартным kube-scheduler'ом. Конечно, было бы гораздо проще с нуля написать «облегчённый» планировщик, который отлично масштабируется в ограниченных сценариях, но терпит крах в реальных условиях. В планировщике Kubernetes заложено много мощной логики, накопленной за годы работы в самых разных production-средах. Избавляться от неё ради иллюзии «быстрого» планировщика было бы неправильно.
Поэтому постараемся максимально сохранить функциональность и внутреннее устройство kube-scheduler'а. Что же мешает сделать его более масштабируемым?
kube-scheduler хранит состояние всех узлов и запускает цикл со сложностью O(n*p), где каждый под сопоставляется с каждым узлом. Сначала он отфильтровывает узлы, которые вообще не подходят поду. Затем проводится скоринг каждого оставшегося узла, чтобы оценить, насколько хорошо он подходит поду. После этого под отправляется на узел с наивысшей оценкой (или на случайный среди лучших, если оценки совпали).
Примечание
У kube-scheduler'а есть пара встроенных механизмов для повышения производительности:
При большом количестве подходящих узлов он оценивает лишь их часть (в крупных кластерах она может снижаться до 5 %).
Он распараллеливает фильтрацию неподходящих узлов и расчёт оценок для конкретного пода.
Очевидно, что эта задача хорошо параллелится. И надо признать, планировщик успешно с ней справляется. Но ему всё равно приходится повторять эти вычисления для всех узлов. К счастью, их можно не просто распараллелить, а распределить.
Базовая архитектура: шардирование по узлам
Это очень похоже на классический паттерн «распределение/сбор» (scatter/gather), используемый в распределённом поиске. Представьте, что каждый узел — это документ в базе данных, а под — поисковый запрос. Запрос расходится по множеству шардов, каждый из которых отвечает за свою часть документов. Шарды выбирают лучших кандидатов и отправляют их центральному сборщику, который и определяет победителя.

Главное отличие в том, что в поиске документы доступны только для чтения, поэтому запросы можно обрабатывать параллельно без конфликтов. В планировании же принятое решение изменяет состояние «документов» (то есть занимает ресурсы узла). Если окажется, что два пода параллельно запланированы на один и тот же узел, один из них разместится успешно, а второй получит отказ из-за нехватки ресурсов.
Тем не менее в больших кластерах разумно использовать принцип «оптимистичного параллелизма» (optimistic concurrency): предполагать, что поды могут планироваться одновременно без конфликтов. При этом перед фиксацией изменений мы всё равно проверяем результат на наличие конфликта. Если он есть, мы «откатываем» другой под и отправляем его на повторный круг планирования. Вероятность таких коллизий и их влияние на систему крайне малы. Они достаточно незначительны, чтобы параллельная обработка давала огромный выигрыш в пропускной способности, полностью перекрывая затраты на редкие откаты.
Моя первоначальная идея архитектуры распределённого планировщика выглядела следующим образом:


true/false, который определяет, нужно ли привязывать под к этому узлуПолучившийся планировщик представляет собой слегка модифицированную версию стандартного kube-scheduler. В него добавлен кастомный gRPC-эндпоинт для приёма новых подов, изменён код для определения зоны ответственности по узлам, а также создана кастомная точка расширения Permit для отправки выбранного узла обратно в Relay. Permit запускается после фильтрации и оценки узлов, но перед привязкой пода. Плагины Permit возвращают true или false, одобряя или отклоняя планирование пода на конкретном узле.
Такой базовый дизайн очень хорошо себя зарекомендовал. Он не совсем подходит для масштаба в 1 миллион узлов (об этом мы поговорим в следующем разделе), но выглядит гораздо более масштабируемым решением, чем то, которое существует сейчас, и при этом сохраняет все тонкости сложной и проверенной временем логики стандартного планировщика.
Текущий планировщик работает со сложностью O(n×p), где n — количество узлов, а p — количество подов. С ростом n такая сложность становится критической. Шардированный подход помогает справиться с проблемой масштабирования: при наличии n узлов работу можно распределить между r реплик (где r — некоторый делитель n), что делает вычисления гораздо более управляемыми.
Стоит упомянуть одно довольно важное исключение — вытеснение (eviction) подов. Оно происходит, когда нужно запланировать новый под, но свободных ресурсов в кластере не хватает. В таких случаях планировщик сканирует все запущенные в кластере поды, пытаясь найти группу подов с более низким приоритетом, удаление которых освободило бы достаточно места. Честно говоря, я поленился реализовывать эту логику. Если постараться, можно представить, как прошардировать и её, но я этого делать не стал.
Проблемы «длинного хвоста»: суровая реальность крупных распределённых систем
На моем оборудовании один планировщик справлялся с фильтрацией и оценкой 1000 узлов для одного пода примерно за 1 мс. Получается, что за 1 секунду можно распределить 1000 подов по 1000 узлов. Но не забывайте, что наша цель — запустить миллион подов на одном миллионе узлов за 60 секунд. Поскольку общая сложность задачи равна O(n x p) (каждый под нужно оценить для каждого узла), переход от 1000 подов и узлов к 1 млн увеличивает объём необходимой работы не в 1000 раз, а в миллион (1000 × 1000). И даже с учётом того, что у нас есть 60 секунд вместо одной, планировщиков потребуется гораздо больше.
Добавление ретрансляторов (Relay) и распределённый сбор оценок
На самом деле планировщиков потребуется так много, что у одного ретранслятора просто не хватит пропускной способности сети, чтобы вовремя отправить данные всем им. Потребуется несколько ретрансляторов. Более того, для связи со всеми планировщиками структура ретрансляторов должна быть многоуровневой.
Этап сбора данных — когда мы собираем все оценки и определяем победителя — тоже можно распределить. У каждого планировщика и ретранслятора есть эндпоинт для сбора оценок (Score Gather). По хешу имени конкретного пода определяется, какой именно планировщик отвечает за сбор оценок для него.
Вот упрощённый пример схемы с несколькими ретрансляторами. Здесь показан коэффициент ветвления (fanout), равный 3, хотя в реальности использовался fanout 10. Я старался максимально загрузить сетевые карты, не превышая их предельную пропускную способность на передачу. Цель — отправить данные подов объемом 1 млн × 4 КБ за 60 секунд.
Обратите внимание на компоновку: планировщики находятся на разных уровнях.

Борьба с задержками «длинного хвоста»
Изначально я рассчитывал на линейное сокращение времени работы при увеличении числа реплик. На деле же столкнулся с плато: сколько бы новых реплик ни развёртывал, ситуация не менялась или даже ухудшалась. Хотя в среднем на большинство планировщиков приходилось меньше работы, поэтому они заканчивали её быстрее, периодически один или два из них сильно отставали. Это ломало весь процесс, ведь для выбора самого подходящего узла требовались ответы от всех планировщиков без исключения.
Эта проблема подробно описана в известной работе Джеффа Дина из Google под названием «The Tail at Scale». Наши серверы не работают в реальном времени. Они постоянно отвлекаются на фоновые задачи: сбор метрик, обновления, сборку мусора. Сборка мусора — огромная проблема в Golang при написании жёстко координируемого ПО. Она прерывает текущую задачу или задерживает выполнение следующей в очереди. В итоге операция, которая обычно занимает 300 микросекунд, внезапно растягивается до 1 миллисекунды. При большом количестве серверов какой-нибудь из них обязательно попадёт на эту задержку. Если распределённая система жёстко завязана на синхронный ответ от каждого узла, то 99 % серверов будут просто ждать этот отстающий 1 %.
Чтобы побороть эту проблему, я придумал несколько хитростей:
Использовал закрепление (pinning) ядер CPU. Оно позволяет закрепить ядра процессора за процессами конкретного контейнера, избегая постоянного переключения контекста на случайные системные процессы. В kubelet для этого используется политика CPU Manager. Одно только это изменение заметно стабилизировало производительность.
Настроил сборщик мусора (GC). Увеличение параметра GOGC до 100+ снижает суммарные затраты времени на сборку мусора за счёт большего расхода памяти. Высокое значение GOGC вместе с установкой GOMEMLIMIT близко к реальному лимиту памяти контейнера — отличный способ запускать GC только при реальной необходимости. Я выставил GOGC на 800, а GOMEMLIMIT — на 90 % от лимита памяти контейнера.
Перестал ждать отстающих планировщиков. Мы просто не ждём ответов от последних N % реплик, тем самым отсекая «длинный хвост». Но здесь кроется опасность: если узел стабильно тормозит, это может запустить опасный цикл — новые запросы будут перегружать его, пока он окончательно не упадёт.
На самом деле рекомендую прочитать саму работу «The Tail at Scale». В ней рассматривается ещё несколько возможных сценариев и вариантов решения проблемы.
Чего я делать не стал, так это раскидывать рабочие нагрузки на несколько серверов. Можно было бы привязать каждый узел сразу к нескольким шардам, чтобы любой из них мог обрабатывать и оценивать поды для этого узла. Но я побоялся проблем с несогласованностью данных, когда на узел планировалось бы избыточное количество подов (over-subscription) из-за того, что один шард не знает, что другой уже разместил поды на этот узел.
Замена механизма Watch на AdmissionWebHook
Тут до конца не понятно, в чём была проблема. Обычно kube-scheduler узнаёт о новых подах через watch-запрос с фильтром полей (fieldSelector) 'spec.nodeName=' (то есть ищет поды, у которых ещё не прописано имя узла). Но при создании большого количества подов на высокой скорости (больше 5К в секунду) этот watch-поток регулярно зависал на десятки секунд.
Вот один из самых жёстких примеров:

Иногда watch-поток зависал настолько сильно, что планировщик простаивал — хотя в очереди было полно подов, ожидающих планирования.
Чтобы решить эту проблему, я пошёл на радикальное изменение логики работы планировщика. Вместо использования watch'ей я настроил его как ValidatingWebhook. Теперь при создании каждого нового пода kube-apiserver отправлял HTTP-запрос на эндпоинт планировщика прямо во время обработки запроса на создание. Обычно валидирующие вебхуки используются для безопасности — чтобы разрешать или запрещать клиентам установку определённых параметров ресурсов. В моём случае вебхук одобрял абсолютно все поды. Это был просто хитрый способ узнавать о новых подах быстро и синхронно, не связываясь с проблемным watch-потоком.
Результаты
Был поднят кластер из 100 тысяч узлов, после чего измерялось время, необходимое для планирования 100 тысяч подов. Для подов не задавались параметры nodeSelector или affinity.
Каждый планировщик запускался на выделенной виртуальной машине c4d-standard-32: это 32 ядра AMD Turin и 128 ГБ оперативной памяти DDR5. В экспериментах, где использовалось более одной реплики dist-scheduler (распределённого планировщика), также развёртывалась одна выделенная ВМ для dist-scheduler-relay.
Каждая реплика dist-scheduler была настроена на запуск 30 отдельных внутренних планировщиков с параметром параллелизма, равным 2.
В качестве стандартного планировщика использовался kube-scheduler 1.32.3 без модификаций.

Довольно любопытный результат: даже одиночная реплика dist-scheduler справилась со своей задачей заметно лучше стандартного kube-scheduler. При этом изменение настроек параллелизма не дало никакого эффекта — ни в плане производительности, ни по процессорному времени, расход которого остановился примерно на значении 20 (то есть 12 ядер простаивали).
Стоит заметить, что добавление новых реплик dist-scheduler более или менее линейно ускорило процесс. Другими словами, при удвоении числа реплик dist-scheduler время выполнения задачи сокращалось вдвое. Эта тенденция сохранилась и на масштабе в 256 реплик при планировании 1 млн подов, о чём пойдёт речь в следующем разделе.
Эксперименты
1 млн узлов и 1 млн подов с kwok
Тестовый стенд
-
kube-apiserver: kube-apiserver'ы K3s v1.32.4+k3s1 на пяти инстансах c4d-standard-192.
kube-scheduler и kube-controller-manager версии v1.32.4 работают как отдельные процессы на тех же ВМ
feature-gates=kube:BtreeWatchCache=falseБез cloud-controller-manager, traefik или servicelb
etcd: кастомная реализация mem_etcd, работающая на инстансе c4d-highmem-16
-
kubelet
Для распределённого планировщика (dist-scheduler): 285 инстансов c4d-highcpu-32
Для kwok: 7 инстансов c4a-highcpu-32
Оба работают под управлением kubelet из K3s v1.32.4+k3s1
kwok: 100 подов, на которых запущена модифицированная версия kwok v0.6.0
Планировщик подов: 289 реплик (на инстансах с ядрами 8670 AMD Turin) кастомной реализации распределённого планировщика; всего 256 планировщиков и 29 ретрансляторов.
Процедура
Запустить все виртуальные машины. Дождаться полной готовности kwok и распределённого планировщика.
Создать 1 млн узлов с помощью утилиты make_nodes.
Дождаться, пока в базе данных etcd появятся записи для 1 млн узлов и 1 млн Lease'ов.
Создать 1 млн подов с помощью скрипта create-pods.
Дождаться момента, когда для каждого пода будет определено поле
spec.nodeName.
Результаты
На графиках ниже зелёная вертикальная линия указывает на момент создания первого пода, а красная — на момент планирования миллионного пода.
etcd




kube-apiserver






Планировщик



Сравнение kwok и kubelet
До сих пор все эксперименты проводились с использованием kwok вместо реальных kubelet'ов. Насколько это приближено к реальности? Вполне возможно, что профиль нагрузки от kwok сильно отличается от классических kubelet'ов. В таком случае наши тесты не будут отражать поведение реального кластера на 1 миллион узлов, на каждом из которых работает полноценный kubelet.
К сожалению, запуск 1 миллиона kubelet'ов выходит за рамки моего бюджета. Но можно провести эксперимент меньшего масштаба с настоящими kubelet'ами и сравнить профиль их нагрузки с кластером аналогичного размера на базе kwok.
Если хорошенько покопаться в конфигах, можно запустить тест кластера на 100 тысяч kubelet'ов. Идея в том, чтобы одновременно запускать кучу kubelet'ов на одной и той же ВМ. Каждый из них работает в изолированном Linux-неймспейсе хоста, имеет свой собственный IPv6-адрес и диапазон для выделения адресов подам, а также запускает собственную копию containerd для управления вложенными подами.
Развёртывание 100 тысяч kubelet-контейнеров в куче виртуальных машин и управление ими кажется сложной задачей. Если бы только существовала программа для оркестрации всего этого дела… О, точно! Можно же просто сделать Kubernetes Deployment из kubelet'ов!
Из-за общего ядра хоста по-прежнему сохраняются узкие места, где процессы конкурируют за ресурсы. kube-proxy по умолчанию работает с iptables, где любые изменения выполняются с помощью мьютекса. nftables более производителен и лоялен к параллелизму, но всё равно остаётся «бутылочным горлышком». По этой причине эффективнее использовать много мелких виртуальных машин вместо нескольких крупных, чтобы распределить нагрузку и обойти ограничения, связанные с параллелизмом.
Кроме того, чтобы IPv6-подсети каждого kubelet были доступны из облачного провайдера, необходимо передавать пакеты объявлений соседей (neighbor advertisement). Для этого развёртывается ndppd на каждой ВМ в виде DaemonSet'а.
Тестовый стенд
-
kube-apiserver: 6 инстансов c4d-standard-192s, на которых работают kube-apiserver'ы из K3s v1.32.4+k3s1.
Без cloud-controller-manager, traefik или servicelb
etcd: кастомная реализация mem_etcd, работающая на инстансе c4d-highmem-8
-
kubelet:
-
Для сценария kubelet-in-pod: 426 инстансов c4d-highmem-8 (3408 ядер AMD Turin и 27264 ГиБ ОЗУ)
Запущен K3s-версии v1.32.4+k3s1
Для сценария kwok: 2 инстанса c4a-highcpu-32
-
kubelet-pod: используется слегка модифицированный образ K3s v1.32.4+k3s1 с установленной библиотекой libjansson; она добавляет поддержку JSON в nftables, необходимую для корректной работы kubelet.
Процедура
Дождаться загрузки кластера и перехода всех виртуальных машин (узлов) в состояние Ready.
Развернуть Deployment kubelet-как-под (kubelet-as-pod).
Отмасштабировать кластер до 100 тысяч реплик.
Снять графики нагрузки на kube-apiserver и etcd.
Удалить кластер и развернуть его заново.
Развернуть kwok. Создать 100 тысяч узлов kwok.
Снять графики нагрузки на kube-apiserver и etcd.
Результаты
etcd






kube-apiserver










Заключение. Каких размеров может достигать кластер Kubernetes?
Практика показывает, что сам по себе размер кластера играет гораздо меньшую роль, чем интенсивность операций над конкретным типом (Kind) ресурса — в особенности операций create и update. Обработка различных типов ресурсов изолирована: каждый процесс выполняется в своей горутине и защищён собственным мьютексом. Типы ресурсов можно даже шардировать по нескольким кластерам etcd — в этом случае изменения различных сущностей будут масштабироваться независимо друг от друга.
Главным источником операций записи обычно являются обновления Lease'ов — так узлы подтверждают свою работоспособность. Из-за этого лимит масштабирования кластера напрямую зависит от того, насколько быстро система успевает обрабатывать такие обновления.
Стандартная конфигурация etcd на современном оборудовании способна обрабатывать около 50 000 изменений в секунду. При вдумчивом шардировании (если разнести Nodes, Leases и Pods по отдельным кластерам etcd) стандартный etcd справляется с кластером из примерно 500 000 узлов.
Если заменить etcd на более масштабируемое решение, узкое место сместится в сторону кеша подписок (watch cache) в kube-apiserver. В настоящее время доступ к каждому типу ресурса защищён единым мьютексом RWMutex поверх B-дерева. Замена этой структуры данных на хеш-карту позволит обрабатывать порядка 100 000 событий в секунду, чего вполне достаточно для работы 1 миллиона узлов на современном железе. Если и этого недостаточно, можно увеличить интервал Lease'ов, например, до 10 секунд и более, чтобы снизить частоту обновлений.
На предельных масштабах главным тормозом становится сборщик мусора (GC) в Go. Kube-apiserver создаёт и отбрасывает колоссальное количество мелких объектов при парсинге и декодировании ресурсов, что приводит к высокой нагрузке на GC. Простое масштабирование количества реплик kube-apiserver здесь не поможет, так как каждая из них вынуждена подписываться на одни и те же потоки событий.
Как запустить самостоятельно
Подробные шаги по запуску собственного кластера описаны в файле RUNNING.
Дополнительные материалы
How does Alibaba ensure the performance of system components in a 10,000 node Kubernetes cluster
Scaling Kubernetes to 7500 Nodes (OpenAI)
P.S.
Читайте также в нашем блоге: