ScyllaDB часто описывают формулой «Cassandra на C++», но за этой формулой стоит конкретное архитектурное решение, а не просто смена языка реализации: каждый узел кластера — не один процесс с общей кучей и пулом потоков, а --smp N независимых шардов-реакторов, каждый из которых владеет своим ядром CPU, своей памятью и своим куском данных безраздельно. Это меняет не только то, как узел использует железо, но и то, как выглядит хвост латентности под нагрузкой — без знакомой по JVM-мирам паузы «стоп всему миру». Ниже — как это устроено и что показывают живые измерения на трёхузловом кластере.
Это вторая статья серии «ScyllaDB: глубокое погружение», продолжение статьи о моделировании данных. Дальше в серии — compaction, tombstones и repair, consistency levels и LWT, от token ring к tablets и multi-DC, shard-aware драйверы на Go и Java и эксплуатация и итоговый выбор.
Числа ниже — с того же живого трёхузлового кластера scylladb/scylla:2026.2.0, что и в первой статье, на уже загруженной таблице readings (672 000 строк). Стенд — в публичном репозитории примеров (digital-cookbook, каталог scylladb/architecture).
В статье
- Shard-per-core и Seastar
- Нагрузка растекается по шардам
- Латентность без GC-пауз
- Честная оговорка про измерения
Shard-per-core и Seastar
Классическая многопоточная СУБД на JVM (в том числе Cassandra) устроена вокруг общего пула потоков и общей кучи: любой поток может обратиться к любым данным, доступ к разделяемым структурам защищён блокировками, а память для всех потоков собирает один сборщик мусора. ScyllaDB, построенная на асинхронном C++-фреймворке Seastar, устроена принципиально иначе — shard-per-core: один shard = один поток ОС = одно ядро CPU, и каждый shard владеет своим куском данных, своей памятью и своей очередью задач безраздельно.
Практическое следствие — никакой блокировки между шардами на горячем пути. Если два запроса попадают на два разных ядра одного узла, они не конкурируют ни за одну структуру данных: каждый shard обрабатывает свою долю трафика независимо, взаимодействие между шардами (когда оно всё же нужно) идёт через явную message-passing очередь, а не через разделяемую память под мьютексом. Ввод-вывод — асинхронный: shard не блокируется на диске или сети, а переключается на следующую задачу в очереди, пока операция не завершится.
Это напрямую определяет путь записи и путь чтения. Запись сначала попадает в memtable (структура в памяти shard’а, отсортированная по ключу кластеризации) и одновременно — в commitlog на диске для устойчивости к падению процесса до сброса memtable; когда memtable заполняется, она атомарно сбрасывается на диск в виде неизменяемого sstable, а фоновая compaction впоследствии сливает несколько sstable в один, вычищая устаревшие и удалённые данные. Сам механизм LSM-дерева (запись только последовательным добавлением, чтение — слияние нескольких уровней) не специфичен для ScyllaDB — устройство памяти и диска в такой схеме подробно разобрано в статье про LSM-tree и B-treeСкоро. Здесь важнее другое: у каждого shard — своя memtable, свой commitlog, свои sstable этой части данных, поэтому запись и compaction на одном ядре не создают contention на другом.
Путь чтения зеркален: клиент подключается к любому узлу-координатору, но координатор не читает данные сам — он вычисляет, какой shard (по токену партиции) физически владеет нужным диапазоном, и маршрутизирует запрос именно туда. Чтение внутри shard’а идёт через кеш строк и memtable, при промахе — сквозь merge нескольких sstable, но всё это — операции одного ядра над своими структурами, без блокировки соседних шардов.
(device_id, day, event_time)"] --> R["Координатор узла:
вычисляет владеющий shard по токену"] subgraph NODE["Один узел ScyllaDB (--smp 2)"] direction LR subgraph S0["Shard 0 = ядро 0"] direction TB M0["memtable"] L0["commitlog"] T0["sstable"] M0 --> T0 end subgraph S1["Shard 1 = ядро 1"] direction TB M1["memtable"] L1["commitlog"] T1["sstable"] M1 --> T1 end end R -->|"токен → shard 0"| S0 R -->|"токен → shard 1"| S1 style S0 fill:#c9e4c5,stroke:#5b8a5e style S1 fill:#f4d9c6,stroke:#c67a4a
flowchart TB
C["Клиентский запрос
(device_id, day, event_time)"] --> R["Координатор узла:
вычисляет владеющий shard по токену"]
subgraph NODE["Один узел ScyllaDB (--smp 2)"]
direction LR
subgraph S0["Shard 0 = ядро 0"]
direction TB
M0["memtable"]
L0["commitlog"]
T0["sstable"]
M0 --> T0
end
subgraph S1["Shard 1 = ядро 1"]
direction TB
M1["memtable"]
L1["commitlog"]
T1["sstable"]
M1 --> T1
end
end
R -->|"токен → shard 0"| S0
R -->|"токен → shard 1"| S1
style S0 fill:#c9e4c5,stroke:#5b8a5e
style S1 fill:#f4d9c6,stroke:#c67a4a
Нагрузка растекается по шардам
Утверждение «shard-per-core» легко произнести, но интереснее проверить его на реальном трафике: если нагрузка действительно распределяется по ядрам, а не оседает на одном, это должно быть видно в счётчиках самого сервера. Сценарий -scenario shard-distribution (architecture/) снимает per-shard счётчик scylla_database_total_reads (монотонный counter с :9180/metrics) до и после 20 000 точечных чтений по случайным ключам уже загруженной readings, на всех трёх узлах кластера (--smp 2 — по два шарда на узел).
Дельты (after − before) по шардам:
| Узел | Shard 0 (Δ) | Shard 1 (Δ) | Сумма по узлу |
|---|---|---|---|
| scylla1 | 7128 | 6855 | 13 983 |
| scylla2 | 6834 | 6042 | 12 876 |
| scylla3 | 6831 | 7467 | 14 298 |
На каждом узле оба шарда получили сопоставимое число чтений — соотношение близко к 50/50, ни один узел не свалил всю нагрузку на одно ядро. Это и есть практическое подтверждение shard-per-core: 20 000 запросов, приходящих на кластер, физически разошлись по шести независимым ядрам-реакторам, а не встали в очередь на одно.
Здесь стоит остановиться на цифре, которая на первый взгляд не сходится: сумма дельт по каждому узлу (~12 900–14 600) заметно больше, чем n/3 (20000/3 ≈ 6667), хотя запросов всего 20 000 на три узла. Разгадка не в ошибке измерения, а в consistency level: таблица читается с Consistency: Quorum при RF=3 — каждое чтение должно получить ответ от не менее двух из трёх реплик, значит каждый логический запрос физически исполняется минимум на двух узлах одновременно. 20 000 запросов при RF=3 и QUORUM дают суммарно около 40 000 физических обращений к шардам по кластеру — отсюда и сумма на узел заметно выше «наивных» n/3.
Латентность без GC-пауз
Отсутствие общей кучи и общего сборщика мусора — не абстрактное архитектурное преимущество, а конкретная, измеримая разница в хвосте латентности. Сценарий -scenario latency прогоняет 50 000 последовательных точечных чтений через один коннекшн (без параллелизма — так измерение отражает latency отдельного round-trip, а не эффект очереди на клиенте) и печатает клиентские тайминги:
Успешных чтений: 50000/50000 (ошибок: 0)
p50: 1.5468ms
p99: 1.7115ms
max: 4.283ms
p99/p50 ≈ 1.11 <= K=5Ключевое здесь — не абсолютные цифры (они зависят от железа и топологии стенда), а форма хвоста: p99/p50 ≈ 1.11, max — единицы миллисекунд, ни одного выброса в сотни миллисекунд за все 50 000 запросов. Прогон воспроизведён трижды подряд, соотношение p99/p50 стабильно держится в диапазоне 1.09–1.13 — ни разу не подскочило.
Серверная сторона (nodetool proxyhistograms на scylla1, тот же момент) показывает тот же порядок соотношения: p50 = 179µs, p99 = 253.75µs (p99/p50 ≈ 1.42). Разница между клиентскими и серверными абсолютными числами — около 1,3 мс — это сетевой оверхед моста Docker Desktop между контейнером-клиентом и кластером, а не задержка внутри ScyllaDB: сравнивать стоит форму распределения (оба соотношения — единицы, не десятки), а не вычитать одно число из другого как «накладные расходы сервера».
Стенд ассертит p99 <= K·p50 с K=5 — это осознанный запас, а не подгонка под фактические 1,09–1,13: на кластере с фоновой compaction или repair хвост может уехать заметно дальше идле-цифры, и K=5 закладывает это. Даже с таким запасом инвариант проходит с большим отрывом — реальное соотношение более чем в 4 раза строже допустимого.
Здесь и раскрывается связь с архитектурой из первого раздела. В JVM-реализациях СУБД возможен особый класс пауз: полный сборщик мусора останавливает все потоки приложения — stop-the-world, во время которой узел не отвечает никому, и хвост латентности может «выстрелить» на порядки (в литературе по GC приводят p99 в сотни миллисекунд при p50 в единицы, то есть p99/p50 порядка десятков-сотен крат). Насколько это проявится на практике — сильно зависит от версии и настройки сборщика, профиля нагрузки и конфигурации: современная JVM/GC-настройка эту картину заметно меняет, и «периодически всё останавливается» — не приговор, а один из возможных режимов. Здесь важна честная рамка: 1,09–1,13 — реальные числа этого стенда, а десятки-сотни крат для JVM-style GC приведены как иллюстративный класс из общей практики, а НЕ как результат замера Cassandra под тем же профилем — такого замера в этой серии нет, и число не стоит читать как «мы измерили Cassandra». Что архитектура ScyllaDB действительно устраняет — это именно глобальную GC-паузу: на C++/Seastar нет общей кучи, которую пришлось бы останавливать целиком, каждый shard управляет своей памятью сам. Это не значит, что хвост латентности вообще неоткуда взяться — он по-прежнему может расти от I/O, compaction, repair, reactor stalls, шума планировщика и хоста; shard-per-core убирает один конкретный, но крупный источник stop-the-world задержки, а не все причины хвоста разом.
не серверная задержка"| SERVER style CLIENT fill:#f4d9c6,stroke:#c67a4a style SERVER fill:#c9e4c5,stroke:#5b8a5e
flowchart LR
subgraph CLIENT["Клиент (n=50000, один коннекшн)"]
direction TB
CP50["p50: 1.5468ms"]
CP99["p99: 1.7115ms"]
CMAX["max: 4.283ms"]
end
subgraph SERVER["Сервер (nodetool proxyhistograms)"]
direction TB
SP50["p50: 179µs"]
SP99["p99: 253.75µs"]
end
CLIENT -->|"~1.3ms разница = Docker-мост,
не серверная задержка"| SERVER
style CLIENT fill:#f4d9c6,stroke:#c67a4a
style SERVER fill:#c9e4c5,stroke:#5b8a5e
Честная оговорка про измерения
Два момента, которые важно проговорить прямо, а не спрятать в сноску.
Во-первых, выбор метрики для замера shard-distribution — не случайность. У ScyllaDB есть и scylla_reactor_utilization — gauge, мгновенная доля занятости CPU реактора в момент снятия — и scylla_database_total_reads — монотонный counter числа обработанных чтений. Gauge шумит от фоновой активности кластера между снапшотами «до» и «после» и плохо подходит для чистого before/after-сравнения; counter даёт однозначную дельту, которая не зависит от момента снятия. Поэтому выбран именно total_reads, а не reactor_utilization.
Во-вторых, сумма дельт по узлу, превышающая n/3, — не ошибка измерения и не баг стенда, а прямое следствие RF=3 + Consistency: Quorum: каждый логический запрос клиента физически исполняется минимум на двух из трёх узлов, поэтому сумма по одному узлу и не обязана совпадать с «наивной» третью от общего числа запросов. Понимание того, зачем нужна репликация и как она физически считается по узлам, дальше пригодится в статье про consistency levels и LWT — там та же арифметика применяется уже к операциям записи и compare-and-swap.
Версии на стенде: образ scylladb/scylla:2026.2.0, CQL/release_version — 3.0.8, --smp 2 на узел, датасет — та же readings из первой статьи серии (672 000 строк). Дальше в серии — что происходит внутри shard'а, когда данные удаляют и когда sstable сливаются.
Комментарии