title: Статья 05 — Репликация и шардирование source: подготовлено 2026-09-17, TASK-45.45; каркас — brief.md §05; развёртка knowledge-base.md §5–6 (мост в §7 и §9); источники — DDIA 2-е изд. гл. 6–7 (TASK-45.39), Alex Xu гл. 5–6 (TASK-45.38), эталоны classic-designs.md §1–3
Статья 05: репликация и шардирование
Пятая статья серии — два способа размножить данные, когда одна нода перестаёт справляться. Критерий «умею, если» из брифа: могу обосновать схему репликации и ключ шардирования для ленты/чата и сказать, что происходит при отказе ноды. Обе половины критерия — про хранилище этапа 4 (детали + БД) и про фокус нашей секции целиком: отказоустойчивость проверяется именно здесь, на узлах с данными.
Место раздела в каркасе: статья 03 (03-scaling.md) закончилась словами «стейт — в правильных местах»; статья 04 — какие это хранилища по классам; эта — как эти хранилища переживают рост и отказы. Связка с соседями: лаг репликации и кворумы — это гарантии согласованности (раздел 09, углубление), failover и поведение при отказах — материал этапа 6 (раздел 08). Здесь — верхний уровень с рабочими формулировками; глубина — ../knowledge-base.md §5–6 и DDIA гл. 6–7.
1. Два способа размножить данные — и почему они вместе
Открывающая пара определений DDIA (гл. 6–7), на которой держится весь раздел:
- Репликация — копии одних и тех же данных на нескольких нодах: ближе к пользователю (латентность), переживать отказы части нод (доступность, durability), больше нод обслуживают чтение (read-масштаб).
- Шардирование — разрезание данных: каждая запись принадлежит ровно одному шарду, когда данные или write-throughput не влезают в одну ноду.
Ответ на вопрос «репликация или шардирование?» почти всегда — оба: каждый шард сам реплицирован (физический узел держит несколько шардов: для одних — лидер, для других — реплика). Так и рисуем на схеме: кластер из нод, внутри — шарды, у каждого шарда — лидер и реплики. Раздел отвечает на два вопроса по очереди: как согласуются копии (§2–5) и как выбирается разрез (§6–8).
Важная оговорка ДДИА: репликация ≠ бэкап. Реплика повторяет и наши ошибки (удалил — удалилось везде); бэкап хранит старые снапшоты, чтобы откатиться во времени. Проговорить, если интервьюер путает.
2. Репликация single-leader: механика и выбор синхронности
2.1 Механика
Дефолт индустрии: один лидер (primary) принимает все записи и пишет их в лог; фолловеры (read-реплики) тянут этот лог (WAL / binlog / change stream) и применяют записи в том же порядке. Чтение — с лидера или любой реплики; запись — только через лидера. Так работают PostgreSQL, MySQL, Mongo, Kafka, Raft-системы. На схеме из статьи 03 это блок «БД: primary + реплики» — здесь он разворачивается.
Как без даунтайма добавить реплику (типовой вопрос-зонд): снапшот лидера в момент T (без глобальной блокировки) → копия на новую ноду → реплика подтягивает из лога лидера все изменения после позиции снапшота (LSN в Postgres, GTID/binlog coordinates в MySQL) → догнала (catch-up) → работает дальше. Ключевое слово — позиция в логе: лог с позициями понадобится ещё дважды (read-your-writes и CDC).
2.2 Sync / async / semi-sync
Главный трейдофф репликации — ждать ли ACK реплики перед ответом клиенту:
| Режим | Что получаем | Чем платим |
|---|---|---|
| Sync | реплика гарантированно догнала запись: лидер упал — данные живы | недоступная/медленная реплика блокирует все записи; all-sync непрактично |
| Semi-sync (ждём ACK одной реплики) | данные минимум на двух нодах при разумной латентности | чуть дороже записи; компромисс-дефолт для «деньги и важное» |
| Async | латентность записи не зависит от реплик, лидер пишет даже если все отстали | лаг (от секунд до минут под нагрузкой) и потери хвоста при failover |
Формулировка на секцию: «semi-sync по умолчанию: лидер + одна синхронная реплика — выживает отказ одной ноды без потери подтверждённых записей; полностью async — там, где допустима eventual consistency». Почему нельзя all-sync: чем больше нод, тем вероятнее, что одна из них сейчас недоступна — запись встанет целиком.
2.3 Что именно реплицируем
Три формата лога (нужно на уровне «назвать и знать цену»): statement-based (прислать SQL) — ломается на now()/rand() и порядке исполнения; WAL-shipping (физический лог движка) — быстрый, но привязан к версии storage-движка — мешает zero-downtime upgrades; logical / row-based лог (строка + новые значения) — переносимый компромисс, а ещё — фундамент CDC (Debezium): внешние системы читают те же события. Это мост в статью 07 (outbox/CDC).
3. Replication lag: что ломается и сессионные гарантии
Цена async-репликации — replication lag: чтение с отставшей реплики возвращает старое состояние. Это eventual consistency: когда-нибудь догонит, но когда — гарантии нет. DDIA разбирает три аномалии — их нужно узнавать по имени и лечить:
- Read-your-writes («сохранил профиль — перезагрузил — изменения нет»): пользователь видит устаревшую версию своей же записи. Лечение: читать потенциально-свои данные с лидера (свой профиль редактирует только владелец — правило дёшево); либо окно после записи (минуту после апдейта — с лидера); либо — правильнее — по позиции в логе: клиент запомнил LSN своей записи, реплика, не дотянувшая эту позицию, не обслуживает его чтение.
- Monotonic reads («время идёт назад»: первый запрос ушёл на свежую реплику, второй — на отставшую, комментарий появился и исчез): лечение — приклеить пользователя к одной реплике (хеш user_id), а не случайный выбор.
- Consistent prefix reads (ответ прилетел раньше вопроса): причинно связанные записи должны читаться в порядке записи; в шардах усугубляется (разные шарды — разный лаг). Лечение — причинно связанные данные в одном шарде.
Итог-формулировка: «асинхронная репликация даёт eventual consistency; сессионные гарантии (read-your-writes, monotonic reads) я обеспечиваю на уровне приложения — маршрутизацией чтений». Мониторинг лага — обязателен (SLO на секунды, алерт при отставании). Если лаг неприемлем и функционально, а лечить маршрутизацией дорого — следующий шаг разговор о stronger consistency (раздел 09), не «притворяться, что async — это sync».
4. Альтернативы single-leader: multi-leader и leaderless
4.1 Multi-leader
Несколько пишущих узлов, изменения гоняются между всеми. Зачем: мульти-ДЦ (лидер в каждом регионе — запись не платит межрегиональный RTT при живом регионе) и офлайн-клиенты (календарь/заметки пишут локально, синхронизируются позже). Плата — конфликты записей: две стороны обновили одно поле. Решения: last-write-wins по таймстампу (просто, но тихо теряет данные — часы-то врут), CRDT/OT (для коллаб-редакторов), ручное слияние (отдать обе версии пользователю). Отдельная мина: автоинкременты и sequence конфликтуют между лидерами. Вердикт Клеппмана — «dangerous territory»; наша позиция на секции: «multi-leader избегаю без необходимости: межрегиональную запись честнее решать single-leader в регионе-лидере, а конфликты — это отдельный класс багов».
4.2 Leaderless (Dynamo-стиль)
Лидера нет: координатор пишет ключ в W из N реплик, читает из R из N и берёт самую свежую версию по отметкам времени/версиям (параллельные версии ловятся векторными часами — раздел 09). Условие W + R > N гарантирует пересечение множеств: хотя бы одна из R прочитанных реплик видела запись — «читаем свежее». Настройка под кейс (tunable consistency): W=N, R=1 — строгая запись, дешёвое чтение; W=1, R=N — наоборот; середина — баланс. Инфраструктура: read repair (свежий читатель заодно чинит отставшую реплику), anti-entropy фоном (Dynamo: Merkle-деревья — сравнение реплик за O(log n) обменом хешей), членство и детект отказов — gossip (никакого центрального детектора = нет его SPOF). При недоступности «своих» узлов ключа — sloppy quorum + hinted handoff: пишем на любые доступные узлы с пометкой-владельцем, позже отдаём законному хозяину.
Две оговорки, отличающие сильный ответ. Первая: кворум ≠ linearizability — W+R>N даёт свежесть чтения после завершённой записи, но гонки параллельных записей и read repair он не сериализует (доказательство — DDIA гл. 10 / KB §7); строгая согласованность — это single-leader или консенсус. Вторая: кворум-конфигурация — это цена латентности: W+R>N при меж-ДЦ репликах превращает каждую запись/чтение в ожидание самого медленного региона (PACELC).
5. Failover: что происходит при отказе ноды с данными
Отказ реплики — скучный и правильный сюжет: она не в обслуживании, читаем с остальных; вернулась — по своему логу подтянула пропущенное (catch-up). Отказ лидера — самое интересное место раздела. Автоматический failover — три шага:
- Детект. Heartbeat'ы; лидер молчит таймаут (~30 с) — считается мёртвым. Ловушка таймаута: длинный — дольше простой, короткий — ложные failover при всплеске нагрузки или сетевом чихе (ложный failover сам является инцидентом — при уже больной системе делает хуже). Детект — кворумом реплик, не одной нодой.
- Выборы нового лидера. Голосование большинства реплик или решение контролера; кандидат — самая свежая реплика (по позиции в логе): теряем минимум подтверждённых записей. Согласовать выборы у всех узлов — задача консенсуса (Raft; эталон: managed-СУБД переключают primary именно так).
- Реконфигурация. Клиенты пишут новому лидеру; остальные реплики тянут лог с него. Вернувшийся старый лидер обязан стать фолловером.
Ловушки — называть все три:
- Потеря async-хвоста. Непрореплицированные подтверждённые записи самого свежего кандидата просто теряются («вы думали, что записали»). Худший кейс из DDIA — GitHub: отстающая реплика стала лидером, её автоинкремент отстал → переиспользование PK → расхождение MySQL и Redis → приватные данные чужим пользователям. Мораль: автоинкремент в распределённой системе — отдельная мина (снежинка/ULID — раздел 10).
- Split brain. Два узла считают себя лидером и оба пишут → данные теряются/портятся. Лечение — fencing: у лидера есть эпоха/term; действия старого (меньшая эпоха) отвергаются всеми; страховочный предохранитель — при обнаружении двух лидеров гасим обоих (лучше недоступность, чем порча).
- Разогрев. Новый лидер первый раз читает горячие страницы — кэши холодные; лаги и шторм ретраев от клиентов, не дождавшихся таймаута (jitter!).
Позиция для секции: «автоматика быстрее, но ложный failover дороже простоя — для критичных данных я предпочитаю полуручной failover с чётким runbook'ом; и всегда fencing». Отказ целого шарда — см. §10; отказ ДЦ — раздел 08 (GeoDNS + продвижение реплик в живых ДЦ).
6. Шардирование: последняя ступень и главное решение
6.1 «А надо ли?» — лестница до шардов
Шардирование ломает транзакции, JOIN и глобальные ограничения — поэтому оно последняя ступень, а не первая реакция на рост:
индексы → вертикальный масштаб (железо) → read-реплики → кэш → денормализация/проекции → шарды.
Сигнал «пора»: storage или write-throughput не влезает в ноду после того, как отработали предыдущие ступени (эталон Яндекса: ~100 ГБ при 500k RPS чтения → «данные не шардировать, а реплицировать» — KV-зеркала; и обратный пример: посты 70 ТБ/год → шардированная БД). Проговорить лестницу обязательно — это признак зрелости: «сначала дёшево, потом сложно».
6.2 Шард-ключ — решение №1
Ключ шардирования выбирается по access pattern, а не по тому, что «есть уникальное»: все горячие запросы должны попадать в один шард. Цена ошибки максимальна: сменить ключ после — переливка всех данных. Контрпример-канон: чат, зашардированный по message_id — каждое чтение диалога = scatter-gather по всем шардам; правильный ключ — conversation_id (все сообщения диалога вместе, порядок внутри). Кросс-шардные транзакции и JOIN — избегать (денормализация; 2PC — осторожно, раздел 09).
7. Стратегии раскладки: от range до consistent hashing
| Стратегия | Как | Плюс | Минус |
|---|---|---|---|
| Range | диапазоны ключа (даты, ID-отрезки) | range-запросы работают | hot spot на свежем крае: всё пишется в шард «сегодня» |
| Hash | hash(key) mod N |
равномерность | range-запросов нет; смена N двигает почти всё |
| Consistent hashing | кольцо: хешируем ключи и узлы в одно пространство, ключ обслуживает ближайший по часовой | при добавлении/уходе узла переезжает ~1/N ключей | без виртуальных нод — неравномерно |
| Directory | lookup-сервис: карта ключ→шард | гибкость (неровные данные, переезды) | сама карта — точка отказа/кэширования |
Два обязательных дополнения к consistent hashing. Виртуальные ноды: каждый физический узел — 100–200 точек на кольце → распределение близко к равномерному, и уходим/приходим «по кусочку» с разных участков кольца. Применения — шардирование БД/кэша, Redis-кластер, routing в CDN/ДЦ. И анти-паттерн mod N: hash(key) % N при смене N пересчитывает почти все отображения — лавина промахов и перегрузка (снежная история про кэш из Сюя, гл. 5). Честная оговорка: у consistent hashing есть цена отладки; при малых N таблица соответствий (lookup) проще — знать как трейдофф.
8. Жизнь шардов: hot partition, ребалансировка, индексы, routing
Hot partition (селебрити). Один перегретый ключ перекашивает шард: миллионник пишет пост — его шард горит, остальные скучают. Средства: salting — ключ + случайный суффикс, записи размазываются по слотам, чтение агрегирует; выделение горячего в отдельный шард/сервис (селебрити — отдельный путь fan-out on read, эталон ленты §9); локальный кэш на ноде (read-mostly). В эталоне ленты так и проговариваем: «кэш по user_id через consistent hashing + реплики горячих ключей».
Ребалансировка без даунтайма. Вопрос «как перекладывать шарды» имеет канонический ответ: двигаем не данные, а границы. Приёмы: заранее фиксированное число маленьких шардов (миграция = перенос нескольких готовых шардов на новый узел; метаданных больше — плата); либо партиций на узел втрое больше, чем нужно (запас под рост); либо динамическое расщепление переполненного шарда. Сам переезд — live-миграция с dual-write/dual-read + backfill и переключением чтения после выравнивания. Никогда — mod N (§7).
Вторичные индексы (классическая ловушка). Индекс не по шард-ключу бывает local (каждый шард ищет сам → запрос по всем шардам, scatter-gather, дорого) и global (сам распределён по ключу индекса → поиск в один прыжок, но отстаёт от записи и уязвим к hot spot'ам по индексируемому полю). Формулировка: «вторичные запросы у меня редкие — local-индексы + scatter-gather; иначе — глобальный индекс ценой лага».
Request routing («где мои данные?»): слой маршрутизации (координатор) + метаданные членства в координационном сервисе (ZooKeeper/etcd) + кэш на клиенте; изменение членства — gossip/notifications. На секции достаточно: «координатор знает карту ключ→шард, клиент кэширует; изменение членства — событие». Сюда же мультиарендность: tenant_id в ключе шарда, шумные соседи — выделенные шарды.
9. Прикладное: лента и чат — критерий брифа
Критерий «умею, если» требует обосновать решения для двух сервисов — собираем всё выше на материале эталонов (../classic-designs.md §2–3):
Лента новостей (инстаграм-класс). Репликация: primary + реплики чтения, async с обеспеченными read-your-writes — свой пост синхронно вставляется в свою ленту на Post API (и L1-кэш), чужие ленты — eventual с лагом в секунды. Шард-ключи по access pattern: посты — по post_id (hash; источник правды, равномерность, чтение поста — по ID); соцграф — отдельная БД по user_id (горячий запрос постинга — «подписчики автора» — одним шардом); кэш лент — Redis-кластер по user_id через consistent hashing (лента = один ключ user_feed:<id>). Hot partition — селебрити: их не пушим (fan-out on read), кэшируем агрессивно; горячие ключи — реплики ключа. Числа-обоснование: 70 ТБ/год постов → шардированная БД; 130 ГБ горячих лент (80/20) → кэш-кластер без шардирования «в лоб».
Мессенджер. Репликация: сообщения диалога — primary+реплики; онлайн-доставка — stateless-шлюзы с реестром user_id → gateway_id (Redis, TTL+heartbeat) — отдельно от данных. Шард-ключ: conversation_id (все сообщения диалога в одном шарде, порядок per-dialog seq; чтение диалога — один шард). message_id достаточно уникального в рамках канала (составной ключ (channel_id, message_id)); глобальный snowflake — если нужна сортировка/шардирование по времени. Отказ шлюза — reconnect с jitter (шторм 500k переподключений — тот же thundering herd, раздел 06).
Формула ответа на «обоснуйте»: «ключ = гранулярность горячего чтения, репликация = semi-sync/async с явной сессионной гарантией, hot key — отдельный план» — три предложения, три раздела статьи.
10. Что происходит при отказе ноды (сводная таблица)
Прогон «узел → отказ → что происходит → что видно в метриках» — обязательный финал разговора о данных (этап 6):
| Узел | Что происходит | Механизм | Что в метриках |
|---|---|---|---|
| Реплика чтения (follower) | Клиенты не замечают: LB/роутер читает с остальных | catch-up по своему логу после возврата | lag остальных реплик ↑ (нагрузка перераспределилась) |
| Лидер БД | Короткая недоступность записи → failover | детект (heartbeat, кворум) → выборы самой свежей реплики → реконфигурация; потери async-хвоста; fencing старого | время failover; ошибки записи; после — replication lag 0 |
| Нода кэша (Redis) | Волна промахов по «своим» ключам | consistent hashing: виртуальные ноды переразложились, ~1/N ключей переезжает; miss-шторм гасим single-flight (раздел 06) | hit rate ↓, latency ↑, RPS на БД ↑ |
| Лидер шарда | Шард без записи до продвижения реплики | replica promotion шарда; чтение — с реплик шарда / через кэш (эталон ленты: «посты шарда не открываются — чтение смягчает кэш») | ошибки по скоупу шарда, не по всем |
| Шлюз/воркер (stateless) | Ротация без потери | stateless: вывели из ротации (статья 03) | 5xx-всплеск до вывода |
| Целый ДЦ | Регион недоступен | GeoDNS переливает трафик; лидеры шардов за ДЦ → продвижение реплик в живых ДЦ (majority — раздел 08) | доступность по регионам, failover-счётчики |
Универсальная последовательность при любом отказе: детект → изоляция/переключение → деградация вместо падения → метрика, по которой видно. Отдельная строчка для интервьюера: «split brain и потерю async-хвоста я не лечу, я их предотвращаю — semi-sync + fencing + выборы самой свежей реплики».
11. Ошибки этого раздела
Подмножество топ-10 ошибок (../methodology.md §4), относящееся к репликации и шардированию:
| Ошибка | Как выглядит | Противоядие |
|---|---|---|
| «Зашардировать сразу» | Шарды при 100 ГБ и read-heavy | Лестница §6.1: индексы → реплики → кэш → … → шарды |
| Шард-ключ мимо access pattern | Чат по message_id, каждый диалог — scatter-gather |
Ключ = гранулярность горячего чтения (§6.2, §9) |
| mod N при ребалансировке | Добавили узел — переложили почти всё | Consistent hashing + виртуальные ноды (§7) |
| Async без сессионных гарантий | «Изменения потерялись» после записи | Read-your-writes: лидер / LSN / окно (§3) |
| Multi-leader как дефолт | «Пишем во все регионы» без стратегии конфликтов | «Избегаю без необходимости» + LWW теряет данные (§4.1) |
| Автоинкремент в распределённой системе | GitHub-кейс: переиспользование PK после failover | Snowflake/ULID/составные ключи (раздел 10) |
| «Кворум = строгая консистентность» | W+R>N «значит линейноизуемо» | Кворум даёт свежесть, не linearizability (§4.2, раздел 09) |
| Hot partition без плана | Селебрити кладёт шард | Salting / выделенный шард / кэш + реплики ключа (§8) |
| Failover «как-нибудь» | Нет ответа на split brain и потерю хвоста | Три шага + три ловушки — проговорить (§5) |
12. Что назвать на секции (чек-лист раздела)
Обязательные фразы и действия, по которым видно, что раздел освоен. Полный чек-лист всех этапов — ../methodology.md §7; канон — ../knowledge-base.md §5–6.
На этапе 4 (детали + БД): - При первом же блоке БД: «primary + реплики, semi-sync: переживаем отказ одной ноды без потери подтверждённых записей». - «Чтение с реплик — eventual с лагом; read-your-writes закрываю чтением своих данных с лидера / по позиции в логе» — одной фразой. - Перед шардированием — лестница: «сначала индексы, вертикаль, реплики, кэш; шарды — когда storage/write не влезает в ноду» + число-обоснование (70 ТБ/год). - «Шард-ключ = access pattern: лента — посты по post_id, соцграф и кэш лент по user_id; чат — по conversation_id; менять ключ после — переливка всего». - «Раскладка — consistent hashing с виртуальными нодами: смена узла двигает ~1/N ключей» + «селебрити — hot partition: salting / выделенный шард / fan-out on read». - Если есть вторичные запросы: «local-индексы + scatter-gather для редких, global — ценой лага для частых».
На этапе 6 (эксплуатация): - Пройти таблицу §10: реплика → лидер → нода кэша → лидер шарда → ДЦ, каждой строке — метрика. - Failover-скетч: «детект по heartbeat с кворумом → выборы самой свежей реплики → fencing старого лидера; async-хвост может потеряться — поэтому semi-sync». - «Мониторю replication lag как SLO, алерт при отставании на секунды» — лаг это не тюнинг, это дежурная метрика. - Отказ ноды кэша связать с stampede (single-flight) — мост в раздел 06.
Сквозные формулы раздела: - «Репликация — копии, шардирование — разрезы; в бою — оба: каждый шард реплицирован». - «Кворум W+R>N даёт свежесть, но не linearizability» — если сказал, дальше про консенсус спросят вас, а не вы оправдаетесь. - «Ключ шардирования = гранулярность горячего чтения».
13. Самопроверка «умею, если»
Критерий раздела из брифа: могу обосновать схему репликации и ключ шардирования для ленты/чата и что происходит при отказе ноды. Проверяется прогоном с закрытыми материалами: для ленты и для чата вслух по §9 — «ключ = гранулярность горячего чтения» с числами-обоснованием (70 ТБ/год постов → шарды; 100 ГБ при 500k RPS → реплики, не шарды); затем таблица отказов §10 от руки — все шесть строк с механизмами и метриками; затем трейдофф-мины: sync/async/semi-sync одной фразой каждый, кворумы с оговоркой про linearizability, ловушки failover (async-хвост, split brain, ложный детект) — поимённо. Контрольное время: полный ответ «репликация + шарды + отказы» для одного сервиса — 5–7 минут этапа 4–6. Протокол тренировок и рубрика 0–3 — ../methodology.md §6; сверка с эталонами — ../classic-designs.md §2–3 (лента, чат) и §1 (URL shortener: «реплицировать, а не шардировать»). Если лента и чат обосновываются без пауз, а на «а что если упадёт…» рука идёт к таблице, а не к фантазии — раздел готов.
Связанные документы
00-overview.md— обзорная статья по всем 12 разделам брифа (TASK-45.40)03-scaling.md— лестница масштабирования: репликация и шарды как её ступени 4–5; stateless-ноды и ротация при отказах02-estimations.md— числа-обоснования «пора/не пора» (storage, 80/20), на которые опирается лестница §6.1../brief.md— бриф: раздел 05 с критерием «умею, если»../knowledge-base.md— §5 репликация, §6 шардирование, §7 консистентность (кворумы, векторные часы, fencing), §12 карта DDIA../methodology.md— §2 карточки этапов 4–6, §4 топ-10 ошибок, §6 протокол тренировок, §7 одностраничная карточка на секцию../classic-designs.md— §1 URL shortener («реплицировать, не шардировать», KV-зеркало), §2 лента (шард-ключи по user_id/post_id, селебрити), §3 мессенджер (conversation_id, registry)../materials/ddia-2ed/ch-06-replication.md·../materials/ddia-2ed/ch-07-sharding.md— источники: полный текст гл. 6–7 DDIA 2-го изд.../materials/alex-xu-vol1/ch-05-consistent-hashing.md·../materials/alex-xu-vol1/ch-06-key-value-store.md— источники: consistent hashing и Dynamo-хранилище- Соседние статьи:
04-data-models-storage.md(какие хранилища режем/реплицируем) ·06-caching.md(кэш-кластер и stampede при отказе ноды) ·08-fault-tolerance.md(мульти-ДЦ, majority-коммит) ·09-consistency.md(кворумы и линейзируемость глубже)