От token ring к tablets: репликация и multi-DC в ScyllaDB

Как устроено распределение данных — от классического token ring к tablets, как ведёт себя multi-DC репликация при отказе ДЦ и что происходит с tablets при decommission узла

Предыдущая статья серии закончилась на том, что LOCAL_QUORUM и QUORUM на одном датацентре различаются на шум измерения — 14 микросекунд, — и что по-настоящему они расходятся только там, где датацентров больше одного. Здесь этот момент наступает: два датацентра в топологии Scylla (multi-DC-семантика без настоящего межрегионального RTT — оба ДЦ стенда живут на одном хосте, об этом честно ниже), живая репликация между ними, честный отказ при потере целого ДЦ и tablets — единица распределения данных, которая заменила собой классический token ring и умеет мигрировать между узлами прямо во время работы кластера, без остановки и без ручной пересборки кольца.

Это пятая статья серии «ScyllaDB: глубокое погружение», продолжение статьи о моделировании данных, статьи об архитектуре shard-per-core, статьи о compaction, tombstones и repair и статьи о consistency levels и LWT. Дальше в серии — shard-aware драйверы на Go и Java и эксплуатация и итоговый выбор.

Стенд этой статьи — не тот же трёхузловой кластер, что в предыдущих четырёх: под multi-DC поднят второй, полностью отдельный живой кластер — два датацентра по два узла (DC1: dc1a, dc1b; DC2: dc2a, dc2b), своя сеть, свой набор seed-узлов в обоих ДЦ разом, чтобы gossip слил их в один логический кластер, а не в два независимых. Основной одноузловой стенд серии при этом не тронут — оба кластера работают параллельно на одном хосте. Код и compose-файл — в публичном репозитории примеров (digital-cookbook, каталог scylladb/topology + scylladb/ops/topology-demo.sh).

Токен-кольцо, tablets и multi-DC в ScyllaDB: репликация DC1 в DC2, отказоустойчивость через LOCAL_QUORUM, миграция tablets 16 в 32 при выводе узла

В статье

Токен-кольцо и tablets: новая единица распределения

Классическая схема распределения данных в Cassandra-совместимых базах — консистентное хеширование: у кластера есть кольцо токенов, каждый узел владеет одним или несколькими диапазонами этого кольца, а партиция попадает на узел по хешу своего ключа партиции. Сама механика кольца — виртуальные узлы (vnodes), их роль в балансировке и болезненность классического ребаланса при добавлении/выводе узла — общая тема для всего семейства таких хранилищ, не специфика ScyllaDB, и отдельно разобрана в статье про шардирование и consistent hashingСкоро. Здесь важнее то, чем ScyllaDB отличается от классической схемы vnodes.

Tablets — новая единица распределения данных, которая в 2026.2.0 (и на этой версии, и на образе, зафиксированном для всей серии) является дефолтной стратегией для новых keyspace вместо vnodes. Вместо статичного владения диапазонами токенов на уровне узла, назначенного один раз при построении кольца, каждая таблица делится на независимые tablets — сравнительно небольшие, динамически управляемые единицы данных со своим списком реплик. Ключевое отличие от vnodes не в терминологии, а в поведении: перераспределение tablets — это лёгкая, инкрементальная операция, которая происходит прямо во время работы кластера (добавление узла, вывод узла, перекос нагрузки), без того ручного и болезненного ребаланса токен-диапазонов, которым славится классическая vnode-схема. Ниже в этой статье tablets показаны не в теории, а в живой миграции: 16 реплик, переехавших с одного узла на другой без остановки кластера и без потери ни одной строки.

RF (replication factor) на multi-DC keyspace задаётся не одним числом, а отдельно на каждый датацентр — NetworkTopologyStrategy реплицирует данные в пределах каждого ДЦ независимо от других: keyspace telemetry_mdc, использованный ниже, объявлен как {'class':'NetworkTopologyStrategy','DC1':2,'DC2':2} — по две реплики в каждом из двух датацентров, четыре реплики партиции суммарно на четырёх узлах кластера.

Честная находка, специфичная для этой сборки: команды nodetool tablets и nodetool cluster tablets, которые звучат как естественный способ посмотреть распределение tablets, в 2026.2.0 не существуютnodetool help не содержит tablets вовсе, а nodetool cluster знает только repair, cleanup и snapshot. Рабочий источник — прямое чтение системной таблицы system.tablets: столбец replicas хранит list<frozen<tuple<uuid,int>>>, по одному элементу (host_id, shard) на каждую реплику диапазона tablet, а host_id каждого узла выясняется отдельным запросом SELECT host_id FROM system.local напрямую к этому узлу — без риска перепутать источник ответа. Тот же system.tablets уже пригодился стенду про repair на tablet-keyspace — там nodetool repair -pr тоже оказался несовместим с tablets и потребовал nodetool cluster repair.

Multi-DC живьём: репликация и цена QUORUM

nodetool status на живом двух-ДЦ кластере подтверждает топологию буквально:

Datacenter: DC1
UN  dc1a  ...
UN  dc1b  ...
Datacenter: DC2
UN  dc2a  ...
UN  dc2b  ...

Сценарий -scenario multidc (n=10000) прогнал запись и чтение через keyspace telemetry_mdc ({DC1:2, DC2:2}), с координатором, жёстко закреплённым в нужном датацентре через DataCenterHostFilter — драйвер физически не открывает пул соединений к узлам другого ДЦ, чтобы измерение не зависело от политики выбора хоста:

Фаза A: 10000 записей LOCAL_QUORUM (координатор DC1, ждёт 2/2 реплик DC1)
Записано: 10000/10000 (ошибок: 0), LOCAL_QUORUM write p50=1.486ms p99=1.624ms

Фаза B: чтение записанных строк LOCAL_QUORUM с координатором DC2
Прочитано из DC2 (LOCAL_QUORUM): совпало=10000, расхождение=0, ошибок чтения=0

Фаза C: 10000 НОВЫХ записей QUORUM (координатор DC1, нужен кворум ВСЕХ 4 реплик — минимум 3 из 4)
Записано: 10000/10000 (ошибок: 0), QUORUM write p50=1.564ms p99=1.769ms

Первый результат — не число, а факт: все 10000 строк, записанные LOCAL_QUORUM в DC1 (координатор ждёт подтверждения только от двух реплик своего датацентра), сразу видны при чтении LOCAL_QUORUM из DC2 — совпадение 10000 из 10000, расхождение ноль. LOCAL_QUORUM не откладывает саму отправку данных в удалённый ДЦ — координатор рассылает запись всем четырём репликам одновременно, он лишь не ждёт подтверждения от DC2, прежде чем ответить клиенту. Репликация между двумя датацентрами этого кластера (в терминах Scylla-топологии) подтверждена не по документации, а по 10000 живых строк.

Второй результат — цена QUORUM относительно LOCAL_QUORUM, которую предыдущая статья обещала, но не могла показать на одном ДЦ. QUORUM для этого keyspace требует кворума всех четырёх реплик кластера — минимум 3 из 4, а значит минимум одно подтверждение обязательно должно прийти из DC2. Результат: QUORUM дороже LOCAL_QUORUM на ~78 микросекунд (~5%) по p50 и на ~144 микросекунды (~9%) по p99 — та же качественная надбавка, что и на одном ДЦ в прошлой статье, но здесь она наконец не тонет в шуме измерения, потому что QUORUM теперь обязан дождаться подтверждения из второго датацентра, а не остаётся в пределах одного. «Второй» здесь — топологический, а не географический: оба ДЦ по-прежнему на одном хосте (см. оговорку сразу ниже), но кворум уже честно пересекает границу между датацентрами Scylla-топологии.

Честная оговорка здесь обязательна и звучит так же, как в предыдущей статье: оба датацентра этого стенда физически сидят на одном Docker-хосте, за одним сетевым мостом — это не настоящий межрегиональный RTT, который на реальном multi-region развёртывании обычно измеряется единицами-десятками миллисекунд, а не десятками микросекунд. Направление эффекта (QUORUM ≥ LOCAL_QUORUM) архитектурно верно и воспроизведено безукоризненно, но абсолютная величина надбавки — 5% по p50, 9% по p99 — артефакт локальной топологии стенда, а не оценка того, во сколько раз QUORUM подорожает при реальном разносе датацентров по разным городам или континентам. На реальном multi-region кластере эта разница была бы измеримо больше — и именно поэтому LOCAL_QUORUM, а не QUORUM, обычно выбирают умолчанием для операционной нагрузки multi-DC развёртывания: цена ожидания удалённого ДЦ на каждой операции там уже не единицы процентов.

DC failover: LOCAL_QUORUM выживает, EACH_QUORUM честно отказывает

Следующий шаг стенда останавливает оба узла DC2 (dc2a и dc2b) — полный отказ целого датацентра, а не одного узла из него, при том что оба узла к этому моменту ещё живы (это важно: decommission ниже необратимо выводит один из узлов DC2, и после него смоделировать «упал весь ДЦ из двух узлов» уже не получится — поэтому порядок шагов в стенде фиксированный).

-- LOCAL_QUORUM (DC1) при упавшем DC2 --
LOCAL_QUORUM write rc=0 read rc=0

-- EACH_QUORUM при упавшем DC2 --
Unavailable exception ... message="Cannot achieve consistency level for cl EACH_QUORUM. Requires 2, alive 0"
info={'consistency': 'EACH_QUORUM', 'required_replicas': 2, 'alive_replicas': 0}
EACH_QUORUM write rc=2

LOCAL_QUORUM в DC1 продолжает работать без единой ошибки, пока весь DC2 недоступен, — координатор и обе реплики, от которых он ждёт подтверждения, всё это время находятся в живом датацентре, и потеря соседнего ДЦ для него попросту не видна. EACH_QUORUM — более строгий уровень, требующий кворум в каждом датацентре кластера, включая мёртвый, — отказывает мгновенно и с буквальной причиной от самого сервера: required_replicas: 2, alive_replicas: 0, не таймаут и не догадка клиента, а сервер сам называет, чего именно не хватает. Это ровно та развилка, ради которой EACH_QUORUM вообще существует: он дороже и строже QUORUM в мирное время и полностью недоступен, если хотя бы один датацентр кластера отвалился целиком, — сознательный компромисс за более сильную гарантию «подтверждено в каждом ДЦ», а не универсальный уровень для multi-DC развёртывания. После восстановления обоих узлов DC2 кластер синхронно возвращается к 4×UN — падение датацентра не оставило кластер в частично сломанном состоянии.

Tablets-ребаланс: decommission защищён репликацией

Последний шаг стенда демонстрирует то, ради чего tablets вообще стоило вводить как отдельную тему статьи, — живую миграцию реплик при выводе узла, без остановки кластера. Первая попытка сделать это буквально — nodetool decommission dc2b при keyspace, объявленном как {DC1:2, DC2:2} — заканчивается отказом:

$ docker exec dc2b nodetool decommission
error executing POST request ...: std::runtime_error (Decommission failed.
  See earlier errors (Unable to find new replica for tablet ... when draining
  {...}. Consider adding new nodes or reducing replication factor.

Причина честная и архитектурная, не ошибка стенда: если dc2b уйдёт из кластера, в DC2 останется один узел, а keyspace требует две реплики именно в этом датацентре — уместить их физически некуда. ScyllaDB сам отказывается выполнять decommission, который сломал бы заявленную репликацию, и кластер после отказа остаётся целым — 4×UN, никакого промежуточного, наполовину выведенного состояния. Реалистичный обходной путь — тот же, что применяют перед выводом предпоследнего узла датацентра в проде: сперва понизить RF затронутого ДЦ явным ALTER KEYSPACE:

ALTER KEYSPACE telemetry_mdc
  WITH replication = {'class':'NetworkTopologyStrategy','DC1':2,'DC2':1};

После этого nodetool decommission dc2b завершается успешно (rc=0). Саму миграцию демонстрирует отдельный keyspace {DC1:1, DC2:1} (RF=1 в каждом ДЦ с самого начала, чтобы демонстрация не требовала лишнего ALTER), таблица на 5000 детерминированных строк, tablet_count=32. Распределение реплик по узлам, прочитанное напрямую из system.tablets до и после decommission:

ДО decommission (tablet_count=32, 32 tablet-строки в system.tablets):
   dc1a  : 16 tablet-реплик
   dc1b  : 16 tablet-реплик
   dc2a  : 16 tablet-реплик
   dc2b  : 16 tablet-реплик
   ИТОГО: 64 (= 32 tablet-строки × RF 2)

nodetool decommission dc2b: rc=0

ПОСЛЕ decommission (тот же table_id, 3 живых узла):
   dc1a  : 16 tablet-реплик
   dc1b  : 16 tablet-реплик
   dc2a  : 32 tablet-реплики
   ИТОГО: 64 (= 32 tablet-строки × RF 2)

Все 16 tablet-реплик, ранее лежавших на dc2b, переехали на dc2a — единственный оставшийся узел DC2 вырос ровно с 16 до 32 реплик, забрав себе всё, что раньше делилось пополам между двумя узлами датацентра. dc1a и dc1b decommission не затронул вовсе — по 16 реплик до и после. Сумма реплик по кластеру — 64 и до, и после — ни одна не потерялась, ни одна не задублировалась, чистая миграция единицы распределения данных, которая произошла во время нормальной работы кластера, без остановки и без ручного перестроения кольца токенов.

flowchart TB subgraph BEFORE["До: DC2 = 2 узла, {DC1:2, DC2:2}, tablet_count=32"] direction LR subgraph DC1B["DC1"] A1["dc1a — 16 реплик"] A2["dc1b — 16 реплик"] end subgraph DC2B["DC2"] B1["dc2a — 16 реплик"] B2["dc2b — 16 реплик"] end end subgraph AFTER["После: ALTER KEYSPACE DC2:1, decommission dc2b (rc=0)"] direction LR subgraph DC1A["DC1"] C1["dc1a — 16 реплик"] C2["dc1b — 16 реплик"] end subgraph DC2A["DC2"] D1["dc2a — 32 реплики"] D2["dc2b — выведен из кластера"] end end BEFORE -->|"ALTER KEYSPACE ... DC2:1
nodetool decommission dc2b"| AFTER style BEFORE fill:#c9e4c5,stroke:#5b8a5e style AFTER fill:#f4d9c6,stroke:#c67a4a style D2 fill:#f4c6c6,stroke:#b4552f style D1 fill:#e8c9a0,stroke:#c67a4a

flowchart TB
  subgraph BEFORE["До: DC2 = 2 узла, {DC1:2, DC2:2}, tablet_count=32"]
    direction LR
    subgraph DC1B["DC1"]
      A1["dc1a — 16 реплик"]
      A2["dc1b — 16 реплик"]
    end
    subgraph DC2B["DC2"]
      B1["dc2a — 16 реплик"]
      B2["dc2b — 16 реплик"]
    end
  end

  subgraph AFTER["После: ALTER KEYSPACE DC2:1, decommission dc2b (rc=0)"]
    direction LR
    subgraph DC1A["DC1"]
      C1["dc1a — 16 реплик"]
      C2["dc1b — 16 реплик"]
    end
    subgraph DC2A["DC2"]
      D1["dc2a — 32 реплики"]
      D2["dc2b — выведен из кластера"]
    end
  end

  BEFORE -->|"ALTER KEYSPACE ... DC2:1
nodetool decommission dc2b"| AFTER style BEFORE fill:#c9e4c5,stroke:#5b8a5e style AFTER fill:#f4d9c6,stroke:#c67a4a style D2 fill:#f4c6c6,stroke:#b4552f style D1 fill:#e8c9a0,stroke:#c67a4a
Tablets-миграция при decommission dc2b: 16 реплик переезжают на dc2a, сумма реплик кластера не меняется

Где заканчивается multi-DC и начинается общий consistent hashing

Всё, что показано выше, — специфика ScyllaDB: tablets как единица распределения, конкретная арифметика LOCAL_QUORUM/QUORUM/EACH_QUORUM на двух датацентрах, конкретное защитное поведение decommission при недостаточной репликации. Общая механика консистентного хеширования — кольцо токенов, виртуальные узлы, почему добавление одного узла в наивной схеме перемешивает почти все ключи, а не только свою долю, — устроена одинаково у всего семейства хранилищ, использующих этот подход, и разобрана отдельно в статье про шардирование и consistent hashingСкоро: она объясняет фундамент, на котором классический token ring вообще стоит, а эта статья показывает, чем ScyllaDB от этого фундамента в итоге отошла.

Граница с драйверами тоже стоит того, чтобы назвать её явно. Здесь координатор запроса выбирался вручную и жёстко — DataCenterHostFilter, чтобы гарантированно измерить нужный ДЦ. В реальном приложении этот выбор делает клиентская библиотека сама, на основании политики маршрутизации, и от того, насколько она умеет учитывать не только датацентр, но и шард внутри узла, зависит отдельная часть производительности — следующая статья серии измеряет именно это: сколько реально даёт шард-осведомлённая маршрутизация в Go- и Java-драйверах поверх того же кластера.

Версии на стенде: образ scylladb/scylla:2026.2.0, CQL/release_version3.0.8, второй, отдельный от основного, multi-DC кластер (compose/multidc.yml) — 2 датацентра × 2 узла, NetworkTopologyStrategy, полностью одноразовый: по завершении демонстрации кластер удалён вместе с volume, а основной трёхузловой стенд серии остался нетронутым.

Обсуждение в Telegram

Присоединиться →

Комментарии