Масштабируемость

Интуитивно каждый представляет, что означает «масштабироваться». Это значит делать больше. Это значит удовлетворять растущий спрос. Это значит обрабатывать миллиард запросов, или «бесконечное» количество запросов, или бесконечный объём данных, или бесконечное число клиентов. Это способность растить приложение в 10x или 100x раз год за годом без переписывания кода.

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

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

Вопрос «Масштабируется ли это горизонтально?» недостаточно точен. Лучше спросить: «В каких отношениях оно масштабируется горизонтально?» Дело в том, что существует три довольно независимых измерения масштабируемости:

  1. Вычисления (compute): сколько транзакций можно обработать
  2. Хранилище (storage): сколько данных можно сохранить
  3. Сеть (networking): сколько соединений и какую полосу пропускания можно поддерживать

Примечание: кстати, это ровно те 3 «фундаментальных компонента», на которые мы ссылались в keynote объявления версии 1.0.

Чтобы увидеть, о чём речь, рассмотрим несколько систем баз данных, которые все предоставляют примерно PostgreSQL-совместимые интерфейсы, но имеют кардинально разные архитектуры.

Например:

  • Postgres в большей части своей — база данных одного узла. Postgres не масштабирует вычисления, хранилище или сеть по горизонтали. Можно, конечно, масштабировать Postgres по горизонтали, развернув множество экземпляров Postgres, но с точки зрения кода Postgres он в основном не знает об этих других экземплярах. Единственное исключение — реплики для чтения, которые позволяют вручную направлять читателей на реплику. Это помогает частично масштабировать сеть и вычисления, но с оговорками о консистентности чтения после записи и производительности. Сам Postgres не имеет понятия кластера первичных узлов, не может запускать транзакции или запросы через них или маршрутизировать вас на нужный узел. Можно, конечно, написать ПО для этого, и именно так люди масштабируют Postgres по горизонтали, но как это делать — остаётся упражнением для читателя.

Анимированная архитектура Postgres: писатели встают в очередь позади одного первичного узла, который держит вычисления и хранилище вместе, WAL выходит на вручную настроенные реплики для чтения, читатели распределяются по ним; вычисления и сеть частично масштабируются только для чтений, хранилище не масштабируется

  • Neon — модифицированный вариант Postgres, масштабирует хранилище по горизонтали. Таблицы Neon поддерживаются объектным хранилищем со структурой локального кеша страниц, что делает доступ к данным быстрым и эффективным. Задержка доступа к странице может быть выше при промахах кеша, но эта архитектура даёт базам данных Neon практически бесконечное хранилище. Neon не автоматически масштабирует вычисления и сеть по горизонтали. Как Postgres, все транзакции записи должны идти через один первичный узел. Однако, как вариант Postgres, Neon тоже поддерживает реплики для чтения, чтобы горизонтально масштабировать вычисления и сеть для рабочих нагрузок чтения, хотя это сопровождается аналогичными оговорками о консистентности.

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

Анимированная архитектура Neon: писатели встают в очередь позади одного Postgres compute, WAL течёт вниз через pageservers в бездонное объектное хранилище, которое видимо растёт, страницы подтягиваются по требованию, реплика для чтения обслуживает читателей; хранилище масштабируется, а вычисления и сеть частично масштабируются только для чтений

  • CockroachDB (и Spanner) в принципе масштабирует хранилище по горизонтали, распределяя данные каждой таблицы по кластеру так, что часть каждой таблицы, называемая «range» (диапазон), хранится на каждой машине и обычно реплицируется на два других машины. Если писатели не пытаются одновременно изменять один и тот же диапазон, это также может масштабировать сеть по горизонтали. Он имеет симметричную архитектуру, то есть любой узел кластера может обслуживать любой SQL запрос, как чтения, так и записи. Наконец, для вычислений, которые хорошо параллелизируются (например, аналитика или записи в несвязанные ключи), это также масштабируется по горизонтали. Узел, к которому вы подключаетесь, вычисляет план запроса и возвращает результаты, но и операции чтения, и записи выполняются как часть распределённой транзакции по кластеру на основе диапазонов.

Анимированная архитектура CockroachDB: три симметричных узла каждый хранит реплицированные диапазоны, записи попадают на любой узел, каждая запись реплицируется на кворум других узлов и собирает подтверждения перед коммитом; вычисления, хранилище и сеть все масштабируются, при условии что транзакции не конкурируют

Можно ли найти святой Грааль?

Так что, если CockroachDB может масштабироваться во всех трёх измерениях, он должен быть лучше Postgres и Neon во всех аспектах, верно? Может быть, он делает что-то особенное с теоремой CAP[1] или атомными часами?

К сожалению, здесь нет магии. Хотя CockroachDB — современное чудо и масштабирует определённые вычисления по горизонтали, не все вычисления могут быть масштабированы по горизонтали. Горизонтальная масштабируемость — это действительно вопрос параллельных вычислений: можно ли разбить вычисление на части, которые разные компьютеры могут выполнять одновременно? Ответ часто «нет», и проблема в том, что с CockroachDB вы платите огромную стоимость координации горизонтальной масштабируемости даже когда конкуренция за данные вынуждает вас выполнять обновления одно за другим. То, что CockroachDB выигрывает в горизонтальной масштабируемости, он теряет в вертикальной масштабируемости и даже больше того.

Неприятная правда в том, что вместо того чтобы делать больше с большим числом компьютеров, горизонтальная масштабируемость часто может означать делать меньше с большим числом компьютеров: концептуально то, что один компьютер может сделать за 1 миллисекунду, 10 компьютеров могут сделать за 100 миллисекунд.

Есть две основные проблемы с подходом CockroachDB к горизонтальной масштабируемости:

  1. Данные, вовлечённые в каждую транзакцию, редко находятся в одном месте на одной машине.
  2. Ни одна транзакция не имеет исключительного доступа к диапазону, поэтому каждая транзакция платит стоимость распределённого управления одновременностью, и конкурирующие транзакции должны ждать, прерваться или повториться.

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

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

Примечание: Spanner частично решает эту проблему с помощью «table-interleaving», которое позволяет вам в сущности сказать Spanner совмещать связанные таблицы. Это может драматически улучшить производительность в простых случаях. CockroachDB поддерживал table interleaving несколько лет, но удалил его в v21.2, решив, что преимущества слишком малы, чтобы оправдать сложность.

Вторая проблема — это общая проблема параллелизации произвольных вычислений: координация под конкуренцией. Даже ванильный Postgres сталкивается с той же проблемой, просто в меньшем масштабе. Вместо координации транзакций в распределённой системе, ему нужно координировать транзакции между несколькими ядрами. Postgres может запускать транзакции на многих ядрах CPU, но как только эти транзакции касаются одних и тех же данных, им нужно время на координацию. Сам CPU должен координировать доступ к сырой памяти через кеш L1, L2 и L3. Другой писатель, изменяющий ту же строку кеша, может инвалидировать вашу локальную копию и заставить ядра синхронизироваться. Postgres должен координировать сами транзакции: кто владеет блокировкой, какие версии строк видны, в каком порядке коммитятся транзакции и нужно ли конкурирующей работе ждать, прерваться или повториться.

Анимированная диаграмма двух ядер, пишущих в ту же строку кеша, масштабировать во времени 0.2 секунды за наносекунду: ядро 1 владеет строкой в состоянии Modified и пишет каждую наносекунду; запрос чтения-для-владения ядра 2 пропускает L1 и L2 и посылается вниз в общий L3, чей директорий подслушивает ядро 1; копия ядра 1 инвалидируется в I состояние и изменённые данные двигаются вниз через его L2, через L3, и вверх в кеши ядра 2, прибывая в состояние Modified после примерно 100 наносекунд, во время которых оба ядра останавливаются, и затем ядро 2 пишет на полной скорости

Например, представьте маленькую OLTP транзакцию, которая включает строку пользователя и связанную строку метаданных. Во-первых, тогда как в однопроцессном Postgres обе строки будут на одной машине, в кластере CockroachDB с N узлами и случайными первичными ключами вероятность того, что они находятся на одном лидере, примерно 1/N. То, что могло быть ~3 микросекундной критической секцией на одном ядре, теперь 1 миллисекундная распределённая критическая секция с несколькими сетевыми запросами. Даже в случае 1/N, когда строки совмещены, транзакция всё равно не может освободить свои блокировки до тех пор, пока запись не будет реплицирована на кворум, поэтому она держит их примерно столько же, сколько круговой обход. Во-вторых, и что более важно, другая транзакция, которая хочет прочитать или изменить те же строки, теперь должна ждать 1 миллисекунду для нашей транзакции коммититься или прерваться. Эти две проблемы мультипликативно усиливают друг друга. Горячий ключ с ~3 µs критической секцией пропускает ~300,000 конкурирующих транзакций в секунду; при 1 ms это ~1,000 tps. Контринтуитивно, однопоточное решение в этом сценарии 300x более «масштабируемо»!

Примечание: Заметьте, что при 1 ms распределённых коммитах даже просто 1% транзакций, конкурирующих на одной строке, делает каждый горизонтально масштабируемый кластер медленнее, чем один ядро. Это просто закон Амдала, применённый к горизонтальной масштабируемости. Конкурирующие транзакции должны выполняться серийно, и общая пропускная способность никогда не может превышать серийную скорость разделённую на долю серийных транзакций: (1 / 1 ms) / 1% = 100,000 TPS. Неважно сколько ядер вы добавите! По сущности есть только один сценарий, в котором кластер обгоняет одно ядро по пропускной способности: у вас нет конкуренции И вы платите более $3,600 в месяц (основано на предположениях модели о оверхеде). Это сценарий, в котором находятся некоторые компании, но это нишевый сценарий.

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

Но, даже для вычислений, которые могут выполняться параллельно, люди часто драматически недооценивают, сколько дополнительного оборудования может быть необходимо просто чтобы преодолеть стоимость координации. Как грубая иллюстрация, задержка L1 кеша может быть около 20x ниже, чем L3. Если параллелизация рабочей нагрузки превращает дешёвые, локальные обращения к кешу в частую синхронизацию и кросс-ядерную коммуникацию, совершенно правдоподобно, что вам может понадобиться 10+ ядер просто чтобы сопоставить производительность одного аккуратного, оптимизированного под кеш ядра. И это даже не учитывает гораздо более высокую стоимость обращения в основную память, когда больше рабочих наборов и метаданные, такие как bookkeeping MVCC, выталкивают полезные данные из кеша. И мы всё ещё не рассмотрели сетевые запросы и сериализацию, требуемые для чего-то вроде распределённого MVCC bookkeeping и всех промахов кеша, которые они вызовут. Вы уверены, что хотите платить за 10 ядер, когда 1 ядро справится? Может быть, для подмножества проблем мы в конце концов достигаем масштабируемости, но какой ценой?

Для по сути серийных вычислений, вы можете переместить это на более быструю машину, оптимизировать код или реплицировать результаты для отказоустойчивости, но вы не можете сделать десять машин работающих быстрее, неважно как вам бы этого хотелось. Каждый программист, который читал The Mythical Man-Month, интуитивно знает это. Добавление программистов к проекту не означает писать ПО быстрее; это означает писать его медленнее, но с большим числом встреч. 10 авторов не могут написать роман быстрее, чем 1, особенно если они координируют через почту.

Мораль рассказа такова: горизонтальная масштабируемость — это не просто свойство, которое система имеет или не имеет. Это прежде всего свойство рабочей нагрузки.

Spacetime

Давайте поговорим о Spacetime.

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

Есть много способов реализовать:

ctx.db.myTable.insert({ name: 'Tyler' });

или:

SELECT * FROM my_table

API ничего не говорит о том, где живут эти данные, какая машина выполняет транзакцию или сколько машин задействовано за кулисами. Например, Convex построил слой транзакций поверх существующих хранилищ, таких как MySQL или Postgres. В принципе, мы могли бы реализовать Spacetime с флотом серверов модулей Spacetime и огромным кластером CockroachDB как хранилищем. Может быть, вас удивит узнать, что это довольно близко к тому, как мы начали. Первый прототип Spacetime был построен на Postgres + Kafka!

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

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

Может быть, удивительно, но сегодня модель выполнения для баз данных Spacetime — однопоточная по дизайну. Интуитивно, особенно для инженеров с ограниченным опытом оптимизации кеша, это может звучать как «плохая новость». Изначально, мы тоже начали с этого предположения: ранние версии нашего кастомного движка БД использовали параллельное выполнение, включённое MVCC транзакциями. Проблема в том, что когда мы действительно измерили, мы обнаружили, что однопоточное выполнение просто перепревосходит параллельное выполнение. В смысле, это, вероятно, стоило нам более $1M просто выяснить, что мы должны просто использовать большую блокировку. Не невозможно, что мы могли бы реализовать параллельное выполнение в рамках одной БД в будущем, которое достигает лучшей производительности в ограниченных ситуациях, но только без чрезвычайно осторожной инженерии и измерения производительности, чтобы убедиться, что оно не вводит больше сложности и оверхеда, чем оно стоит.

Кривая Гаусса мем: 'Просто используй один ядро'

Так это означает, что Spacetime не может масштабироваться по горизонтали? Отнюдь нет.

Горизонтальное масштабирование Spacetime

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

Отрывок из Hewitt и Baker, мотивирующий модель акторов высокопараллельными машинами многих независимых монопроцессоров

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

Следующий этап уже имеет дату: асинхронная межбазовая коммуникация и многоуровневое хранилище выходят 31 октября 2026 года как часть нашего запуска масштабируемости: Spacetime Continuum.

Шаг Возможность Статус
1 Обеспечить реплицированные БД с мировым уровнем производительности Сегодня
2 Обеспечить опыт шардирования мирового уровня (async IDC) 31 октября 2026
3 Масштабировать хранилище каждой БД по горизонтали (многоуровневое хранилище) 31 октября 2026
4 Масштабировать сеть каждой БД по горизонтали (реплики для чтения) Планируется
5 Реализовать кросс-БД транзакции (sync IDC) Планируется
6 Реализовать внутри-БД секционирование Планируется

Перед деталями, есть две версии Spacetime, которые стоит различать:

  • SpacetimeDB Standalone, однопроцессная версия, доступная на GitHub.
  • SpacetimeDB Cloud, собственническая, кластеризованная версия.

Реплицированные базы данных с мировым уровнем производительности

Мы уже подробно говорили о том, как каждая БД Spacetime достигает мирового уровня, однопоточную производительность, но важно заметить, что SpacetimeDB Cloud также поддерживает распределённую репликацию конечного автомата для БД.

Это означает, что хотя каждая БД однопоточна, она необязательно работает только на одной машине. Это звучит контринтуитивно, но общая идея в том, что мы можем реплицировать то, что делает один поток на несколько машин так, что если машины отказывают, БД остаётся доступной. Это означает, что каждая БД в SpacetimeDB Cloud работает как распределённая система.

Анимированная диаграмма Viewstamped Replication: клиенты посылают транзакции однопоточному первичному на узле 1, который выстраивает в трубу PREPARE сообщения резервным на узлах 2 и 3 и собирает PREPARE_OK подтверждения; когда узел 1 отказывает, резервные выполняют изменение представления, обмениваясь STARTVIEWCHANGE и DOVIEWCHANGE сообщениями, узел 2 становится первичным представления 2 и посылает STARTVIEW, и клиенты маршрутизируют к нему

Важно заметить, что с нашей выстроенной в трубу реализацией, мы продемонстрировали в нашем тестировании, что реплицированные БД достигают той же пропускной способности, что и нереплицированные БД (примерно 300k TPS для транзакций тестирования), при условии достаточной пропускной способности сети между узлами и памяти для глубины трубы.

Анимированная диаграмма выстроенной в трубу репликации: журнал транзакций лидера течёт через четыре водометки, current, applied, decided, и durable; лидер применяет новые запросы немедленно в начале, лидер выстраивает PREPARE сообщения в трубу для двух резервных, первое PREPARE_OK подтверждение обратно продвигает водометку decided, первое DURABLE_OK продвигает водометку durable, ответы выходят клиентам, прослушивающим уровни applied, decided, и durable, и блокировка БД держится только для транзакции, выполняемой в начале

Заметьте, что confirmedReads(true), настройка по умолчанию для Spacetime, конфигурирует клиентов читать только после того, как транзакции достигнут дurable статуса на кластере.

Опыт шардирования мирового уровня

Шардирование вашей БД даёт вам лучшую производительность для транзакций, которые работают исключительно в одном шарде. Распределённые транзакции дают вам способность запускать транзакции, которые охватывают шарды. Почему мы не можем иметь лучшее из обоих миров? Что если бы большинство транзакций работали в одной машине и мы дали бы программисту инструменты организовать их данные так, чтобы большинство транзакций происходили в этой машине? Только в редких случаях нам нужна была бы распределённая транзакция.

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

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

Spacetime уже предоставляет инструменты для управления этой архитектурой. БД можно создавать как потомки других БД. Процедуры могут публиковать БД и вызывать функции на них. Корневая БД может отслеживать идентичности и состояние БД под ней.

Несколько наших клиентов работают сотни или тысячи БД таким способом и BitCraft также масштабируется таким образом. Одна корневая БД поддерживает глобальные данные, тогда как региональные БД обрабатывают разные части мира и разные группы игроков. Транзакции в одном регионе остаются быстрыми и локальными, тогда как регионы сами выполняются параллельно.

Мы планируем сделать эту модель значительно проще через первоклассную межбазовую коммуникацию (IDC). Асинхронная IDC позволит одной БД вызывать функции на другой БД типобезопасным способом.

Анимированная диаграмма кластера Spacetime с шестью узлами, каждый хостит несколько БД; асинхронные межбазовые сообщения прыгают между БД через узлы, и каждая БД продолжает выполняться независимо

На высоком уровне, TypeScript API будут выглядеть примерно так:

// Async IDC: send a type-safe message to another database.
// This will result in exactly one reducer call on the target DB.
ctx.db.receivePlayer.insert({
    msgId: 0n,
    target: regionDb,
    player,
});

Масштабировать хранилище каждой БД по горизонтали

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

Однако это ограничение не фундаментально для Spacetime или ключ к его невероятной производительности. Мы рассматриваем память, диск и объектное хранилище как естественное расширение модели кеша CPU. Духовно мы рассматриваем память как L4 кеш, диск как L5 кеш и объектное хранилище как L6 кеш. Эта модель кеширования — время-честный способ получить максимальную производительность от одного писателя. Это часто называют «многоуровневым хранилищем».

Пирамида многоуровневого хранилища: иерархия кеша CPU L1 через L3, расширенная Spacetime с памятью как L4, NVMe диском как L5 и объектным хранилищем как L6; каждый уровень вниз больше, дешевле и медленнее, с памятью доступной сегодня и диском и объектными уровнями выходящими октябре 2026

Строки кеша и предвыборка позволяют вам упреждающе вытягивать партии данных из более дешёвого, более медленного хранилища в более дорогое, более быстрое хранилище. Предоставление многоуровневого хранилища БД позволяет вам амортизировать стоимость промахов кеша и достигнуть производительности, которая в целом близка к более быстрому, более дорогому хранилищу, одновременно сохраняя масштабируемость более большого, более дешёвого хранилища.

Не будет ли разбиение на страницы блокировать однопоток?

Если каждая транзакция выполняется на одном потоке, может показаться, что один промах кеша к объектному хранилищу остановит всю БД на 50 миллисекунд. Это было бы как работать библиотекой и каждый раз, когда клиент кладёт книгу на заказ, вы делаете всю линию клиентов стоять там и ждать две недели, пока книга придёт, перед тем как обработать следующего клиента. Очевидно глупо. Правило, которое делает многоуровневое хранилище совместимым с однопоточным выполнением просто: не ждите промаха кеша, держа блокировку. Закажите книги асинхронно и обработайте следующего клиента, пока вы ждёте!

К счастью, это хорошо изученная проблема. H-Store назвал технику anti-caching: выполните транзакцию нормально, и в момент, когда она касается данных, которые не резидентны, прервите её, получите данные асинхронно и запустите другие транзакции, пока вы ждёте. Когда данные приходят, запустите транзакцию снова. Прерывание почти бесплатно, потому что ничего не коммитилось, и каждое переисполнение стоит микросекунд CPU против миллисекунд I/O, которые оно избегает сериализации. Calvin использовал аналогичный трюк, который он назвал reconnaissance queries: сухой запуск транзакции против недавного снимка, чтобы открыть, что она читает, предвыборка этого, затем выполнить по-настоящему. TigerBeetle запускает явную фазу предвыборки перед синхронным применением каждой партии транзакций.

Анимированная диаграмма прерывания и перезапуска: транзакции текут через однопоточный исполнитель лидера, читая горячие страницы в памяти с быстрыми циклами туда-обратно; одна транзакция касается страницы, которая не в памяти, прерывается с промахом страницы и паркует в слоте ожидания внутри лидера, пока страница получается асинхронно из объектного хранилища через диск в память; исполнитель продолжает выполнять другие транзакции на протяжении всей выборки, и когда страница приходит припаркованная транзакция переисполняется и коммитится

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

Это стратегия, которая чрезвычайно способная команда в TigerBeetle называет «диагональное масштабирование».

Таблицы диска и объектного хранилища выходят 31 октября 2026. Мы рассматриваем это как критическое улучшение масштабируемости Spacetime. С табличками диска, мы можем значительно увеличить лимиты хранилища и уменьшить цены хранилища данных. С табличками объектного хранилища, мы можем практически удалить лимиты хранилища и позволить вам хранить теоретически неограниченное количество данных в одной БД.

Масштабировать сеть каждой БД по горизонтали

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

Как CockroachDB, SpacetimeDB Cloud также имеет симметричную архитектуру, что означает, что вы можете подключиться к любому узлу и SpacetimeDB Cloud обеспечит обработку вашего запроса записи или прокси его на нужное место. Снаружи это что-то вроде одного гигантского компьютера. Для писателей, это позволяет нам «вентилировать-в» подключения. Клиенты подключаются к любому узлу, и этот узел будет мультиплексировать те подключения в одно подключение к узлу, который хостит БД.

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

«Согласованная» здесь означает, что реплика применяет журнал транзакций лидера в том же полном порядке, поэтому подписка видит ровно ту же последовательность обновлений, просто отложенную. Для point-in-time SQL запросов, лидер назначит каждому запросу позицию в журнале транзакций без исполнения, и реплика не ответит на запрос до тех пор, пока она не применила через этот оффсет. Чтения поэтому линеаризуемы в отношении каждой записи, и стоимость лидера — только раздавание порядковых номеров.

Анимированная диаграмма вентилирования-в вентилирования-из: клиенты писателя подключаются к любому узлу, который мультиплексирует их транзакции к однопоточному лидеру; журнал транзакций лидера текёт к планируемым согласованным репликам для чтения, которые вентилируют обновления подписок читателям

Кросс-БД транзакции

Синхронная IDC предоставляет встроенную двухфазную коммит-машинерию, требуемую для выполнения настоящих распределённых транзакций вида, которые предоставляет CockroachDB. Это позволит транзакции охватывать несколько БД в кластере, когда требуется настоящая атомарность.

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

// Sync IDC: call another database inside the current transaction.
ctx.at(regionDb).reducers.receivePlayer(player);

Есть очевидное напряжение здесь. Если каждая БД — одна нить позади одной большой блокировки, затем синхронный вызов из БД A в БД B держит блокировку B до тех пор, пока транзакция A не коммитится или прерывается, что занимает круговой обход сети. Наш план — держать выполнение однопоточным, но переввести MVCC для этих распределённых транзакций, чтобы позволить одновременным транзакциям идти вперёд, пока круговой обход в полёте.

Анимированная диаграмма выстроенного в трубу двухфазного коммита: в раунде 1 клиентская транзакция на БД A вызывает БД B, получает PREPARED, и обе БД коммитятся в памяти и освобождают свои блокировки; в раунде 2 БД обмениваются PREPARED TO PERSIST и COMMIT PERSIST и реплицируют входы журнала своим резервным без держания любых блокировок, пока одновременные транзакции продолжают выполняться на протяжении

Мы разработали выстроенный в трубу протокол двухфазного коммита, который коммитит транзакции в памяти сначала и никогда не держит блокировку БД во время диск I/O, и мы модель-проверили его свойства безопасности в TLA+. Детали стоят отдельного поста.

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

Внутри-БД секционирование

Несколько БД — это естественная граница масштабирования, но они также требуют вам управления несколькими модулями, обновления их схем независимо и выяснения, как переуравновешивать нагрузку, пока она меняется. Следующий шаг для адресации этих UX проблем — введение секций в одной логической БД.

Секции позволят одному модулю БД содержать несколько независимо выполняющихся шардов. Это сохраняет один артефакт развёртывания и позволяет изменениям схемы применяться атомарно по всей БД, пока транзакции для разных секций выполняются одновременно.

Секция в Spacetime — это другая вещь от диапазона в CockroachDB. Диапазон — это решение размещения: он говорит, какая машина хранит часть таблицы, но любая транзакция все ещё может касаться любого диапазона, и каждая запись реплицируется и координируется одним и тем же способом неважно где она приземлится. Секция Spacetime — это единица выполнения. Она владеет своими данными, она имеет собственный серийный журнал транзакций, и код редьюсера, который работает на этих данных, работает внутри неё. Разработчики выбирают границы секции, чтобы совпасть со структурой их рабочей нагрузки, так что обычная транзакция никогда не оставляет свою секцию, и Spacetime может двигать целые секции между машинами, чтобы переуравновешивать нагрузку без ослабления того гарантия.

Анимированная диаграмма внутри-БД секционирования: одна логическая БД, построенная из одного модуля, охватывает две машины и содержит четыре секции, каждая со своим журналом транзакций; клиентские транзакции маршрутизируют к их секции по ключу и выполняются одновременно, редкая кросс-секционная транзакция запускает двухфазный коммит между двумя секциями, и переуравновешивание движет целой секцией от машины 1 к машине 2, пока всё продолжает работать

Транзакция, которая никогда не оставляет её секцию, не платит ничего за существование других. Это принцип нулевой стоимости: вы не должны платить за горизонтальную масштабируемость до тех пор, пока вы её используете.[2]

Поток решений для транзакции Spacetime: транзакции, которые остаются в одной БД или секции берут локальный быстрый путь, доступный сегодня; работа кросс-границы использует async IDC, выходящую октябре 2026, если она не должна быть синхронно атомарной, в таком случае планируемая распределённая транзакция платит стоимость координации

Так, масштабируется ли это?

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

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

Tyler Cloutier
Cofounder, Clockwork Labs


[1] Короткое отступление о теореме CAP для нердов распределённых систем. Теорема CAP, также известная как Теорема Брюэра, просто гласит, что никакая система в общем не может быть всеми тремя: согласованной, доступной и отказоустойчивой к разделению. Система может быть 0, 1 или 2 из этих вещей, но никогда 3.

«Согласованная» в этом случае означает, что каждое чтение получает самую свежую запись или ошибку. На практике, это свойство достигается правильной реализацией распределённой репликации конечного автомата (т.е. некоторый вариант Paxos/Viewstamped Replication).

«Доступная» означает, что каждый запрос к неотказавшему узлу получает ответ без ошибки.

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

Если я услышу ещё одного человека расплывчато ссылающегося на теорему CAP как «доказательство» того, что БД может или не может быть горизонтально масштабируемой, я просто могу сойти с ума. По какой-то причине есть распространённое неправильное понимание, что теорема CAP, также известная как Теорема Брюэра, каким-то образом ограничивает масштабируемость OLTP систем БД. Например, я был спрошен инвесторами и инженерами, может быть, половину дюжины раз, как мы «обошли теорему CAP», чтобы получить числа тестирования Spacetime. На самом деле, теорема CAP не делает никаких комментариев по этому поводу.

Моя теория — что это неправильное понимание появилось из этой публикации и этой статьи от Eric Brewer Google, оригинального формулятора конъектуры (теперь теорема: она была позже доказана в этой статье Seth Gilbert и Nancy Lynch), конкретно, потому что Spanner известен как знаменитая горизонтально масштабируемая SQL БД. Может быть, она попала в воздух, что, потому что эта статья ссылается на обе теорему CAP и Spanner, теорема CAP должна быть каким-то образом важна для масштабируемости. Однако, вы заметите, что оригинальная статья Spanner не упоминает теорему CAP вообще.

В любом случае, теорема CAP только релевантна для БД постольку, поскольку лучшее, что вы можете сделать — это сделать выбор между разработкой CP системы или AP системы. В CP системе вы всегда согласованы, но иногда недоступны из-за отказа сети, и в AP системе вы всегда доступны, но иногда несогласованны из-за отказа сети.

Так, потому что мы не можем сделать сети безошибочными, в смысле теорема CAP просто задаёт вопрос: вы хотите быть правильным или доступным? Теорема CAP говорит абсолютно ничего о пропускной способности, задержке или масштабируемости.

Независимо, на практике для OLTP БД, нет даже настоящего выбора. OLTP БД с сильными гарантиями согласованности, включая Spanner, должны быть CP системы, потому что несогласованность означает, что ваши пользователи видят неправильные результаты (устаревшие чтения, немонотонные чтения, конфликтующие записи, которые в редких случаях могли бы затопить весь ваш интернет-магазин).

Через эту линзу, вы теперь можете ясно видеть, что теорема CAP не интересна или релевантна для масштабируемости. Действительно интересный вопрос: как мы можем построить горизонтально масштабируемую, согласованную (т.е. правильную) OLTP БД систему?

Согласованность в CAP близко связана с Изоляцией в ACID (смешивающе, Согласованность в ACID описывает что-то совершенно другое). CAP согласованность — линеаризуемость: операции появляются происходить в одном порядке, который уважает реальное время. ACID изоляция, в её сильнейшем, — сериализуемость: транзакции появляются происходить в некотором серийном порядке. БД, которая даёт вам обе — «строго сериализуемая», что предлагает Spanner.

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

Теперь мы также можем понять, почему Spanner использует атомные часы. Атомные часы позволяют Spanner назначить линеаризуемый порядок транзакциям, которые не имеют взаимных зависимостей данных (т.е. они тривиально параллелизуемы) и могут быть выполняться на разных сторонах земли. Часы позволяют Spanner сказать, какая транзакция произошла сначала без необходимости посылать сообщения вокруг света, чтобы решить относительный порядок транзакций, которые иначе не нуждались бы в координации.

[2] Что по поводу запросов и транзакций, которые охватывают секции? Мы определённо могли бы в конце концов их поддержать. Распределённые SQL запросы через секции восстановили бы общую горизонтальную масштабируемость, предложенную системами такими как CockroachDB, при сохранении быстрого локального пути Spacetime для транзакций, которые остаются в одной секции.

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

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

Статическая зарисовка двух исследовательских форм для запросов через секции БД: опция A, scatter-gather SQL, вентилирует один клиентский запрос через шлюз вовне к каждой живой секции и объединяет результаты, поэтому чтения касаются пути записи; опция B течёт изменения секции в БД кластера аналитики только-чтения, которую клиенты запрашивают вместо, держа чтения с горячих ядер

Для этой, я думаю, лучше всего слушать-наших-клиентов подход и видеть, что они действительно нуждаются.