← назад к разделу

Чат или лента статусов заказа на одном экземпляре сервиса работает без хитростей: клиент держит соединение, сервер по нему пишет. Потом трафик растёт, вы поднимаете три реплики в Kubernetes — и часть сообщений перестаёт доходить. Ошибок нет, соединения живы, а адресат молчит.

Причина простая: WebSocket-соединение живёт на конкретном поде, а событие «пришло новое сообщение» рождается там, куда балансировщик отправил запрос отправителя. Это другой под, и у него в памяти нет соединения получателя. Статья про то, как связать поды между собой шиной событий, что при этом настроить в Kubernetes и как жить с тем, что долгие соединения рвутся при каждом выкате.

Главное

  • Соединение привязано к поду, событие рождается где угодно: без общей шины между подами доставка работает только на одной реплике.
  • Шина — топик в Kafka: событие публикуется один раз, каждый под читает его своей группой потребителей и отдаёт тем соединениям, которые держит сам.
  • Ключ сообщения — идентификатор получателя: так порядок событий одного пользователя сохраняется.
  • Kubernetes ничего не знает про долгие соединения: таймауты Ingress, выкат и масштабирование по CPU нужно настраивать под них отдельно.
  • Каждый выкат рвёт все соединения пода: закрывайте их мягко и по очереди, а на клиенте переподключайтесь с задержкой и разбросом.
  • Разрыв — не авария, а часть протокола: у клиента есть номер последнего события, у сервера есть откуда дочитать пропущенное.
  • Доставка «хотя бы раз»: события могут прийти дважды, клиент отбрасывает повторы по идентификатору.
Обязательно

Почему «просто добавить реплик» ломает доставку

Для обычного HTTP реплики взаимозаменяемы: запрос попал на любой под, тот сходил в базу и ответил. Состояние живёт в базе, под ничего не помнит. У WebSocket иначе. После рукопожатия TCP-соединение остаётся открытым, и объект соединения лежит в памяти конкретного пода. Отправить сообщение пользователю может только тот под, который держит его сокет.

Теперь пользователь А пишет пользователю Б. Запрос А балансировщик отправил на под 1. Под 1 сохранил сообщение и ищет соединение Б у себя в памяти. А соединение Б живёт на поде 3. С тремя репликами доставка удаётся примерно в трети случаев — ровно тогда, когда оба попали на один под.

Два ложных решения, которые приходят первыми. Привязка сессий на балансировщике (sticky sessions) не помогает: она отправляет запросы одного пользователя на один под, но отправитель и получатель — разные пользователи. Одна реплика решает проблему ценой отказа от масштабирования и отказоустойчивости: выкат этой реплики — пауза для всех.

Настоящее решение — перестать доставлять напрямую. Под, где родилось событие, публикует его в общую шину. Каждый под слушает шину и доставляет события тем, кого держит сам.

Шина событий между подами

Шиной удобно делать топик в Kafka, особенно если она уже есть в системе. Схема:

  1. Сервис сохранил сообщение в базу и опубликовал событие в топик chat-events с ключом — идентификатором получателя.
  2. Каждый под сервиса подписан на топик. У каждого пода своя группа потребителей: имя группы равно имени пода. Так каждый под получает все события топика, а не долю.
  3. Под смотрит в свою карту «пользователь → соединения». Если получатель здесь — пишет в сокет. Если нет — молча пропускает событие: его получит тот под, где соединение есть.

Отдельная группа на под — не ошибка, а суть широковещания. Обычно группу делают одну на сервис, чтобы разделить работу между подами; здесь работу делить нельзя, потому что не известно заранее, на каком поде нужное соединение. Поэтому топик читают все.

Три детали, без которых схема протекает. Для новой группы ставьте чтение с конца (auto.offset.reset=latest): свежий под после выката не должен перечитывать историю и рассылать вчерашние сообщения. Группы погибших подов Kafka сама забудет через срок хранения смещений (по умолчанию семь дней), но при частых выкатах их будут сотни — назовите группы по имени пода, чтобы в мониторинге было видно, чьи они. И ключ сообщения: события одного получателя с одним ключом попадают в одну партицию, и порядок «сообщение, потом его удаление» не перепутается.

Сколько это стоит? Каждый под читает весь поток. При десяти подах и тысяче событий в секунду брокер отдаёт десять тысяч событий в секунду — для Kafka это немного, а под тратит время только на поиск в своей карте соединений. Схема перестаёт устраивать на сотнях подов или при тяжёлых событиях; тогда переходят к адресной доставке, о ней ниже.

Что меняется в Kubernetes

Kubernetes видит поды и порты, а не долгие соединения. Четыре места, где его настройки по умолчанию против вас.

Ingress и таймауты. WebSocket проходит через Ingress как обычный HTTP-запрос с заголовком Upgrade, и большинство контроллеров его поддерживают. Но таймаут простоя у них считан на короткие запросы: у nginx-контроллера это шестьдесят секунд, и молчащее соединение он закроет. Поднимайте proxy-read-timeout и proxy-send-timeout до часа и посылайте ping каждые полминуты — и как сигнал жизни, и чтобы соединение не считалось простаивающим.

Service без привязки. Привязка сессий на Service не нужна: после рукопожатия соединение и так живёт на одном поде. Она понадобится только если клиент рядом с WebSocket делает HTTP-запросы, которые обязаны попасть на тот же под. Лучше не проектировать так: все данные, нужные любому поду, должны лежать в базе или в шине.

Масштабирование не по CPU. Тысяча висящих соединений почти не грузят процессор, зато держат память и файловые дескрипторы. Автомасштабирование по CPU не заметит перегрузки. Масштабируйте по своей метрике «соединений на под» и поднимите предел открытых файлов в контейнере.

Выкат. Rolling update по очереди убивает старые поды. Каждый убитый под рвёт все свои соединения, клиенты переподключаются разом — и приходят на оставшиеся поды волной. Что с этим делать — следующий раздел.

Выкат без шторма переподключений

Когда Kubernetes останавливает под, он посылает процессу SIGTERM и ждёт terminationGracePeriodSeconds (по умолчанию тридцать секунд), после чего убивает. Нормальное завершение с WebSocket выглядит так:

  1. Под перестаёт принимать новые соединения: проверка готовности отвечает отказом, и Service убирает его из адресов. Убирание занимает секунды, поэтому в preStop ставят паузу в пять-десять секунд — иначе новые рукопожатия ещё летят в умирающий под.
  2. Под закрывает соединения с кодом 1001 (Going Away) — не все разом, а порциями: тысяча соединений за десять секунд, а не за одну.
  3. Клиент, получив 1001, переподключается через случайную задержку: базовая секунда, умноженная на попытку, плюс разброс. Без разброса тысяча клиентов постучится в одну и ту же миллисекунду.

Период ожидания ставьте с запасом на порционное закрытие: шестьдесят секунд для нескольких тысяч соединений. И ограничьте скорость выката: maxUnavailable: 1 и бюджет прерываний (PodDisruptionBudget), чтобы при переезде узлов кластер не остановил половину подов сразу.

Разрывы и дочитывание

Соединение рвётся не только при выкате: сеть телефона, сон ноутбука, перезапуск Ingress. Пока клиент был отключён, события шли — и Kafka их доставила подам, у которых этого клиента уже не было. Чтобы ничего не потерять, нужны две вещи.

У каждого события — номер по порядку в потоке пользователя, клиент запоминает последний полученный. После переподключения он присылает этот номер, и сервер отдаёт всё, что после него. Откуда отдаёт? Не из Kafka: смещения там считаются по партициям, а не по пользователям, и «дай события пользователя Б после номера 118» топик не умеет. Дочитывают из хранилища, где события уже лежат по-своему: таблица сообщений с порядковым номером или список последних событий пользователя в Redis. Kafka — для живой доставки, хранилище — для восполнения.

Из этого следует режим доставки «хотя бы раз»: событие могло уйти по сокету за миг до разрыва и повториться при дочитывании. Клиент отбрасывает повторы по номеру события. Строго один раз в такой системе не бывает, и обещать его не нужно.

Дополнительно: при первом чтении можно пропустить

Глубже: адресная доставка через реестр соединенийрасширенное

Когда подов много или события тяжёлые, широковещание расточительно: девяносто девять подов из ста читают событие впустую. Альтернатива — знать, где соединение. Поды пишут в общий реестр (Redis): «пользователь Б → под 3», с временем жизни и продлением по ping. Отправитель смотрит в реестр и публикует событие уже адресно: в топик с ключом пода или напрямую в под по внутреннему адресу.

Цена — реестр как ещё одна точка правды, которая расходится с реальностью. Под умер без записи о выходе — в реестре остался мёртвый адрес, пока не истечёт время жизни. Пользователь открыл два окна — соединений два, на разных подах, реестр должен хранить список. Поэтому реестр заводят, когда широковещание измеримо дорого, а не заранее.

Вариант «одна группа потребителей на сервис, партиции по пользователю» тоже кажется естественным: тогда события пользователя Б читает один под. Но соединение Б установил не тот под, который получил его партицию, а тот, куда его отправил балансировщик, и придётся пересылать между подами. А при каждом ребалансе партиции переезжают, соединения — нет. Схема рабочая, но сложнее реестра.

Глубже: медленные клиенты и обратное давлениерасширенное

Под получил из шины событие для тысячи соединений, а десять из них на медленной сети. Запись в сокет блокируется или копит очередь. Если очередь на соединение не ограничена, один клиент за городом съест память пода. Правило: у каждого соединения своя очередь с пределом, при переполнении — закрыть соединение с кодом 1013 или выбросить старые события, если для них есть дочитывание. Читать из Kafka быстрее, чем пишем в сокеты, нельзя: лаг потребителя вырастет, и события начнут приходить поздно всем.

Глубже: набросок конфигурациирасширенное

Deployment с мягким завершением и пределом на выкат:

spec:
  replicas: 3
  strategy:
    rollingUpdate: { maxSurge: 1, maxUnavailable: 1 }
  template:
    spec:
      terminationGracePeriodSeconds: 60
      containers:
        - name: gateway
          lifecycle:
            preStop:
              exec: { command: ["sh", "-c", "sleep 8"] }
          readinessProbe:
            httpGet: { path: /ready, port: 8080 }

Ingress для долгих соединений:

metadata:
  annotations:
    nginx.ingress.kubernetes.io/proxy-read-timeout: "3600"
    nginx.ingress.kubernetes.io/proxy-send-timeout: "3600"

Потребитель шины, одинаковый на любом языке:

group.id = имя пода            # каждый под читает весь топик
auto.offset.reset = latest     # новый под не перечитывает историю

на каждое событие из chat-events:
    соединения = карта[событие.получатель]
    если пусто — пропустить
    для каждого соединения:
        положить событие в его очередь (с пределом)
        если очередь полна — закрыть соединение кодом 1013

Что почитать дальше