Job 2026 md

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 — три шага:

  1. Детект. Heartbeat'ы; лидер молчит таймаут (~30 с) — считается мёртвым. Ловушка таймаута: длинный — дольше простой, короткий — ложные failover при всплеске нагрузки или сетевом чихе (ложный failover сам является инцидентом — при уже больной системе делает хуже). Детект — кворумом реплик, не одной нодой.
  2. Выборы нового лидера. Голосование большинства реплик или решение контролера; кандидат — самая свежая реплика (по позиции в логе): теряем минимум подтверждённых записей. Согласовать выборы у всех узлов — задача консенсуса (Raft; эталон: managed-СУБД переключают primary именно так).
  3. Реконфигурация. Клиенты пишут новому лидеру; остальные реплики тянут лог с него. Вернувшийся старый лидер обязан стать фолловером.

Ловушки — называть все три:

  • Потеря 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 (кворумы и линейзируемость глубже)