Назад к блогу

Что внутри Redis Cluster: слоты, госсип и выбор лидера при сбое

Что внутри Redis Cluster: слоты, госсип и выбор лидера при сбое

Redis Cluster часто воспринимают как «просто шардированный Redis», но за этим скрывается несколько независимых механизмов: хеширование ключей в слоты, обмен состоянием между узлами через gossip и выбор нового мастера при отказе. Разбор показывает, как каждый из этих механизмов устроен на уровне протокола — от вычисления номера слота битовой маской до голосования реплик и остановки записи при потере большинства.

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

Как ключ превращается в слот

Слот ключа — это младшие 14 бит CRC16 от ключа. Берутся именно младшие биты, а не остаток от деления, потому что 16384 — это степень двойки: младшие 14 бит дают то же значение, что и CRC16(key) mod 16384, но вычисляются одной битовой маской 0x3FFF без операции деления. В документации та же формула записана как HASH_SLOT = CRC16(key) mod 16384.

Если в ключе есть шаблон {...}, хешируется не весь ключ, а только часть между первой { и первой } справа от неё. Сначала ищется первая {, затем первая } после неё. Если { не найдена, либо } отсутствует, либо между ними ничего нет, хешируется весь ключ. Иначе берётся хеш подстроки между фигурными скобками.

Алгоритм останавливается на первом непустом фрагменте между фигурными скобками, потому что берёт именно первую { и первую } после неё. Для ключа foo{bar}{zap} хешируется bar, так как алгоритм останавливается на первом совпадении { и }. Это позволяет держать связанные ключи в одном слоте, задавая общий префикс в фигурных скобках.

Почему слотов именно 16384

Ключевое пространство кластера разбито на 16384 слота, что задаёт верхнюю границу размера кластера в 16384 мастер-узла. Слот — это младшие 14 бит crc16, то есть маска 0x3FFF — 14 единичных бит.

Число слотов напрямую определяет размер битовой карты слотов. Карта объявлена как unsigned char owner_not_claiming_slot[CLUSTER_SLOTS / 8];, то есть по одному биту на слот. Установка и снятие бита выполняются побитово по позиции слота: вычисляется номер байта как позиция, делённая на 8, и номер бита внутри байта как остаток от деления позиции на 8, после чего нужный бит выставляется.

Так как слотов 16384, карта занимает 16384/8 байт. Именно эти биты узел рассылает в gossip, отмечая, какими слотами он владеет.

MOVED: узел не проксирует, а перенаправляет

Когда узел получает команду для ключа, слот которого он не обслуживает, он не пересылает команду сам, а отвечает клиенту перенаправлением. Ответ формируется как -MOVED <hashslot> <endpoint>:<port>, где hashslot — слот, вызвавший перенаправление, а endpoint и port берутся из узла-получателя. Пример такого ответа: -MOVED 3999 127.0.0.1:6381.

Узлы не проксируют команды, потому что, согласно спецификации, они перенаправляют клиентов к нужным узлам, обслуживающим соответствующую часть ключевого пространства. Клиент затем сам переотправляет запрос на указанный endpoint и port, а при необходимости обновляет карту слотов.

Endpoint в MOVED может быть IP-адресом, hostname или пустым. Пустой endpoint означает, что у сервера неизвестный адрес, и клиент должен отправить следующий запрос на тот же endpoint, но с указанным портом.

Если слот не привязан ни к какому узлу, клиент получает ответ CLUSTER_REDIR_DOWN_UNBOUND, а не MOVED. Для заблокированного клиента, если слот не обслуживается этим узлом и не импортируется, отправляется MOVED к узлу-владельцу слота.

Миграция слота: MIGRATING, IMPORTING и ASK

Миграция слота выполняется вручную через CLUSTER SETSLOT. На источнике команда CLUSTER SETSLOT <slot> MIGRATING <node ID> проверяет, что слот принадлежит этому узлу, находит целевой узел и записывает его как цель миграции. На приёмнике CLUSTER SETSLOT <slot> IMPORTING <node ID> проверяет, что слот ещё не принадлежит этому узлу, находит узел-источник и записывает его как источник импорта.

Затем ключи переносятся командой MIGRATE. Она специально разрешена к локальному выполнению, если слот в состоянии миграции или импорта. Когда клиент обращается к ключу, которого на источнике уже нет, но слот помечен как migrating, узел возвращает ASK и указывает узел-приёмник, потому что ключ мог уже переехать. Если на источнике есть часть ключей, но не все, возвращается ошибка TRYAGAIN. В отличие от неё, ASK и MOVED — это перенаправления: ASK указывает клиенту разово обратиться к другому узлу за конкретным ключом, не меняя карту слотов, а MOVED сообщает, что слот окончательно переехал, и клиент должен обновить карту слотов. Для ещё не перенесённых ключей, которые на источнике есть, редиректа нет — их обслуживает сам источник.

Состояние миграции или импорта определяется по наличию записи о целевом или исходном узле для этого слота.

Как клиент обрабатывает ASK

При получении ASK-редиректа клиент должен отправить только ту команду, которая была перенаправлена, на указанный узел, но продолжать отправлять последующие запросы на старый узел. Перед самой командой клиент обязан отправить ASKING, чтобы целевой узел принял запрос для слота, который он ещё не обслуживает постоянно. ASKING устанавливает у клиента флаг CLIENT_ASKING и отвечает OK.

ASK не означает обновление карты слотов, потому что миграция слота ещё не завершена; клиент не должен обновлять локальные таблицы, чтобы не начать постоянно обращаться к новому узлу раньше времени. Если клиент всё же обновит карту раньше, это не проблема: без ASKING целевой узел ответит MOVED и вернёт клиента к старому узлу. После завершения миграции слота старый узел отправит MOVED, и тогда клиент может окончательно переназначить слот.

Мультиключевые операции во время переселения

Во время переселения слотов мультиключевые операции, такие как MSET, могут временно становиться недоступными. Если все ключи существуют и по-прежнему хешируются в один и тот же слот (либо на исходном, либо на целевом узле), операция остаётся доступной. Если ключи не существуют или в процессе переселения разделены между исходным и целевым узлами, операция возвращает ошибку -TRYAGAIN, и клиент может повторить её позже.

В коде при миграции слота, если не все ключи найдены, но некоторые существуют, возвращается CLUSTER_REDIR_UNSTABLE, что и приводит к ошибке TRYAGAIN. Как только миграция указанного хеш-слота завершена, все мультиключевые операции для этого слота снова становятся доступными.

Cluster bus: отдельный порт для служебного обмена

Кластерная шина (cluster bus) — это отдельный порт, который узел объявляет другим узлам для служебного обмена. Он вычисляется так: если задан cluster-announce-bus-port, берётся он; иначе если задан cluster_port, берётся он; иначе к клиентскому порту прибавляется CLUSTER_PORT_INCR. Клиентский порт по умолчанию — это TLS-порт при включённом TLS-кластере, иначе обычный порт.

Узлы всегда принимают соединения на порту кластерной шины и даже отвечают на ping, даже если пингующий узел не является доверенным. Но все остальные пакеты отбрасываются, если отправитель не считается частью кластера. Узел принимает другого в кластер только двумя способами: если тот представился сообщением MEET (команда CLUSTER MEET ip port), либо если уже доверенный узел распространит о нём gossip-сообщение.

Кого и как часто узел пингует

Узел отправляет PING не всем подряд, а выборочно. clusterCron выполняется 10 раз в секунду. Раз в 10 итераций (примерно раз в секунду) узел пингует один случайный узел. Для выбора проверяются 5 случайных узлов, и пингуется тот, у кого самое старое время pong_received. Кандидаты отбираются случайно, но пропускаются узлы без соединения, с уже активным пингом, сам узел и узлы в состоянии HANDSHAKE.

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

Что несёт gossip-секция PING/PONG

В заголовке сообщения отправителя описывают: имя узла-отправителя, его IP (если задан cluster-announce-ip, иначе поле остаётся нулевым для автообнаружения), основной и вторичный клиентские порты и порт cluster bus. Слоты отправителя передаются в битовой карте: туда копируются слоты мастера, где мастер — это сам узел, если он мастер, иначе его мастер. Эпоху конфигурации несёт configEpoch (это configEpoch мастера), текущую эпоху узла — currentEpoch, состояние кластера с точки зрения отправителя — в поле state, а флаги узла — в поле flags.

При приёме gossip-секции узел по каждому элементу ищет узел по имени. Если отправитель — мастер, добавляется или удаляется failure report по флагам FAIL/PFAIL, и вызывается проверка на переход в FAIL. Для живого узла без активного пинга и без отчётов об отказе обновляется pong_received, если присланное время не из будущего (допуск 500 мс) и больше текущего. Для недоступного узла с другим адресом обновляются ip, порты и сбрасывается флаг CLUSTER_NODE_NOADDR. Неизвестный узел добавляется в словарь, если отправитель известен, нет флага CLUSTER_NODE_NOADDR и узел не в чёрном списке.

Слоты и configEpoch обрабатываются отдельно: слот перепривязывается к отправителю, если он не занят или configEpoch прежнего владельца меньше senderConfigEpoch. При этом слоты, бывшие у прежнего мастера, помечаются как мигрировавшие.

PFAIL и FAIL: локальная догадка против решения кластера

Узел помечает другого как PFAIL (возможный отказ), когда по этому узлу есть активный ping и задержка превышает NODE_TIMEOUT. В коде это состояние отслеживается полем ping_sent, и проверка идёт только если ping_sent != 0.

Просто не получить ответ недостаточно, потому что любой входящий трафик по cluster bus считается доказательством живости. Поэтому задержка вычисляется как минимум из задержки pong и задержки данных, и PFAIL ставится только если эта задержка превышает cluster_node_timeout.

PFAIL — это лишь локальная информация узла о других узлах, и для признания узла отказавшим её нужно повысить до FAIL. Повышение выполняет markNodeAsFailingIfNeeded, которая требует, чтобы узел был в PFAIL и чтобы число отчётов об отказе от мастеров (включая себя, если мы мастер) достигло кворума:

int needed_quorum = (server.cluster->size / 2) + 1;

То есть большинства мастеров.

Как PFAIL превращается в FAIL

Переход выполняется в markNodeAsFailingIfNeeded, который сначала проверяет, что узел уже в PFAIL и ещё не в FAIL. Число отчётов о сбое берётся из clusterNodeFailureReportsCount, и если сам узел мастер, он добавляет себя как голосующего.

Отчёты учитываются только за окно времени: clusterNodeCleanupFailureReports удаляет отчёты старше server.cluster_node_timeout * CLUSTER_FAIL_REPORT_VALIDITY_MULT. Документация уточняет, что большинство мастеров должно сигнализировать PFAIL или FAIL в течение NODE_TIMEOUT * FAIL_REPORT_VALIDITY_MULT, и коэффициент равен 2 — то есть дважды NODE_TIMEOUT.

Когда порог достигнут, узел снимает флаг PFAIL, ставит FAIL и записывает fail_time. Затем он рассылает сообщение FAIL всем доступным узлам и планирует обновление состояния и сохранение конфигурации.

Когда FAIL снимается обратно

FAIL снимается в функции clearNodeFailureIfNeeded, которая вызывается только когда узел помечен FAIL, но снова доступен. Для реплик и мастеров без слотов FAIL снимается сразу при восстановлении связи. Для мастеров, которые всё ещё владеют слотами, FAIL снимается только если прошло больше server.cluster_node_timeout * CLUSTER_FAIL_UNDO_TIME_MULT времени с момента fail_time, то есть никто не выполнил повышение реплики и слоты остаются за этим мастером.

Разница в том, что реплики не участвуют в повышении, а мастера без слотов не участвуют в кластере и ждут конфигурации, поэтому их FAIL можно снять сразу. Мастер со слотами должен подождать, чтобы дать возможность выполнить failover, и лишь затем вернуться в кластер.

Как реплика решает начать выборы

Реплика начинает выборы, когда её мастер в состоянии FAIL, мастер обслуживал ненулевое число слотов и связь репликации была разорвана не дольше заданного пользователем времени, чтобы данные были достаточно свежими.

Перед стартом выборов реплика ждёт задержку:

DELAY = 500 milliseconds + random delay between 0 and 500 milliseconds + REPLICA_RANK * 1000 milliseconds

ранг реплики: реплика с самым свежим смещением имеет ранг 0, следующая — ранг 1 и так далее.

Фиксированная пауза нужна, чтобы дождаться распространения состояния FAIL по кластеру, иначе реплика может попытаться получить голоса, пока мастера ещё не знают о FAIL и отказываются голосовать. Случайная задержка десинхронизирует реплики, чтобы они не начали выборы одновременно. Вклад от ранга штрафует реплики с менее свежими данными, чтобы более обновлённые реплики пытались получить голоса раньше.

Голосование: условия кандидата и условия мастера

Кандидат (реплика) должен выполнить: его мастер в состоянии FAIL, мастер обслуживал ненулевое число слотов, а реплика была отключена от мастера не дольше заданного пользователем времени. Также реплика должна увеличить свой currentEpoch и разослать FAILOVER_AUTH_REQUEST всем мастерам, ожидая ответы не более двух NODE_TIMEOUT (но не менее 2 секунд).

Мастер голосует, только если он сам является мастером, обслуживающим хотя бы один слот. Далее он проверяет, что currentEpoch запроса не меньше его currentEpoch, и что он ещё не голосовал в этой эпохе. Также мастер голосует за реплику, только если её мастер помечен как FAIL (или есть флаг CLUSTERMSG_FLAG0_FORCEACK для ручного переключения), и если с момента предыдущего голосования за реплику этого мастера прошло не менее NODE_TIMEOUT * 2. Наконец, мастер проверяет, что configEpoch запрошенных слотов не меньше configEpoch мастеров, обслуживающих те же слоты в текущей конфигурации.

Эпоха конфигурации как логические часы

configEpoch — это номер версии конфигурации узла, который мастер всегда рассылает в пакетах ping и pong вместе с битовой картой обслуживаемых слотов. При создании нового узла он равен нулю.

currentEpoch — это текущая эпоха узла. Победитель выборов увеличивает счётчик currentEpoch и получает новый уникальный configEpoch: реплика увеличивает currentEpoch как первый шаг для участия в выборах, а после получения авторизации создаётся новый уникальный configEpoch, и реплика становится мастером с этим новым configEpoch. Таким образом, currentEpoch — это счётчик эпохи, который ведёт узел и увеличивает при выборах, а configEpoch — номер версии конфигурации конкретного мастера, который он рассылает вместе с битовой картой своих слотов.

Коллизия эпох разрешается так: если два мастера имеют одинаковый configEpoch, узел с лексикографически меньшим Node ID присваивает себе наибольшую известную эпоху плюс один.

Эпоха — это логические часы, а не время: она определяет, что информация с большей эпохой побеждает информацию с меньшей. Значения configEpoch и currentEpoch гарантированно сохраняются и fsync-ятся на диск при обновлении до продолжения операций узлом.

Почему без кворума мастеров failover не происходит

Автоматический failover запускается только тогда, когда реплика видит своего мастера в состоянии FAIL, а для перевода мастера в FAIL нужно, чтобы его пометило большинство мастеров. Предварительные условия для failover: реплика должна быть слейвом, её мастер должен быть помечен FAIL (или это ручной failover), не должен быть выставлен запрет failover, и мастер должен обслуживать ненулевое число слотов.

Далее реплика рассылает FAILOVER_AUTH_REQUEST всем мастерам и ждёт ответы максимум дважды NODE_TIMEOUT (но не менее 2 секунд), а мастер голосует только если сам обслуживает хотя бы один слот. Победа требует ACK от большинства мастеров: needed_quorum вычисляется как (server.cluster->size / 2) + 1.

Поэтому если кворум мастеров потерян (большинство недоступно), реплика не может собрать нужное число голосов, и повышение не происходит, даже если сама реплика жива и её данные достаточно свежи. Если большинство не набрано за период двух NODE_TIMEOUT (но не менее 2 секунд), выборы прерываются, и новая попытка будет через NODE_TIMEOUT * 4 (но не менее 4 секунд). Когда лишь меньшинство мастеров пометило узел как FAIL, повышение реплики не происходит, и каждый узел снимает FAIL по правилам выше.

Что видит клиент во время сбоя и промоушена

Во время сбоя мастера клиент, обращающийся к стороне меньшинства раздела, получает отказ в записях: сторона меньшинства начинает отказывать в записях, как только истечёт NODE_TIMEOUT без контакта с большинством. До этого момента записи принимаются, но могут быть потеряны: все записи, выполненные на стороне меньшинства до этого момента, могут быть потеряны.

На стороне большинства кластер снова становится доступным после NODE_TIMEOUT плюс несколько секунд на выборы и переключение (failover обычно выполняется за 1–2 секунды). Ошибки возможны и после выбора нового мастера, потому что клиент с устаревшей таблицей маршрутизации может записать старому мастеру до того, как кластер превратит его в реплику нового мастера. Кроме того, после исправления раздела записи ещё некоторое время отклоняются, чтобы другие узлы успели сообщить об изменениях конфигурации.

Сетевой разрыв: отказ, возврат и потеря записей

Узел, отрезанный от большинства мастеров, перестаёт принимать запросы: clusterUpdateState подсчитывает доступных мастеров и, если их меньше кворума (size/2)+1, переводит состояние кластера в CLUSTER_FAIL, что означает отказ в обслуживании.

При восстановлении связи узел не сразу становится доступным: если он мастер и был в меньшинстве, clusterUpdateState возвращает управление без смены состояния, пока не истечёт задержка возврата, ограниченная значениями от 500 до 5000 мс, чтобы успеть получить обновление конфигурации.

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

Что из этого следует на практике

  • Ключи с общим префиксом в {} попадают в один слот. Фигурные скобки хешируются не сами, а содержимое между первой { и первой } — это единственный способ гарантированно держать связанные ключи вместе для мультиключевых операций.
  • Клиент обязан уметь обрабатывать MOVED и ASK по-разному. MOVED означает, что слот переехал окончательно, и карту слотов нужно обновить. ASK означает, что миграция идёт, и обновлять карту рано: перенаправленную команду нужно предварить ASKING, а последующие запросы слать на старый узел.
  • Мультиключевые операции могут временно возвращать TRYAGAIN во время переселения. Это не ошибка клиента, а признак того, что ключи оказались разделены между источником и приёмником; операцию нужно повторить позже.
  • Размер кластера ограничен 16384 мастерами — это прямое следствие того, что слот занимает 14 бит, а битовая карта слотов рассылается в каждом gossip-сообщении.
  • Отказ узла — это решение большинства, а не наблюдение одного узла. PFAIL ставится локально по таймауту активного пинга, но FAIL требует кворума (size/2)+1 отчётов от мастеров за окно в два NODE_TIMEOUT. Без большинства мастеров ни один узел не будет помечен FAIL, и реплика не начнёт выборы.
  • Реплика с самым свежим смещением репликации получает фору. Ранг добавляет к задержке старта по 1000 мс на каждую более свежую реплику, поэтому именно она с наибольшей вероятностью запросит голоса первой.
  • При потере большинства кластер перестаёт принимать записи. Сторона меньшинства отказывает в записях через NODE_TIMEOUT после потери контакта с большинством, а записи, принятые до этого, могут быть потеряны. После восстановления связи узел выжидает задержку возврата (500–5000 мс), прежде чем снова обслуживать запросы.
  • Подтверждённая клиенту запись может быть потеряна. Если мастер подтвердил запись, но не успел её реплицировать и был повышён failover'ом на реплике без этой записи, запись исчезает — это цена асинхронной репликации, а не сбой алгоритма.

Где смотреть в коде

Источники

Похожее