Главная / Статьи / От модуля к кольцу хешей: масштабирование флота кэш-серверов Node.js без проблем с нагрузкой

От модуля к кольцу хешей: масштабирование флота кэш-серверов Node.js без проблем с нагрузкой

Узнайте, почему метод шардинга hash-mod-N приводит к сбоям в работе баз данных при изменении узлов, как сравниваются алгоритмы хеширования rendezvous, jump и ring, и как создать сбалансированный взвешенный кольцевой структуру в Node.js.

4407 слов

Кластер кэша, который стабильно работал месяцами, может перестать обслуживать базу данных в течение нескольких минут после рутинной замены — добавления одного узла. Причиной обычно является одна строка кода клиента, которая выбирает сервер по формуле hash(key) % N. В этом руководстве подробно объясняется, почему эта строка терпит неудачу, сравниваются четыре серьезных альтернативных решения, показано, как создать кольцо консистентного хэша высокого качества для производственных целей в Node.js, а также перечислены операционные подводные камни, от которых сам алгоритм не сможет вас защитить.

Паттерн сбоев

Представьте здоровый кластер из четырех узлов кэша с коэффициентом успешных запросов 94%. Из-за роста нагрузки перед сезонным пиком инженер добавляет пятый узел. Это изменение конфигурации, состоящее из двух строк, тщательно внедряется в рабочее время. Примерно через полторы минуты база данных заполняется на 100% CPU, и сайт становится недоступным.

Никто не поступил небрежно. Клиент просто направлял ключи так, как всегда делал:

const node = nodes[hash(key) % nodes.length];

Это выражение очень равномерно распределяет ключи, поэтому кажется корректным. Проблема заключается в том, что происходит сразу после изменения значения nodes.length, и именно здесь начинается остальная часть этого руководства. В нем рассматриваются четыре отдельных проблемы, анализируются реальные альтернативы, после чего создается и тестируется структура кольца в Node.js.

Четыре проблемы, скрытые за одной идеей

Консистентное хеширование часто представляют как один трюк. На практике оно решает четыре разных проблемы, и реализация, учитывающая только первую из них, всё равно потерпит неудачу в производственных условиях из-за остальных трех.

Проблема 1: изменение значения N перемещает почти все ключи

При использовании формулы hash(key) % N изменение значения N приводит к перемещению почти всех ключей, а не лишь некоторых из них.

Полезно проанализировать цифры. Ключ с хешем 1,000,003 соответствует узлу 3 при использовании оператора % 4, и к счастливой случайности он также соответствует узлу 3 при использовании оператора % 5. Ключ с хешем 1,000,004 соответствует узлу 0 при использовании оператора % 4 и узлу 4 при использовании оператора % 5. Между этими двумя соответствиями нет никакой связи, поэтому ключ остается на прежнем месте лишь случайно, с вероятностью примерно 1 к N.

При анализе миллиона ключей становится очевидно: переход с 8 на 9 узлов приводит к перемещению 88,93% ключей, а в кластере из ста узлов добавление одного узла делает недействительными примерно 99% кэша.

Обратите внимание, в каком направлении развивается эта тенденция. Чем больше растет система, тем более разрушительными становятся каждые шаги масштабирования. Это сбой, скрытый до тех пор, пока бизнес процветает.

Каждая перенаправленная ключ-запись считается неудачей, каждая неудача приводит к запросу в базу данных, и все они поступают в течение нескольких секунд. База данных, предназначенная для обработки 6% запросов, которые обычно не находят нужных данных, вдруг получает почти все такие запросы.

Проблема 2: тот же процесс перераспределения, но без плана

По крайней мере, первая проблема возникает тогда, когда вы решаете масштабировать систему. Вторая — это тот же процесс, запускаемый сбоем, но в самый некстати момент.

У узла исчерпывается память, хост останавливается работу или сетевое разделение скрывает один узел от половины всей сети. Количество узлов снижается с 8 до 7, и каждый клиент самостоятельно перенаправляет примерно 87% своих ключей на оставшиеся узлы.

Ситуация сейчас критическая: исчезло 12,5% ёмкости кэша, база данных сталкивается с 87%-ным уровнем неудачных запросов, а семь оставшихся узлов берут на себя трафик потерянного узла одновременно с пополнением своих данных почти всеми ключами. Частым результатом является сбой второго узла, что вынуждает проводить ещё одну полную перенастройку, что в свою очередь приводит к сбою третьего узла.

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

Проблема 3: наивное кольцевое строение сильно дисбалансировано

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

Однако при наивной реализации баланс нагрузки соблюдается плохо. Узлы располагаются там, где падают их хэши, поэтому расстояния между ними случайны, а случайные расстояния редко бывают одинаковыми. При размещении по одному узлу на каждый хэш одно из измерений показало, что один узел контролировал 45% пространства ключей, в то время как другой — всего 13%; разница составила 3,4 раза, причем никаких ошибок не было.

Этот дисбаланс также является постоянным. Он обусловлен самими именами узлов, поэтому тот же узел, например cache-04, остается активным до тех пор, пока его не переименуют, и любой, кто изучает код, видит, что он работает ровно так, как написан.

При более крупных масштабах ситуация ухудшается. При наличии по одной точке кольца на каждый из 8 узлов самый загруженный узел обрабатывал 434% от своей справедливой доли, в то время как самый слабо загруженный — всего 3,7%. Это фактически эквивалентно одному перегруженному серверу и семи неактивным.

Задача 4: все клиенты должны согласиться

Самая сложная проблема связана с авторитетностью. Кто-то должен соотносить ключ с узлом, причем это соответствие должно быть одинаковым на всех машинах, которые запрашивают информацию.

Одним из вариантов является использование сервиса-координатора, хранящего авторитетную карту узлов. В этом случае либо каждый поиск требует двух проходов по сети, либо клиенты кэшируют карту, что требует способа её аннулирования. Если у двух клиентов на мгновение окажутся разные версии карты, один может записать user:42 в узел A, в то время как другой будет читать его из узла B. Ничего не теряется, но это, пожалуй, ещё хуже: теперь существует два возможных значения.

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

Конкретный сценарий и варианты

Чтобы сделать сравнение более наглядным, рассмотрим следующую систему.

Система. API электронной коммерции хранит данные сессий и профилей в пуле узлов Redis. Максимальная нагрузка составляет примерно 40 000 запросов в секунду к 25 миллионам ключей при коэффициенте успешных запросов 94%. База данных существует только благодаря обработке тех 6% запросов, которые не находятся в кэше. (Чтобы ознакомиться с основами паттернов кэширования, посмотрите Основы кэширования в Redis.)

Требования:

  1. Увеличить количество узлов с 8 до 12 перед началом продаж без возникновения всплеска неудачных запросов
  2. Выдерживать потерю узла при ограниченном допустимом радиусе воздействия
  3. Сохранять равномерную нагрузку, чтобы ни один узел не работал более чем на 120% от своей нормы
  4. Исключить любые координаторы из пути чтения
  5. Учитывать различия в аппаратном обеспечении: у некоторых узлов 64 ГБ, у других — 16 ГБ, и они не должны нести одинаковую нагрузку

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

Вариант А: хеширование по модулю

hash(key) % N обеспечивает идеальный баланс, требует всего одной инструкции и не использует память.

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

Выбирайте его тогда, когда значение N действительно никогда не меняется, например при распределении задачи среди фиксированного числа работников или при разделении данных внутри одного процесса.

Вариант Б: координатор и таблица поиска

Сохраняйте явное соответствие между диапазонами ключей и узлами в хранилище, таком как etcd или ZooKeeper. Системы вроде Vitess и HBase работают примерно таким образом.

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

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

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

Вариант C: хеширование с встречей (HRW)

Метод с наибольшим случайным весом сравнивает каждый ключ со каждым узлом и выбирает победителя:

function rendezvous(key, nodes) {
  let best = null, bestScore = -1;
  for (const node of nodes) {
    const score = mix(hash(key), hash(node));
    if (score > bestScore) { bestScore = score; best = node; }
  }
  return best;
}

Это весь алгоритм. Здесь нет кольца, нет виртуальных узлов, нет отсортированной структуры, и ничего не нужно перестраивать при изменении состава.

По двум наиболее важным показателям он также превосходит кольцевую структуру. Переход с 8 на 9 узлов привёл к перемещению 11,09% ключей, что меньше теоретического минимума в 11,11%, причём баланс оставался почти идеальным без какой-либо настройки.

Недостаток заключается в сложности O(N) при каждом поиске, так как каждый ключ хешируется относительно каждого узла, и эта стоимость быстро растёт по мере увеличения числа узлов.

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

Вариант D: консистентное хеширование с прыжком

Этот алгоритм, опубликованный Google в 2014 году, занимает примерно десять строк, не требует памяти, обеспечивает почти идеальное сбалансирование и перемещает минимальное количество ключей.

Его ограничение заключается в структуре: он относит ключ к номеру контейнера в диапазоне [0, N) и не учитывает идентичность узлов. Контейнеры могут добавляться или удаляться только в конце диапазона; нет способа удалить узел с номером 3 из середины, сохраняя при этом стабильность всего остального.

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

Вариант E: кольцо хешей с виртуальными узлами

Это классический дизайн. Узлы и ключи делит один круговой пространство адресов, причем ключ принадлежит первому узлу, найденному по часовой стрелке от него.

Выбирайте его тогда, когда количество узлов достаточно велико, чтобы линейные затраты на поиск точек встречи стали значительными, и когда вам также нужны узлы с весами и возможность удаления любого узла.

Выбор под конкретную ситуацию

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

Как работает кольцо

Отложим в сторону массивы и остатки. Представьте круг, пронумерованный от 0 до 2³² − 1, который заканчивается у верхней точки.

Весь алгоритм состоит из двух правил:

  1. Хэш-кодируйте название каждого узла на этом круге, так что cache-01 окажется там, куда указывает его хэш-код.
  2. Хэш-кодируйте каждый ключ на том же круге, затем двигайтесь по часовой стрелке. Узел, к которому вы первым доберетесь, владеет этим ключом.

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

Почему изменения состава остаются локальными

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

Ключи, находящиеся вне этой дуги, не затрагиваются и попадают к своему прежнему владельцу точно так же, как раньше. Новый участник в среднем занимает примерно 1/(N+1) дуги, поэтому доля ключей меняется. При изменении от 8 до 9 этот показатель составил 11,06%, при минимальном значении 11,11%, тогда как при использовании алгоритма модуля при том же изменении и наборе ключей показатель составил 88,93%

Удаление работает в обратном направлении: дуга удаляемого узла переходит к его следующему по часовой стрелке узлу. При сокращении количества узлов с 8 до 7 было перемещено 12,60% ключей, что близко к теоретическим 12,50%. Влияние ограничено и управляемо, при этом остальные шесть узлов остаются нетронутыми.

Виртуальные узлы устраняют дисбаланс

Возвращаемся к задаче 3: четыре узла в четырех случайных положениях создают очень неравномерные дуги.

Решение на удивление простое: не размещайте каждый узел только один раз. Разместите его 160 раз под 160 производными именами вроде cache-01#0 и cache-01#1. Каждое производное имя попадает в разное место, поэтому каждый физический узел имеет 160 небольших рассеянных дуг вместо одной большой, и закон больших чисел сглаживает ситуацию.

При тестировании на миллионе ключей на 8 узлах баланс постепенно улучшается по мере увеличения количества копий. Значение 160 является стандартным по умолчанию, потому что именно там кривая улучшений примерно выравнивается, однако значение 500 все равно даёт заметно лучшие результаты. На одну точку кольца требуется примерно 12 байт (4 байта для указания положения плюс 8 байт для ссылки на владельца), так что восемь узлов с 500 копиями занимают менее 50 КБ. Если для вас баланс важнее, чем объём памяти, выбирайте большее значение; это расчёт, который лишь немногие команды делают.

Виртуальные узлы также практически делают взвешивание бесплатным. Узел с вдвое большей памятью получает вдвое больше баллов, а значит и примерно вдвое больший объем трафика. При коэффициентах взвешивания 4:4:1:1 полученные показатели составили 40,6%, 41,1%, 9,4% и 9,0%, тогда как идеальные значения — 40/40/10/10.

Реализация кольца в Node.js

Реализация сводится к четырем шагам: хэширование имён виртуальных узлов, их размещение, сортировка и использование бинарного поиска для нахождения владельца ключа. Приведённый ниже класс хранит информацию о членах в структуре Map, пересоздаёт отсортированные массивы при изменении состава и предоставляет метод get(key) для поиска.

export class ConsistentHashRing {
  #positions = new Uint32Array(0); // sorted ring positions
  #owners = []; // owners[i] owns #positions[i]
  #nodes = new Map(); // id -> { weight, points }

  constructor({ replicas = 160, hash = defaultHash } = {}) {
    if (replicas < 1) throw new RangeError("replicas must be >= 1");
    this.replicas = replicas;
    this.hash = hash;
  }

  addNode(id, weight = 1) {
    if (typeof id !== "string" || id.length === 0)
      throw new TypeError("node id must be a non-empty string");
    if (weight <= 0) throw new RangeError("weight must be > 0");
    if (this.#nodes.has(id)) return this;
    this.#nodes.set(id, {
      weight,
      points: Math.max(1, Math.round(this.replicas * weight)),
    });
    this.#rebuild();
    return this;
  }

  removeNode(id) {
    if (this.#nodes.delete(id)) this.#rebuild();
    return this;
  }

  #rebuild() {
    const pairs = [];
    for (const [id, { points }] of this.#nodes) {
      for (let i = 0; i < points; i++)
        pairs.push([this.hash(`${id}#${i}`), id]);
    }
    pairs.sort((a, b) => a[0] - b[0]);
    this.#positions = Uint32Array.from(pairs, (p) => p[0]);
    this.#owners = pairs.map((p) => p[1]);
  }

  /** Index of the first ring point >= h, wrapping to 0. */
  #successor(h) {
    const pos = this.#positions;
    let lo = 0,
      hi = pos.length;
    while (lo < hi) {
      const mid = (lo + hi) >>> 1;
      if (pos[mid] < h) lo = mid + 1;
      else hi = mid;
    }
    return lo === pos.length ? 0 : lo;
  }

  get(key) {
    if (this.#positions.length === 0) return null;
    return this.#owners[this.#successor(this.hash(key))];
  }
}

Три аспекта дизайна заслуживают внимания.

Параллельные массивы вместо массива объектов. Хранение координат в Uint32Array позволяет использовать бинарный поиск в компактной, непрерывной памяти, которая остается в кэше процессора. При 8 узлах и 160 копиях получается 1 280 точек, объем примерно 5 КБ, а поиск требует около 11 сравнений.

Эффект обхода границы в lo === pos.length ? 0 : lo. Ключ, хеш-значение которого превышает последнюю точку, относится к первому узлу на круге. Игнорирование этого условия — самая распространенная ошибка при реализации таких структур: результат верен почти для любого ключа, но по какой-то причине ошибочен для немногих ключей, находящихся близко к верхнему пределу диапазона.

Пересоздание структуры при изменении состава, а не при поиске. Изменения в составе происходят редко, тогда как поиски выполняются десятками тысяч раз в секунду, поэтому периодическое сортирование 1 280 записей не имеет значимой стоимости.

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

getReplicas(key, count = 1) {
  const n = this.#positions.length;
  if (n === 0) return [];
  const wanted = Math.min(count, this.#nodes.size);
  const out = [];
  const start = this.#successor(this.hash(key));
  for (let step = 0; step < n && out.length < wanted; step++) {
    const owner = this.#owners[(start + step) % n];
    if (!out.includes(owner)) out.push(owner);
  }
  return out;
}

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

Тщательный выбор функции хэшинга

Во многих учебниках этот раздел опускается, хотя именно он содержит наиболее показательные результаты всего упражнения.

Большинство реализаций кольцевых структур по умолчанию используют MD5. Этот алгоритм работает, но медленно: кольцо, использующее его, достигало всего 423 000 операций поиска в секунду, причем анализ производительности показал, что почти весь времени тратилось на обработку данных с помощью MD5.

Замена его на FNV-1a — быстрый некриптографический хеш — увеличила пропускную способность примерно в 15 раз. Однако баланс разрушился: стандартное отклонение нагрузки на каждый узел выросло с 10,4% до 30,7%

Анализ вычисленных координат нескольких виртуальных узлов позволяет выявить причину:

cache-01#0 → 4037809751
cache-01#1 → 4021032132
cache-01#2 → 4071364989
cache-01#3 → 4054587370
cache-01#4 → 3970699275

Все они сосредоточены в узкой области около 4,0 миллиарда. FNV-1a обладает слабым эффектом аваланчирования, что означает, что похожие входные данные дают похожие выходные результаты. Названия виртуальных узлов отличаются лишь суффиксом, поэтому вместо рассеивания 160 точек по всему кругу каждый узел сгруппировывает их в один плотный кластер. Оптимизация под скорость незаметно вернула проблему 3.

Решение заключается в пропускании выходных данных FNV через финализатор с смешиванием битов, который является последним этапом алгоритма MurmurHash3:

function fnv1a(str) {
  let h = 0x811c9dc5;
  for (let i = 0; i < str.length; i++) {
    h ^= str.charCodeAt(i);
    h = Math.imul(h, 0x01000193);
  }
  return h >>> 0;
}

// Scrambles the bits so near-identical inputs land far apart.
function fmix32(h) {
  h ^= h >>> 16;
  h = Math.imul(h, 0x85ebca6b);
  h ^= h >>> 13;
  h = Math.imul(h, 0xc2b2ae35);
  h ^= h >>> 16;
  return h >>> 0;
}

export const defaultHash = (str) => fmix32(fnv1a(str));

fmix32 чередует сдвиги, операции XOR и умножения, чтобы изменение любого входного бита распространялось на все выходные биты. Math.imul выполняет настоящее умножение 32-битных целых чисел, а >>> 0 преобразует результат обратно в беззнаковое 32-битное число, соответствующее типу Uint32Array. Благодаря примерно десяти дополнительным операциям комбинированный хеш работал лучше, чем MD5, при этом его скорость была примерно в 13 раз выше.

Этот урок актуален не только для хеширования: когда вы заменяете какой-либо компонент на более быстрый, необходимо измерять те характеристики, которые вы не оптимизировали. FNV-1a — это вполне приемлемый хеш; он просто не подходит для этой задачи, и в его описании нет никаких предупреждений об этом.

Преобразование «кольца» в клиент кэша

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

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

export class ShardedCache {
  #ring; #clients; #down = new Map();

  constructor(clients, { replicas = 160, failoverDepth = 2, downtimeMs = 10_000 } = {}) {
    this.#clients = new Map(Object.entries(clients));
    this.#ring = new ConsistentHashRing({ replicas });
    for (const id of this.#clients.keys()) this.#ring.addNode(id);
    this.failoverDepth = failoverDepth;
    this.downtimeMs = downtimeMs;
  }

  #markDown(id) {
    this.#down.set(id, Date.now() + this.downtimeMs);
    this.#ring.removeNode(id);
  }

  #reviveExpired() {
    const now = Date.now();
    for (const [id, until] of this.#down) {
      if (now >= until) { this.#down.delete(id); this.#ring.addNode(id); }
    }
  }

  async get(key) {
    this.#reviveExpired();
    for (const id of this.#ring.getReplicas(key, this.failoverDepth)) {
      try {
        return { value: await this.#clients.get(id).get(key), node: id };
      } catch {
        this.#markDown(id);
      }
    }
    return { value: null, node: null, allDown: true };
  }
}

При запуске этого кода на четырёх симулированных узлах, хранящих по 10 000 ключей, с последующим уничтожением одного из них было получено следующее:

Keys per node:  cache-01 2213 | cache-02 2462 | cache-03 2923 | cache-04 2402

Killing cache-02...
  served from cache: 7538
  cache misses:      2462
  hard failures:     0
  healthy nodes:     cache-01, cache-03, cache-04

  24.6% of traffic became a miss.

Эта цифра — это ваш план масштабирования. В флоте из четырех узлов одна неисправность приводит к тому, что дополнительно 24,6% операций чтения падает на базу данных, тогда как в флоте из восьми узлов этот показатель снижается до 12,5%. Если база данных не может выдержать нагрузку в размере 1/N от общего объема операций чтения, поступающих одновременно, то настоящая проблема заключается в ее ёмкости, а не в кэшировании, причем метод консистентного хеширования превращает эту проблему из фатальной в заметную. Также стоит отметить, что жестких сбоев не было: запросы на ключи недоступного узла перенаправлялись на следующий узел, где они становились обычными неудачами.

Проблемы в производственной среде, которые алгоритм не учитывает

Существует несколько проблем, которые находятся вне рамок действия алгоритма и всё равно нанесут вам ущерб.

Несоответствие версий кольца

Кольцо будет эффективно только в том случае, если каждый клиент вычисляет один и тот же результат. Постепенно внедряйте изменения в состав членов, и в течение нескольких минут половина флота будет видеть 8 узлов, а другая половина — 9. Эти две группы будут расходиться во мнениях по примерно 11% ключей. Для кэша это означает небольшое снижение коэффициента успешных запросов; для любого ресурса, принимающего записи, это означает наличие разных данных. Версионируйте набор членов, распространяйте его через единственный канал и указывайте версию в своих метриках, чтобы вы могли наблюдать отклонения, а не только догадываться о них.

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

Не удаляйте узел сразу после одного истечения времени ожидания. Узел, постоянно выходящий и возвращающийся, вызывает «бурю смены узлов», поскольку каждая смена приводит к перемещению 1/N ключей. Для удаления узла необходимо, чтобы в течение определённого временного окна произошло несколько последовательных сбоев, после чего его следует осторожно вновь принять. В приведённом выше примере используется фиксированное наказание в десять секунд; в производственном коде следует применять экспоненциальное замедление попыток и проверку работоспособности перед тем, как вновь разрешить узлу подключение.

Горячие ключи не сбалансированы

Консистентное хеширование равномерно распределяет ключи, но ничего не говорит об запросах. Когда какой-то продукт внезапно становится чрезвычайно популярным, его запись оказывается на одном сервере, и этот сервер перегревается, несмотря на то что система корректно функционирует. Существует два решения: небольшая кэш-память внутри процесса перед системой для самых популярных ключей или вариант ограниченной нагрузки консистентного хеширования, разработанный в Google, который ограничивает объем данных, которые может принять любой узел, и направляет излишек вдоль часовой стрелки.

Переназначение — это не миграция

Всё вышесказанное исходит из предположения, что потеря ключа приводит лишь к невозможности обращения к кэшу. Если кольцо передаёт постоянные данные, то фраза «11% ключей перемещены» означает, что 11% данных необходимо физически скопировать на новый узел, прежде чем их можно будет прочитать там; обычно чтения или записи осуществляются в обоих местах до завершения копирования. Кольцо определяет, что нужно переместить, но само не занимается процессом перемещения. Именно поэтому такие системы, как Vitess, зависят от координатора: им необходимо контролировать процесс миграции, а не просто его рассчитывать.

Количество реплик фактически является постоянным

Изменение значения replicas с 160 на 500 приводит к смещению всех позиций в кольце и перераспределению почти всех ключей, что создаёт столько же проблем, сколько и изменение модуля. Рассматривайте это как единовременное решение по проектированию, принятое до наличия активных данных; если сомневаетесь, выбирайте более высокое значение.

Когда кольцо — не подходящий инструмент

Частью инженерного суждения является умение понимать, когда впечатляющее решение на самом деле неверно.

  • используйте алгоритм хеширования с встречей. Он требует меньше кода, обеспечивает лучшее балансирование и не требует настройки количества копий. Единственным преимуществом такой структуры является скорость поиска O(log N), которая при таком размере имеет второстепенное значение.
  • Переменные контейнеры, где меняется только их количество: используйте алгоритм «jump hash» — он требует всего десяти строк кода и не потребляет памяти.
  • Данные, требующие контролируемой перераспределения: используйте координатор. Только с помощью явного маппинга можно переместить отдельный шард, одновременно контролируя его влияние, чего невозможно достичь с помощью хеширования.
  • N, которое действительно никогда не меняется: достаточно использовать формулу % N. Не стоит создавать сложные механизмы для изменений, которые никогда не произойдут.

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

Основные выводы

  • Основная идея проста: размещать ключи и серверы в одном адресном пространстве, чтобы изменение состава группы влияло только на близлежащие узлы, а не на всё целое.
  • Модульное маршрутизирование приводит к практически полной очистке кэша при любом событии масштабирования или сбое узла, причём ущерб растёт с размером кластера.
  • Хэшинг типа «рендеву» часто является лучшим выбором для небольших флотов кэшей; обращайтесь к кольцевой структуре, когда важны размер, взвешивание узлов и возможность произвольного удаления элементов.
  • Без виртуальных узлов один сервер может обрабатывать объём данных, в несколько раз превышающий его справедливую долю. Решайте количество копий заранее, так как изменение этого показателя позже приведёт к полной перестройке всей структуры.
  • Быстрый хеш с слабым эффектом аваланча может незаметно вновь привести к дисбалансу; всегда измеряйте распределение, а не только скорость.
  • Потеря одного узла приводит к тому, что примерно 1/N запросов направляется в базу данных. Размер базы данных должен учитывать это, изменения в составе версий, удаление узлов следует осуществлять осторожно, а проблемные ключи — обрабатывать отдельно.
  • Сам кольцевой механизм состоит всего из нескольких десятков строк кода. Инженерные решения, обеспечивающие его надежность в производственных условиях, — это всё, что окружает его.

    Связанные материалы