Масштабирование баз данных: Репликация и шардирование

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

Предел вертикального роста: когда база данных становится Bottleneck

Предел вертикального роста: когда база данных становится Bottleneck

В предыдущих модулях мы научились виртуозно масштабировать слой приложения. Благодаря легковесным горутинам, M:N планировщику и Stateless-архитектуре, горизонтальное масштабирование (Scale-Out) Go-сервисов стало тривиальной задачей. Уперлись в CPU? Собираем бинарник, упаковываем в Docker и поднимаем еще 50 подов в Kubernetes за пару секунд. Балансировщик нагрузки равномерно распределит трафик, и система снова задышит.

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

В этой статье мы разберем, почему база данных неизбежно становится главным Bottleneck всей системы, почему мы не можем просто бесконечно покупать для нее серверы мощнее (Scale-Up) и что именно ломается внутри БД под высокой нагрузкой.

Асимметрия распределенных систем

Архитектура большинства современных веб-проектов страдает от фундаментальной асимметрии масштабирования.

Слой приложения (App Tier) спроектирован как Stateless. Он не хранит состояние между запросами, поэтому узлы взаимозаменяемы. Слой данных (Data Tier) — это Stateful. База данных обязана хранить состояние, обеспечивать транзакционную целостность (ACID) и персистентность на диске.

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

Иллюзия спасения через Scale-Up

Первая инстинктивная реакция на перегрузку БД — применить Scale-Up (вертикальное масштабирование). Не хватает ресурсов? Давайте переедем на сервер побольше. Добавим CPU, докинем терабайт RAM, поставим самые быстрые NVMe-диски.

На ранних этапах жизни проекта это абсолютно правильная и экономически обоснованная стратегия. Переезд на более мощный инстанс занимает минуты и не требует изменения ни строчки кода. Но у Scale-Up есть три непреодолимых предела.

1. Физический предел

Закон Мура замедлился. Тактовая частота процессоров практически не растет последние 15 лет — производители наращивают количество ядер. Но базы данных не могут идеально распараллеливать все задачи. Как мы помним из анализа sync.RWMutex, конкурентный доступ к общим структурам данных (например, к индексам БД или таблицам блокировок) создает Contention. Начиная с определенного количества ядер, накладные расходы на синхронизацию потоков внутри самой СУБД начинают съедать весь прирост производительности.

2. Экономический предел

Стоимость вычислительных ресурсов при Scale-Up растет нелинейно. Удвоение мощности на малых объемах стоит копейки, но в топовом сегменте цена взлетает по экспоненте.

Тип инстанса (AWS RDS PostgreSQL) vCPU RAM Примерная цена в месяц Стоимость 1 ГБ RAM
db.m6g.large 2 8 ГБ ~120 USD 15 USD
db.m6g.4xlarge 16 64 ГБ ~950 USD 14.8 USD
db.m6g.16xlarge 64 256 ГБ ~3,800 USD 14.8 USD
db.x2g.16xlarge (Memory Optimized) 64 1024 ГБ ~12,500 USD 12.2 USD
Свой сервер на 4 ТБ RAM 128 4096 ГБ Сотни тысяч USD Эксклюзивное железо

В какой-то момент покупка одного суперкомпьютера становится экономически нецелесообразной по сравнению с покупкой десятка обычных серверов (Scale-Out).

3. Предел отказоустойчивости (SPOF)

Самый мощный сервер в мире все еще подключается к сети через кабель, который можно случайно выдернуть, и питается от блока, который может сгореть. Огромная монолитная база данных — это классический SPOF. Если она упадет, время восстановления (Recovery Time) базы размером в несколько терабайт из бэкапа может занять часы. Для High-Load систем такие простои недопустимы.

Анатомия узкого места: что ломается первым?

Когда база данных достигает предела насыщения (Saturation), деградация происходит стремительно. Давайте разберем три главных ресурса, исчерпание которых убивает БД.

Память и дисковый I/O (Проблема Cache Miss)

База данных работает быстро только тогда, когда данные находятся в оперативной памяти (Buffer Pool / Shared Buffers). Чтение из RAM занимает наносекунды, чтение с SSD — микросекунды (разница в сотни раз).

Пока объем «горячих» данных (Working Set) помещается в RAM, база летает. Но как только объем данных превышает размер памяти, СУБД начинает постоянно вытеснять страницы из памяти на диск и читать новые. Возникает массовый Cache Miss. Дисковая подсистема перегружается, и Latency запросов взлетает в космос.

Процессор (CPU)

В отличие от простых Key-Value хранилищ, реляционные БД выполняют сложную работу:

  • Парсинг и планирование SQL-запросов.
  • Сортировка данных (операции ORDER BY, требующие сложности O(NlogN)O(N \log N)).
  • Агрегации (GROUP BY) и сложные JOIN. Если приложение не использует Batching и страдает от проблемы N+1, оно заваливает базу тысячами мелких запросов, заставляя CPU базы данных непрерывно сжигать такты на парсинг одного и того же SQL.

Пул соединений (Connection Pool Exhaustion)

Это самое коварное узкое место, на котором спотыкаются разработчики на Go.

В Go мы привыкли, что создать 10 000 горутин — это бесплатно (каждая занимает ~2 КБ памяти). Но для классической СУБД вроде PostgreSQL каждое входящее соединение — это отдельный тяжеловесный процесс операционной системы (OS Process), который потребляет от 2 до 10 МБ памяти на свои внутренние структуры и кэши.

Если 50 инстансов Go-сервиса откроют по 100 соединений к БД, база получит 5000 активных процессов. Процессор СУБД уйдет в состояние Thrashing, тратя все время на Context Switch между тысячами процессов, а не на выполнение самих запросов.

Чтобы защитить БД, мы ограничиваем размер пула соединений (например, 200 коннектов на всю базу). Но здесь вступает в игру Закон Литтла: L=λ×WL = \lambda \times W. Если Concurrency (LL) жестко ограничено пулом в 200 соединений, а время обработки запроса базой (WW) из-за нагрузки выросло с 10 мс до 100 мс, то максимальная пропускная способность (λ\lambda) математически ограничивается: 200/0.1=2000200 / 0.1 = 2000 RPS.

Все остальные запросы от Go-приложения будут стоять в очереди на получение соединения из пула (database/sql будет блокировать горутины), пока не отвалятся по таймауту.

Смена парадигмы: от монолита к распределенным данным

Когда Scale-Up исчерпан, а база данных стала узким местом по CPU, памяти или соединениям, у архитектора остается только один путь — применить Scale-Out к слою данных. База данных должна перестать быть черным ящиком на одном сервере и превратиться в распределенную систему.

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

  1. Масштабирование чтения (Replication). В большинстве веб-приложений запросов на чтение (SELECT) на порядки больше, чем на запись (INSERT/UPDATE). Мы можем копировать данные на дополнительные серверы (Реплики), чтобы они обслуживали читающий трафик, разгружая основной узел.
  2. Масштабирование записи и объема (Sharding). Если один сервер больше не может справляться с потоком новых данных или общий объем превысил все разумные пределы дисковых массивов, данные необходимо физически разрезать на части (Шарды) и разнести по разным независимым серверам.

Переход к распределенным данным вернет нам теорему CAP, проблемы Partial Failure и сетевые задержки — все то, что мы обсуждали в пятом курсе, но теперь уже на уровне хранения состояния.

Репликация: масштабирование чтения и обеспечение высокой доступности

Репликация: масштабирование чтения и обеспечение высокой доступности

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

В социальных сетях на один написанный пост приходятся тысячи просмотров. В e-commerce на один оформленный заказ — сотни поисковых запросов и просмотров карточек товаров. Соотношение операций чтения к записи (Read-to-Write ratio) в High-Load системах редко бывает 1:1. Чаще это 10:110:1, 100:1100:1 или даже 1000:11000:1.

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

Паттерн Leader-Follower (Master-Slave)

Чтобы распределить чтение, мы применяем архитектурный паттерн Leader-Follower (исторически известный как Master-Slave).

Идея проста: мы разворачиваем несколько экземпляров базы данных. Один узел назначается главным (Leader). Только он имеет право принимать операции записи (INSERT, UPDATE, DELETE). Все остальные узлы становятся ведомыми (Followers). Их задача — хранить точную копию данных лидера и обслуживать запросы на чтение (SELECT).

Как данные попадают на Followers? Базы данных (например, PostgreSQL или MySQL) используют механизм WAL (Write-Ahead Log). Это бинарный журнал, в который СУБД последовательно записывает все изменения до того, как применить их к файлам данных на диске. Leader непрерывно транслирует этот поток WAL-записей по сети своим Followers. Ведомые узлы получают поток, проигрывают его у себя и таким образом повторяют состояние лидера.

Благодаря этому паттерну мы снимаем нагрузку с основного узла. Если раньше Leader обрабатывал 100% трафика, то теперь он занимается только записью (те самые 1-10% запросов). Чтение мы балансируем между Followers, используя алгоритмы вроде Round Robin или Least Connections, которые мы разбирали в пятом курсе. Если чтения становится больше — мы просто добавляем еще одного Follower в кластер (горизонтальное масштабирование, Scale-Out).

Физика копирования и Replication Lag

Звучит идеально, но данные не умеют телепортироваться. Передача WAL-журнала по сети и его применение на диске Follower-узла требует времени. Это время называется Replication Lag (задержка репликации).

В спокойном состоянии задержка может составлять миллисекунды (ограничиваясь RTT сети). Но при всплеске операций записи, деградации сети или тяжелых аналитических запросах на Follower-узле, задержка может вырасти до секунд или даже минут.

Наличие Replication Lag возвращает нас к теореме CAP. Разделяя данные между узлами (Partition), мы вынуждены выбирать между консистентностью (C) и доступностью (A).

Существует два режима репликации, отражающих этот компромисс:

Режим репликации Как работает Следствие по CAP
Синхронная (Synchronous) Leader ждет подтверждения от Follower, прежде чем ответить клиенту OK. CP. Нулевой Lag, данные никогда не потеряются. Но если Follower недоступен или сеть тормозит — Leader блокирует запись. Доступность падает.
Асинхронная (Asynchronous) Leader пишет данные себе и сразу отвечает клиенту OK. WAL отправляется в фоне. AP. Максимальная скорость и доступность записи. Но появляется Replication Lag.

В High-Load системах в 99% случаев используется асинхронная репликация, так как блокировка записи из-за сетевых морганий недопустима. Это означает, что мы осознанно переходим в парадигму Eventual Consistency (согласованность в конечном счете). Если мы прекратим писать данные, через какое-то время Δt\Delta t (равное Replication Lag) все Followers догонят лидера и состояние системы станет консистентным.

Но пока Δt\Delta t не прошло, система живет в состоянии рассинхрона.

Для решения проблемы "Read-Your-Own-Writes" (прочитай свою запись) на уровне приложения часто применяют эвристики. Например: после того как пользователь обновил профиль, приложение принудительно отправляет все его SELECT запросы на Leader-узел в течение следующих 5 секунд, а запросы остальных пользователей продолжают идти на Followers.

От масштабирования к отказоустойчивости (High Availability)

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

Если у нас есть только один сервер БД, его отказ (сгорел блок питания, упал процесс СУБД) означает полную остановку сервиса. Наличие Followers меняет картину: у нас появляются "горячие резервы" (Hot Standby).

Процесс переключения системы на резервный узел при падении основного называется Failover. Как это происходит:

  1. Система мониторинга замечает, что Leader не отвечает на Heartbeat-запросы.
  2. Один из Followers выбирается новым Лидером (Promotion).
  3. Приложение (или балансировщик, например PgBouncer) перенаправляет пул соединений для записи на новый IP-адрес.

Однако здесь кроется главная опасность асинхронной репликации. Если Leader упал, а Replication Lag в этот момент составлял 2 секунды, то последние 2 секунды транзакций остались только на мертвом диске Лидера. Новый Лидер о них не знает. При Failover в асинхронном кластере потеря небольшого окна данных неизбежна. Это осознанный архитектурный Trade-off между доступностью системы и сохранностью 100% данных в момент аварии.

Реализация на уровне Go-приложения

Как паттерн Leader-Follower меняет архитектуру нашего кода? В монолитной БД у нас был один пул соединений *sql.DB. Теперь инфраструктура данных стала гетерогенной.

На уровне приложения (или промежуточного слоя доступа к данным) мы обязаны разделить пулы. У нас появляется WritePool (указывающий на Leader) и ReadPool (балансирующий запросы между Followers).

Разработчик должен явно принимать решение:

  • Вызываем UpdateUserBalance() \rightarrow берем соединение из WritePool.
  • Вызываем GetProductReviews() \rightarrow берем соединение из ReadPool.

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

Мы успешно масштабировали чтение. Но что произойдет, когда наш сервис вырастет настолько, что сам Leader перестанет справляться с потоком записи? Добавлять новых Followers бессмысленно — они не принимают запись. Развернуть несколько Лидеров? В следующей главе мы разберем топологии репликации и узнаем, почему Multi-Leader архитектура превращает консистентность данных в настоящий кошмар.

Топологии репликации и проблема Consistency в распределенных данных

Топологии репликации и проблема Consistency в распределенных данных

Мы уже научились масштабировать чтение, направляя тяжелые SELECT запросы на ведомые узлы (Followers). Но эта идиллия рушится, как только мы задаем два вопроса. Первый: что произойдет, если единственный Leader физически сгорит? Второй: как обеспечить низкую задержку записи (Latency) для пользователей из Австралии, если наш единственный Leader находится во Франкфурте?

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

Цена автоматического Failover: Split-Brain и Консенсус

В классической топологии Leader-Follower отказ лидера требует процедуры Failover — назначения нового лидера из числа реплик. Если делать это вручную, система будет недоступна для записи минуты или часы. Автоматический Failover кажется логичным шагом, но он скрывает фатальную уязвимость.

Представьте кластер из двух узлов: Узел А (Leader) и Узел Б (Follower). Между ними пропадает сетевая связность, но оба узла продолжают работать и принимать трафик от клиентов. Узел Б перестает получать Heartbeat от Узла А. По логике автоматического Failover, Узел Б решает, что лидер мертв, и провозглашает лидером себя. Теперь в системе два лидера, каждый из которых принимает независимые записи. Это состояние называется Split-Brain (расщепление мозга). Когда сеть восстановится, базы данных невозможно будет слить воедино без потери данных, так как их истории разошлись.

Чтобы избежать Split-Brain, распределенные системы используют алгоритмы консенсуса (Raft, Paxos) и концепцию кворума. Решение о выборе нового лидера не может принять один узел — за него должно проголосовать строго больше половины узлов кластера.

Если кластер состоит из 3 узлов, для выбора лидера нужно 2 голоса. При сетевом разделении (partition) на две части (1 узел и 2 узла), только часть с двумя узлами сможет собрать кворум и выбрать лидера. Одинокий узел поймет, что оказался в меньшинстве, и откажется принимать операции записи, сохранив консистентность ценой частичной недоступности (выбор в пользу CP по теореме CAP).

Multi-Leader: масштабирование записи и гео-распределение

Алгоритмы консенсуса решают проблему надежности, но не решают проблему задержки при гео-распределении. Если пользователи разбросаны по миру, отправка всех UPDATE запросов в один дата-центр (DC) нарушит наш бюджет Latency.

Топология Multi-Leader (Active-Active) подразумевает наличие нескольких узлов, принимающих запись. Обычно в каждом дата-центре разворачивается свой Leader. Пользователь из США пишет в американский Leader, пользователь из Европы — в европейский. Задержка минимальна. Затем лидеры асинхронно обмениваются изменениями друг с другом.

Главная проблема Multi-Leader — Write Conflicts (конфликты записи). Два пользователя одновременно редактируют одну и ту же страницу в Notion. Один запрос уходит в DC в США, другой — в DC в Европе. Обе записи успешны (каждый лидер локально применил транзакцию). Но когда дата-центры попытаются синхронизироваться, возникнет конфликт: какую версию текста сохранить?

Стратегии разрешения конфликтов:

  1. Last Write Wins (LWW) — побеждает запись с самой поздней временной меткой. Это самый популярный, но опасный подход. В распределенных системах часы на разных серверах не синхронизированы идеально (Clock Skew). Запись, сделанная физически позже, может иметь более раннюю метку времени и быть отброшена. При LWW мы осознанно соглашаемся на потерю части данных.
  2. Маршрутизация по ключу (Sticky Routing) — архитектурный обход проблемы. Мы гарантируем, что все записи для конкретного пользователя (или документа) всегда направляются только в один определенный дата-центр. Конфликт становится невозможным, пока этот дата-центр жив.
  3. Разрешение на уровне приложения — база данных сохраняет обе конфликтующие версии (как ветки в Git) и при следующем чтении отдает их приложению. Ваш Go-код должен содержать бизнес-логику слияния (например, объединить товары в корзине) и записать финальный результат обратно.

Leaderless топология: запись повсюду

Пойдя еще дальше, архитекторы создали топологию Leaderless (вдохновленную Amazon Dynamo, реализованную в Cassandra и Riak). В ней вообще нет концепции лидера.

Любая реплика может принимать запросы на чтение и запись. Клиент (или узел-координатор) отправляет запрос на запись сразу всем доступным репликам параллельно. Но если сети нестабильны, некоторые узлы могут пропустить обновление. Как в системе без лидера гарантировать, что при чтении мы не получим устаревшие данные?

Математика Кворума (W+R>NW + R > N)

В Leaderless системах консистентность настраивается математически для каждого отдельного запроса. Вводятся три параметра:

  • NN — общее количество реплик, хранящих данные.
  • WW — количество узлов, которые должны подтвердить успешную запись, чтобы она считалась успешной для клиента.
  • RR — количество узлов, с которых клиент должен параллельно запросить данные при чтении.

Чтобы гарантировать Strong Consistency (строгую согласованность, при которой чтение всегда возвращает последнюю запись), необходимо соблюдать формулу строгого кворума: W+R>NW + R > N

Если сумма узлов записи и узлов чтения больше общего числа реплик, то по принципу Дирихле множества узлов записи и чтения обязательно пересекутся. Хотя бы один узел из тех, к кому мы обратились при чтении (RR), гарантированно участвовал в последней успешной записи (WW) и имеет актуальную версию.

При чтении с RR узлов клиент получит несколько версий данных. Сравнив их внутренние версии (Version Vectors), клиент отбросит устаревшие и возьмет самую свежую.

Изменение WW и RR — это классический архитектурный Trade-off:

  • Если важна скорость записи: ставим W=1W = 1, но для соблюдения кворума придется делать R=NR = N (медленное чтение, так как опрашиваем всех).
  • Если важна скорость чтения: ставим R=1R = 1, но тогда W=NW = N (медленная и хрупкая запись, отказ одного узла блокирует запись).
  • Баланс: обычно выбирают N=3N = 3, W=2W = 2, R=2R = 2. Это дает отказоустойчивость к падению одного узла и приемлемую скорость обеих операций.

Если же мы сознательно нарушаем формулу и делаем W+RNW + R \leq N (например, W=1,R=1W=1, R=1 при N=3N=3), мы получаем сверхбыструю систему, но переходим в режим Eventual Consistency. Мы рискуем прочитать старые данные, пока они фоном не синхронизируются.

Read Repair и Anti-Entropy

Если узел был отключен от сети и пропустил запись, как он получит актуальные данные в Leaderless системе?

Первый механизм — Read Repair (восстановление при чтении). Когда клиент делает запрос с R=2R=2, а узлы возвращают разные версии, клиент (или координатор) замечает, что один из узлов отстал. Прежде чем вернуть ответ пользователю, координатор отправляет актуальные данные на отставший узел. База данных лечит сама себя за счет пользовательского трафика.

Второй механизм — Anti-Entropy (анти-энтропия). Фоновый процесс, который постоянно сравнивает данные между репликами. Чтобы не пересылать гигабайты данных для сравнения, используются структуры вроде деревьев Меркла (Merkle Trees) — они позволяют за логарифмическое время найти конкретные отличающиеся строки и синхронизировать только их.

Итог

Мы рассмотрели три топологии репликации:

  1. Leader-Follower: простота и отсутствие конфликтов, но единая точка отказа и узкое место для записи.
  2. Multi-Leader: отличная производительность для гео-распределенных клиентов, но высокая сложность разрешения конфликтов записи.
  3. Leaderless: максимальная доступность и настраиваемая консистентность через кворумы, но высокие накладные расходы на Read Repair и фоновую синхронизацию.

Репликация решает проблемы доступности и задержки, а также масштабирует чтение. Но ни одна из этих топологий не решает проблему объема данных. Если ваша база данных вырастает до 10 терабайт, репликация лишь создаст несколько копий по 10 терабайт, требуя невероятно дорогих серверов. Для масштабирования объема и пропускной способности записи нам придется разрезать сами данные на части.

Шардирование: горизонтальное масштабирование записи и распределение данных

Шардирование: горизонтальное масштабирование записи и распределение данных

Репликация позволяет выдерживать колоссальные нагрузки на чтение. Если ваш сервис получает 100 000 запросов на выборку профилей в секунду, вы просто добавляете новые узлы-фолловеры. Но что произойдет, если пользователи начнут массово обновлять свои статусы, генерируя 50 000 операций записи в секунду? Или если объем пользовательских данных достигнет 50 терабайт?

В архитектуре Leader-Follower все операции записи проходят через единый узел — Лидера. Более того, каждый узел в кластере репликации хранит полную, стопроцентную копию всей базы данных. Вы не можете преодолеть физический предел IOPS одного диска или вместить 50 ТБ данных на один 10-терабайтный NVMe-накопитель, просто копируя этот накопитель.

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

Анатомия шардирования

Шардирование (Sharding), или секционирование — это архитектурный паттерн, при котором единый логический набор данных разбивается на множество независимых сегментов, каждый из которых размещается на отдельном физическом узле.

Если репликация — это копирование данных, то шардирование — это их разделение.

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

Каждый такой независимый сегмент называется шардом (Shard).

Разбивая базу данных на шарды, мы решаем сразу две фундаментальные проблемы High-Load:

  1. Масштабирование объема (Storage): База размером 50 ТБ, разбитая на 10 шардов, требует от каждого узла хранения лишь 5 ТБ.
  2. Масштабирование записи (Write Throughput): Если система генерирует 50 000 записей в секунду, то при равномерном распределении по 10 шардам каждый узел будет обрабатывать комфортные 5 000 записей. Общая пропускная способность кластера становится суммой пропускных способностей его частей: Wtotal=WshardW_{total} = \sum W_{shard}.

Ключ шардирования: компас для данных

Когда данные лежали на одном сервере, приложение просто открывало TCP-соединение и отправляло INSERT или SELECT. В шардированном кластере перед приложением (или специальным прокси-сервером) встает новая задача: маршрутизация.

Как понять, на каком из 10 серверов лежит профиль пользователя user_id = 42? Опрашивать все серверы подряд — значит убить производительность и создать проблему N+1 на уровне инфраструктуры.

Для мгновенной маршрутизации вводится ключ шардирования (Shard Key) — конкретный атрибут данных, на основе которого принимается детерминированное решение о размещении записи. Чаще всего это уникальный идентификатор сущности (ID пользователя, ID документа) или атрибут, по которому чаще всего происходит фильтрация (например, tenant_id в B2B SaaS системах).

Самый базовый алгоритм маршрутизации — деление по модулю. Мы берем числовой ключ шардирования, делим его на количество доступных шардов NN и используем остаток от деления как номер целевого шарда: ShardIndex=ID(modN)ShardIndex = ID \pmod N.

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

Если выбрать ключ неправильно, можно получить «горячий шард» (Hot Spot) — ситуацию, когда 90% трафика записи летит в один узел, в то время как остальные простаивают. Например, если шардировать логи по дате создания, то весь сегодняшний трафик будет обрабатывать только один сервер, сводя на нет весь смысл шардирования.

Логические и физические шарды

Наивная реализация шардирования жестко привязывает сегмент данных к конкретному физическому серверу (как в формуле ID(modN)ID \pmod N, где NN — число серверов). Это создает огромную проблему при масштабировании: если серверов станет не 3, а 4, математический остаток от деления изменится для большинства ключей. Данные придется массово перемещать между машинами.

Чтобы отвязать логику распределения от «железа», архитекторы вводят уровень косвенности — логические шарды (Partitions).

Вместо того чтобы делить данные на 3 физических сервера, мы делим их, например, на 1000 логических шардов. А уже эти логические шарды распределяются по физическим узлам с помощью таблицы маршрутизации (Routing Table).

  • Сервер А: логические шарды 0–333
  • Сервер Б: логические шарды 334–666
  • Сервер В: логические шарды 667–999

Если сервер А перегружен, мы можем прозрачно для приложения перенести логические шарды 300–333 на новый сервер Г, просто обновив таблицу маршрутизации. Маршрутизация становится двухэтапной: сначала приложение вычисляет логический шард по ключу, а затем смотрит в конфигурации, на каком IP-адресе этот шард сейчас живет.

Симбиоз: Шардированный кластер

Важно понимать, что шардирование не заменяет репликацию. Это ортогональные концепции, которые в High-Load системах всегда работают вместе.

Если шард — это просто один физический сервер, то его падение приведет к безвозвратной потере части пользовательских данных. Поэтому в реальной архитектуре каждый шард представляет собой полноценный Replica Set (набор реплик).

Логический шард №1 обслуживается Лидером А и его Фолловерами А1, А2. Логический шард №2 — Лидером Б и Фолловерами Б1, Б2.

Такая матричная структура дает максимальную гибкость:

  • Не хватает пропускной способности на запись? Добавляем новые шарды (горизонтальное расширение).
  • Не хватает пропускной способности на чтение конкретного шарда? Добавляем ему фолловеров (вертикальное расширение кластера вниз).
  • Упал физический сервер? Срабатывает механизм Failover внутри конкретного шарда, выбирая нового лидера, при этом остальные шарды даже не замечают сбоя.

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

Алгоритмы шардирования: от Range-based до Consistent Hashing

Алгоритмы шардирования: от Range-based до Consistent Hashing

В прошлой главе мы разделили огромную базу данных на логические шарды. Теперь перед нами стоит сугубо алгоритмическая задача маршрутизации: когда микросервис хочет записать строку с UserID = 42, как он (или балансировщик перед БД) за долю миллисекунды должен понять, на какой именно сервер отправить этот запрос?

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

Range-based шардирование: интуитивно, но опасно

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

Представьте, что мы шардируем пользователей по алфавиту:

  • Шард 1: A — H
  • Шард 2: I — P
  • Шард 3: Q — Z

Математически это выражается через простые проверки неравенств: строка попадает в Шард 1, если AKey<IA \leq Key < I.

Сила диапазонных запросов

Главное и почти единственное преимущество Range-based подхода — идеальная работа с запросами по диапазону (Range Queries). Если вам нужно получить всех пользователей на букву «B», база данных точно знает, что все они лежат в Шарде 1. Остальные узлы кластера даже не узнают об этом запросе, их ресурсы сэкономлены.

Проблема Hotspot (Горячей точки)

На практике шардирование по алфавиту используется редко из-за неравномерности данных (фамилий на «S» больше, чем на «X»). Гораздо чаще Range-based применяют для временных рядов (Time-series) — логов, метрик, транзакций, где ключом выступает Timestamp.

  • Шард 1: Январь — Март
  • Шард 2: Апрель — Июнь
  • Шард 3: Июль — Сентябрь

И здесь кроется фатальный архитектурный изъян. Данные генерируются прямо сейчас. Это значит, что 100% операций записи (INSERT) в текущем квартале полетят строго в Шард 3.

Возникает Hotspot (Горячая точка) — ситуация, когда один узел кластера перегружен и работает на пределе CPU и диска, в то время как остальные шарды простаивают, обслуживая лишь редкие запросы на чтение исторических данных. Горизонтальное масштабирование теряет смысл: добавление Шарда 4 не разгрузит Шард 3 прямо сейчас.

Hash-based шардирование: математическая справедливость

Чтобы победить Hotspot, нам нужно «размазать» последовательные данные по всем доступным узлам. Для этого применяется хэширование.

Вместо того чтобы смотреть на само значение ключа, мы пропускаем его через криптографическую или быструю хэш-функцию (например, MurmurHash), получая псевдослучайное, но детерминированное число. Затем вычисляем остаток от деления на количество шардов NN:

Shard=hash(Key)(modN)Shard = hash(Key) \pmod N

Если у нас 3 шарда, то последовательные ID пользователей (101, 102, 103) после хэширования будут распределены примерно так:

  • hash(101) % 3 \rightarrow Шард 2
  • hash(102) % 3 \rightarrow Шард 0
  • hash(103) % 3 \rightarrow Шард 1

Запись идеально сбалансирована. Ни один узел не перегружен.

Цена баланса: Scatter-Gather

За равномерное распределение записи мы платим производительностью чтения диапазонов.

Если мы шардировали таблицу заказов по хэшу от OrderID, то запрос SELECT * FROM orders WHERE OrderID = 105 выполнится мгновенно — мы вычислим хэш и пойдем в конкретный шард.

Но что если менеджеру нужен отчет: SELECT * FROM orders WHERE date BETWEEN '2023-01-01' AND '2023-01-31'? Поскольку мы шардировали по хэшу ID, заказы за январь разбросаны случайным образом по всем узлам кластера.

Системе придется применить паттерн Scatter-Gather (Разброс и Сбор):

  1. Scatter: Отправить этот запрос параллельно на все шарды.
  2. Каждый шард выполнит локальный поиск и вернет свою часть январских заказов.
  3. Gather: Координатор (или API Gateway) дождется ответов от всех, объединит их, отсортирует в памяти и только потом отдаст клиенту.

Consistent Hashing и виртуальные узлы (vNodes)

Hash-based шардирование через остаток от деления (hash(Key)(modN)hash(Key) \pmod N) отлично работает, пока NN (количество шардов) неизменно. Но в пятом курсе, разбирая балансировку нагрузки, мы выяснили, что изменение NN ломает маршрутизацию: если увеличить кластер с 3 до 4 узлов, остаток от деления изменится почти для всех ключей. В контексте базы данных это означает необходимость перенести 75% терабайтных данных по сети на новые места. Это убьет кластер.

Решением выступает Consistent Hashing (Консистентное хэширование). Напомним базовый принцип: мы сворачиваем пространство хэшей в кольцо. И узлы кластера, и сами данные хэшируются и помещаются на это кольцо. Данные приписываются первому узлу, который встретится при движении по часовой стрелке. При добавлении нового узла перемещаются только те данные, которые оказались между новым узлом и его предшественником.

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

В теории кольцо решает проблему решардинга. На практике в базах данных возникает нюанс. Если у нас всего 3-5 физических серверов, их хэши могут распределиться по кольцу неравномерно. Один сервер захватит 60% кольца, а два других — по 20%. Мы снова получим Hotspot, но уже не из-за природы данных, а из-за математической случайности.

Виртуальные узлы (vNodes)

Чтобы сбалансировать кольцо, физические серверы разбивают на Виртуальные узлы (vNodes). Вместо того чтобы хэшировать IP-адрес сервера один раз, мы хэшируем его 100 или 256 раз с разными суффиксами (Node1-001, Node1-002...).

Один физический сервер теперь представлен сотней маркеров на кольце. Они чередуются с маркерами других серверов. Чем больше vNodes, тем идеальнее балансировка данных между физическими машинами. Более того, если у вас есть мощный сервер (64 ядра) и слабый (32 ядра), вы можете дать мощному 200 vNodes, а слабому — 100, реализовав взвешенное шардирование.

Directory-based шардирование: абсолютный контроль

Иногда ни хэш, ни диапазоны не подходят под бизнес-логику. Например, вы хотите, чтобы данные премиум-пользователей лежали на быстрых NVMe-дисках (Шард 1), пользователи из Европы — в дата-центре во Франкфурте (Шард 2), а заблокированные аккаунты — в холодном хранилище (Шард 3). Никакая математическая функция не опишет эти правила.

В этом случае применяется Directory-based шардирование. Создается отдельная служебная база данных (Routing Table), которая хранит явный маппинг:

  • UserID 10 \rightarrow Шард 2
  • UserID 11 \rightarrow Шард 1

Преимущество: абсолютная гибкость. Перенос пользователя на другой шард — это просто обновление одной строки в Routing Table. Недостаток: каждое обращение к БД теперь требует двух запросов (сначала узнать шард, затем пойти за данными). Эта служебная база сама становится узким местом (Bottleneck) и точкой отказа (SPOF), поэтому ее приходится агрессивно кэшировать в памяти приложения.

Сводная таблица алгоритмов

Алгоритм Балансировка записи Запросы по диапазону Сложность решардинга Применение
Range-based Плохая (риск Hotspot) Отличная Средняя Логи, аналитика по датам
Hash-based (Modulo) Отличная Плохая (Scatter-Gather) Очень сложная Статичные кластеры
Consistent Hashing Отличная (с vNodes) Плохая (Scatter-Gather) Легкая Динамические NoSQL БД (Cassandra, DynamoDB)
Directory-based Зависит от логики Зависит от логики Легкая Сложная бизнес-логика, Multi-tenant SaaS

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

Сложности распределенного слоя данных: Join, транзакции и решардинг

Сложности распределенного слоя данных: Join, транзакции и решардинг

Вы разделили базу данных на 10 шардов. Запись масштабировалась линейно, CPU отдыхает, а графики метрик выглядят идеально. Иллюзия победы длится ровно до тех пор, пока продуктовый менеджер не просит добавить на главную страницу простой блок: «Последние заказы пользователя с именами курьеров».

В монолитной БД это решалось одним JOIN. Но теперь пользователи лежат на шарде №2, заказы — на шарде №7, а курьеры — на шарде №4. Физические границы серверов разрушают магию реляционных баз данных.

Шардирование — это архитектурная сделка. Вы получаете бесконечную емкость хранения и пропускную способность записи, но взамен отдаете три фундаментальные возможности классических БД: локальные джоины, ACID-транзакции и легкость изменения топологии. Разберем, как архитекторы платят по этим счетам.

Проблема №1: Cross-shard Joins (Распределенные объединения)

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

Есть три основных паттерна решения этой проблемы.

1. Scatter-Gather на уровне приложения

Мы уже упоминали этот паттерн в прошлой главе для поиска. Для JOIN он работает так же: приложение берет на себя роль координатора.

Сначала Go-сервис делает запрос в шард пользователей, получает данные. Затем он извлекает из них ID заказов и делает параллельные запросы (Fan-Out) в шарды заказов. Наконец, горутина собирает ответы и склеивает их в памяти (Fan-In).

Минусы: Сетевые задержки и риск проблемы N+1. Если нужно отсортировать объединенный результат с пагинацией (LIMIT/OFFSET), придется выкачивать огромные массивы данных в память Go-приложения.

2. Глобальные таблицы (Broadcast)

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

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

3. Денормализация данных

Если Scatter-Gather слишком медленный, а таблица слишком велика для Broadcast, мы идем на нарушение нормальных форм БД — мы дублируем данные.

Если при чтении заказа нам всегда нужно показывать имя пользователя, мы добавляем колонку user_name прямо в таблицу orders. Да, мы тратим больше места на диске. Да, при смене имени пользователем нам придется обновить не одну строку, а тысячи его заказов (аномалия обновления). Но в High-Load системах мы сознательно жертвуем сложностью записи ради скорости чтения (которое происходит в 100 раз чаще).

Проблема №2: Распределенные транзакции

Представьте перевод 100 RUB от Алисы (Шард А) к Бобу (Шард Б). В монолите мы бы открыли транзакцию, списали деньги, начислили деньги и сделали COMMIT. Если сервер упадет посередине, БД сделает ROLLBACK.

В распределенной системе возникает проблема частичного отказа (Partial Failure). Шард А успешно списал деньги, а Шард Б недоступен из-за сетевого разрыва. Деньги исчезли.

Two-Phase Commit (2PC)

Для обеспечения строгой консистентности (Strong Consistency) между шардами используется алгоритм двухфазного коммита (2PC). В нем появляется новая сущность — Координатор.

Алгоритм делится на две фазы:

  1. Prepare (Подготовка): Координатор спрашивает все шарды-участники: «Вы готовы выполнить транзакцию?». Шарды проверяют ограничения (хватает ли баланса), блокируют нужные строки и отвечают «Да».
  2. Commit (Фиксация): Если все ответили «Да», координатор отправляет команду «Фиксируй!». Если хотя бы один ответил «Нет» (или отвалился по таймауту), координатор отправляет «Отмена» (Abort).

Почему 2PC избегают в High-Load?

  • Блокировки: Строки Алисы и Боба заблокированы на всё время сетевого общения. Если координатор «зависнет» между фазами, строки останутся заблокированными бесконечно долго (до ручного вмешательства).
  • Закон Амдала: Скорость транзакции равна скорости самого медленного шарда-участника.

Альтернатива: Паттерн Saga

Поскольку 2PC убивает пропускную способность, современные системы используют паттерн Saga.

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

Если шаг падает, запускаются компенсирующие транзакции в обратном порядке.

  1. Шард А: Списать 100 RUB у Алисы (Успех) \rightarrow Событие «Деньги списаны».
  2. Шард Б: Начислить 100 RUB Бобу (Ошибка: счет заблокирован) \rightarrow Событие «Ошибка начисления».
  3. Шард А (Компенсация): Вернуть 100 RUB Алисе.

Saga не дает строгой консистентности. В промежутке между шагом 1 и 3 Алиса может увидеть нулевой баланс, хотя перевод не удался. Это Eventual Consistency в действии — мы жертвуем временной точностью ради высочайшей производительности и отказоустойчивости.

Проблема №3: Решардинг без даунтайма (Zero-Downtime Migration)

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

Но алгоритм лишь говорит какие данные нужно перенести. Физический перенос 500 ГБ данных под нагрузкой в 10 000 RPS — это инженерный вызов. Мы не можем просто остановить БД на час, скопировать файлы и запустить снова.

Процесс переноса логического шарда (назовем его Чанк X) со старого узла на новый выполняется в несколько этапов, незаметно для пользователей.

Фазы миграции:

  1. Снапшот и копирование: Система делает снимок Чанка X на старом узле и начинает фоновое копирование на новый узел. Клиенты продолжают писать и читать со старого узла.
  2. Нагоняющая репликация (Catch-up): Пока шло копирование, данные изменились. Старый узел начинает отправлять на новый узел свой WAL (Write-Ahead Log) для Чанка X, чтобы новый узел применил пропущенные изменения.
  3. Двойная запись (Dual Write) / Маршрутизация: Когда новый узел почти догнал старый, API Gateway или балансировщик (Routing Table) переключает запись: новые данные пишутся синхронно в оба узла, а чтение пока идет со старого.
  4. Переключение (Cut-over): Чтение переключается на новый узел.
  5. Очистка: Чанк X удаляется со старого узла, освобождая место.

Этот процесс требует сложной оркестрации. В современных распределенных базах данных (таких как CockroachDB, YugabyteDB или MongoDB в режиме Sharded Cluster) эта логика встроена в ядро СУБД. Если же вы шардируете PostgreSQL на уровне приложения, вам придется реализовывать этот конвейер миграции силами инфраструктурной команды.

Итоги

Масштабирование слоя данных — это всегда компромисс. Хотите быстрые джоины? Придется денормализовать данные и усложнить логику обновления. Нужны строгие распределенные транзакции? Готовьтесь к падению пропускной способности из-за 2PC. Хотите эластичности? Придется строить сложные пайплайны решардинга.

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

Архитектурный синтез: комбинирование репликации и шардирования в High-Load

Архитектурный синтез: комбинирование репликации и шардирования в High-Load

Если вы разделите базу данных объемом 50 ТБ на 10 шардов по 5 ТБ, вы решите проблему предела записи и объема. Но что произойдет, если жесткий диск на сервере третьего шарда выйдет из строя? Вы потеряете ровно 10% пользовательских данных. Шардирование без репликации — это не масштабирование, это распределение единой точки отказа (SPOF) на множество серверов.

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

Двумерная матрица данных: Шардированные реплика-сеты

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

  1. Горизонтальная ось (Шардирование): Разделение данных на логические сегменты. Обеспечивает масштабирование объема (Capacity) и пропускной способности записи (Write Throughput).
  2. Вертикальная ось (Репликация): Копирование каждого сегмента. Обеспечивает отказоустойчивость (High Availability) и масштабирование чтения (Read Throughput).

В такой архитектуре логический шард перестает быть одним физическим сервером. Каждый шард — это полноценный кластер Leader-Follower (реплика-сет).

Если у нас есть система из 3 шардов, и каждый шард имеет фактор репликации 3 (1 Leader + 2 Followers), физически наш слой данных состоит из 9 серверов. Запись ключа, принадлежащего первому шарду, попадает на Leader-узел первого шарда, а затем через WAL асинхронно растекается на два его Follower-узла. Остальные 6 серверов в этом процессе не участвуют.

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

Умная маршрутизация: как найти нужный узел

В монолитной базе данных приложение просто открывает TCP-соединение и отправляет SQL-запрос. В двумерной матрице из десятков серверов клиент не может и не должен знать топологию сети. Эту задачу берет на себя промежуточный слой — Smart Router (например, mongos в MongoDB или VTGate в Vitess для MySQL).

Роутер выступает фасадом и выполняет двухэтапную маршрутизацию каждого запроса:

  1. Поиск по оси X (Выбор шарда): Роутер извлекает Shard Key из запроса (например, user_id), применяет алгоритм (хэширование или поиск по Directory-таблице) и определяет, какому логическому шарду принадлежат данные.
  2. Поиск по оси Y (Выбор роли): Роутер анализирует тип запроса. Если это UPDATE или INSERT, запрос направляется строго на Leader-узел выбранного шарда. Если это SELECT, роутер балансирует запрос между Follower-узлами этого же шарда (например, алгоритмом Least Connections), учитывая возможное отставание (Replication Lag).

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

Изоляция отказов: что происходит при падении узла

Главное преимущество синтеза репликации и шардирования — радикальное уменьшение радиуса поражения (Blast Radius) при авариях.

Представим, что в нашем кластере из 3 шардов (по 3 узла в каждом) происходит аппаратный сбой, и физический сервер, являющийся Leader-узлом Шарда №2, сгорает.

Как реагирует система:

  1. Шарды №1 и №3: Продолжают работу без малейших изменений. Запись и чтение пользователей, чьи данные лежат там, не прерываются ни на миллисекунду.
  2. Чтение в Шарде №2: Роутер мгновенно исключает упавший Leader из пула и продолжает отправлять SELECT запросы на два оставшихся Follower-узла. Чтение не деградирует.
  3. Запись в Шарде №2: Происходит кратковременная пауза (обычно несколько секунд). Оставшиеся узлы Шарда №2 инициируют выборы нового лидера. Поскольку их двое, они собирают кворум, один из них становится новым Leader-узлом, и роутер возобновляет запись.

Архитектура предотвращает каскадный отказ. Вместо полной остановки системы (как было бы с одним сервером) или потери данных (как было бы при чистом шардировании), мы получаем лишь секундную недоступность записи для 33% пользователей.

Решардинг в двумерном мире

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

Если Шард №1 переполнен и его нужно разделить на Шард №1A и Шард №1B, инженеры не просто копируют файлы. Сначала в кластер добавляются новые пустые серверы, из которых собираются два новых полноценных реплика-сета. Затем роутер настраивается на паттерн Dual Write: новые данные пишутся и в старый Leader Шарда №1, и в Leader-узлы новых шардов. Фоновый процесс переносит исторические данные.

Как только Replication Lag между старым шардом и новыми становится нулевым, роутер атомарно обновляет свою карту маршрутизации. Старый реплика-сет Шарда №1 выводится из эксплуатации. Благодаря тому, что чтение обслуживается Follower-узлами, пользователи не замечают миграции терабайтов данных, происходящей в фоне.

Итоги модуля

Масштабирование слоя данных — это не поиск «серебряной пули» или идеальной базы данных. Это инженерия компромиссов.

Вы начали с одного сервера, уперлись в потолок ресурсов (Bottleneck) и добавили реплики, чтобы масштабировать чтение. Вы столкнулись с отставанием репликации и научились разделять пулы соединений. Когда объем данных превысил возможности одного диска, вы внедрили шардирование, разрезав базу на части, и столкнулись с болью распределенных транзакций (2PC, Saga) и Cross-shard Joins.

Синтез этих паттернов в единую матрицу шардированных реплика-сетов — это вершина архитектуры реляционных и NoSQL баз данных под высокой нагрузкой. Следующий шаг в ускорении системы лежит уже за пределами дисковых хранилищ — в оперативной памяти и продвинутых паттернах кэширования.